修改mq库
This commit is contained in:
@@ -2384,7 +2384,7 @@ char* find_mp_name_from_mp_id(const char* mp_id)
|
||||
ConsumeStatus myMessageCallbackrtdata(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if(INITFLAG != 1)return 1;//防止崩溃
|
||||
if(INITFLAG != 1)return rocketmq::RECONSUME_LATER;//防止崩溃
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
@@ -2479,7 +2479,7 @@ ConsumeStatus myMessageCallbackrtdata(
|
||||
ConsumeStatus myMessageCallbackupdate(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if(INITFLAG != 1)return 1;//防止崩溃
|
||||
if(INITFLAG != 1)return rocketmq::RECONSUME_LATER;//防止崩溃
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
@@ -2550,7 +2550,7 @@ ConsumeStatus myMessageCallbackupdate(
|
||||
ConsumeStatus myMessageCallbackset(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if(INITFLAG != 1)return 1;//防止崩溃
|
||||
if(INITFLAG != 1)return rocketmq::RECONSUME_LATER;//防止崩溃
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
@@ -2613,7 +2613,7 @@ ConsumeStatus myMessageCallbackset(
|
||||
ConsumeStatus myMessageCallbacklog(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if(INITFLAG != 1)return 1;//防止崩溃
|
||||
if(INITFLAG != 1)return rocketmq::RECONSUME_LATER;//防止崩溃
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
@@ -2676,7 +2676,7 @@ ConsumeStatus myMessageCallbacklog(
|
||||
ConsumeStatus myMessageCallbackrecall(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if(INITFLAG != 1)return 1;//防止崩溃
|
||||
if(INITFLAG != 1)return rocketmq::RECONSUME_LATER;//防止崩溃
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
@@ -2752,7 +2752,7 @@ ConsumeStatus myMessageCallbackrecall(
|
||||
ConsumeStatus myMessageCallbackfile(
|
||||
const rocketmq::MQMessageExt& msg)
|
||||
{
|
||||
if (INITFLAG != 1) return 1;
|
||||
if (INITFLAG != 1)return rocketmq::RECONSUME_LATER;
|
||||
|
||||
if (!should_process_after_start(msg)) {
|
||||
std::cout << "[MQ] skip old message: "
|
||||
|
||||
Reference in New Issue
Block a user