提交 e640ff99 编写于 作者: hebao@lab.ibiz5.com's avatar hebao@lab.ibiz5.com

MQ消费这个支持groupname配置

上级 6f2fbffd
......@@ -20,8 +20,10 @@ import org.springframework.context.annotation.Configuration;
@Configuration
public class MQConsumerConfigure {
@Value("${rocketmq.consumer.groupName:}")
private String groupName;
@Value("${rocketmq.consumer.ruleEngineGroupName:dstCG}")
private String ruleEngineGroupName;
@Value("${rocketmq.consumer.resultsGroupName:resultsCG}")
private String resultsGroupName;
@Value("${rocketmq.consumer.namesrvAddr:}")
private String namesrvAddr;
@Value("${rocketmq.consumer.topics:}")
......@@ -52,8 +54,7 @@ public class MQConsumerConfigure {
//@ConditionalOnExpression("${rocketmq.consumer.isOnOff:off}.equals('on')")
public DefaultMQPushConsumer defaultConsumer() throws MQClientException {
log.info("defaultConsumer 正在创建---------------------------------------");
String groupName = "dstCG";
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(groupName);
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(ruleEngineGroupName);
consumer.setNamesrvAddr(namesrvAddr);
consumer.setConsumeThreadMin(1);
consumer.setConsumeThreadMax(1);
......@@ -65,7 +66,7 @@ public class MQConsumerConfigure {
try {
consumer.subscribe(ruleEngineTopic, "*");
consumer.start();
log.info("consumer 创建成功 groupName={}, topics={}, namesrvAddr={}",groupName,topics,namesrvAddr);
log.info("consumer 创建成功 groupName={}, topics={}, namesrvAddr={}",ruleEngineGroupName,topics,namesrvAddr);
} catch (MQClientException e) {
log.error("consumer 创建失败!");
}
......@@ -75,8 +76,7 @@ public class MQConsumerConfigure {
@Bean
public DefaultMQPushConsumer resultsMQMsgConsumer() throws MQClientException {
log.info("resultsMQMsgConsumer 正在创建---------------------------------------");
String groupName = "resultsCG";
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(groupName);
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(resultsGroupName);
consumer.setNamesrvAddr(namesrvAddr);
consumer.setConsumeThreadMin(10);
consumer.setConsumeThreadMax(20);
......@@ -88,7 +88,7 @@ public class MQConsumerConfigure {
// 设置该消费者订阅的主题和tag,如果订阅该主题下的所有tag,则使用*,
consumer.subscribe(resultsTopic, "*");
consumer.start();
log.info("consumer 创建成功 groupName={}, topics={}, namesrvAddr={}",groupName,topics,namesrvAddr);
log.info("consumer 创建成功 groupName={}, topics={}, namesrvAddr={}",resultsGroupName,topics,namesrvAddr);
} catch (MQClientException e) {
log.error("consumer 创建失败!");
}
......
Markdown 格式
0% or
您添加了 0 到此讨论。请谨慎行事。
先完成此消息的编辑!
想要评论请 注册