the heartbeat reply ledgerupdate is ok
This commit is contained in:
@@ -44,6 +44,8 @@ extern std::string G_ROCKETMQ_TOPIC;//topie
|
||||
extern std::string G_ROCKETMQ_TAG;//tag
|
||||
extern std::string G_ROCKETMQ_KEY;//key
|
||||
|
||||
extern std::string FRONT_INST;
|
||||
|
||||
#ifdef __cplusplus
|
||||
extern "C" {
|
||||
#endif
|
||||
@@ -283,14 +285,14 @@ RocketMQConsumer::~RocketMQConsumer()
|
||||
// 在 RocketMQConsumer 类中新增函数用来设置消费模式
|
||||
void RocketMQConsumer::setConsumerMessageModel(const std::string& topic)
|
||||
{
|
||||
if (topic == G_MQCONSUMER_TOPIC_SET) {
|
||||
/*if (topic == G_MQCONSUMER_TOPIC_SET) {
|
||||
// 设置为普通消费模式
|
||||
if (SetPushConsumerMessageModel(consumer_, CLUSTERING) != 0) {
|
||||
std::cout << "Error setting message model to CLUSTERING for topic: " << topic << std::endl;
|
||||
} else {
|
||||
std::cout << "Set consumer to CLUSTERING for topic: " << topic << std::endl;
|
||||
}
|
||||
} else {
|
||||
} else*/ {
|
||||
// 默认设置为广播消费模式
|
||||
if (SetPushConsumerMessageModel(consumer_, BROADCASTING) != 0) {
|
||||
std::cout << "Error setting message model to BROADCASTING for topic: " << topic << std::endl;
|
||||
@@ -649,7 +651,7 @@ void rocketmq_test_rt()
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_RT);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_RT);
|
||||
std::ifstream file("rt.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
@@ -663,7 +665,7 @@ void rocketmq_test_ud()//用来测试台账更新
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_UD);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_UD);
|
||||
std::ifstream file("ud.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
@@ -677,7 +679,7 @@ void rocketmq_test_set()//用来测试进程控制脚本
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_SET);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_SET);
|
||||
std::ifstream file("set.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
@@ -691,7 +693,7 @@ void rocketmq_test_only()//用来测试进程控制脚本
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_SET);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_SET);
|
||||
std::ifstream file("set_debug.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
@@ -706,7 +708,7 @@ void rocketmq_test_rc()
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_RC);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_RC);
|
||||
std::ifstream file("rc.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
@@ -721,7 +723,7 @@ void rocketmq_test_log()
|
||||
{
|
||||
Ckafka_data_t data;
|
||||
data.monitor_id = 123123;
|
||||
data.strTopic = QString::fromStdString(G_MQCONSUMER_TOPIC_LOG);
|
||||
data.strTopic = QString::fromStdString(std::string(FRONT_INST) + "_" + G_MQCONSUMER_TOPIC_LOG);
|
||||
std::ifstream file("log_test.txt"); // 文件中存储长字符串
|
||||
std::stringstream buffer;
|
||||
buffer << file.rdbuf(); // 读取整个文件内容
|
||||
|
||||
Reference in New Issue
Block a user