mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
This commit is contained in:
+28
-29
@@ -99,38 +99,37 @@ public class ClientManageProcessor implements NettyRequestProcessor {
|
||||
}
|
||||
}
|
||||
|
||||
SubscriptionGroupConfig subscriptionGroupConfig =
|
||||
this.brokerController.getSubscriptionGroupManager().findSubscriptionGroupConfig(
|
||||
consumerData.getGroupName());
|
||||
SubscriptionGroupConfig subscriptionGroupConfig = this.brokerController.getSubscriptionGroupManager()
|
||||
.findSubscriptionGroupConfig(consumerData.getGroupName());
|
||||
boolean isNotifyConsumerIdsChangedEnable = true;
|
||||
if (null != subscriptionGroupConfig) {
|
||||
isNotifyConsumerIdsChangedEnable = subscriptionGroupConfig.isNotifyConsumerIdsChangedEnable();
|
||||
int topicSysFlag = 0;
|
||||
if (consumerData.isUnitMode()) {
|
||||
topicSysFlag = TopicSysFlag.buildSysFlag(false, true);
|
||||
}
|
||||
String newTopic = MixAll.getRetryTopic(consumerData.getGroupName());
|
||||
this.brokerController.getTopicConfigManager().createTopicInSendMessageBackMethod(
|
||||
newTopic,
|
||||
subscriptionGroupConfig.getRetryQueueNums(),
|
||||
PermName.PERM_WRITE | PermName.PERM_READ, hasOrderTopicSub, topicSysFlag);
|
||||
|
||||
if (null == subscriptionGroupConfig) {
|
||||
continue;
|
||||
}
|
||||
if (null != subscriptionGroupConfig) {
|
||||
boolean changed = this.brokerController.getConsumerManager().registerConsumer(
|
||||
consumerData.getGroupName(),
|
||||
clientChannelInfo,
|
||||
consumerData.getConsumeType(),
|
||||
consumerData.getMessageModel(),
|
||||
consumerData.getConsumeFromWhere(),
|
||||
consumerData.getSubscriptionDataSet(),
|
||||
isNotifyConsumerIdsChangedEnable
|
||||
);
|
||||
if (changed) {
|
||||
LOGGER.info(
|
||||
"ClientManageProcessor: registerConsumer info changed, SDK address={}, consumerData={}",
|
||||
RemotingHelper.parseChannelRemoteAddr(ctx.channel()), consumerData.toString());
|
||||
}
|
||||
|
||||
isNotifyConsumerIdsChangedEnable = subscriptionGroupConfig.isNotifyConsumerIdsChangedEnable();
|
||||
int topicSysFlag = 0;
|
||||
if (consumerData.isUnitMode()) {
|
||||
topicSysFlag = TopicSysFlag.buildSysFlag(false, true);
|
||||
}
|
||||
String newTopic = MixAll.getRetryTopic(consumerData.getGroupName());
|
||||
this.brokerController.getTopicConfigManager().createTopicInSendMessageBackMethod(newTopic, subscriptionGroupConfig.getRetryQueueNums(),
|
||||
PermName.PERM_WRITE | PermName.PERM_READ, hasOrderTopicSub, topicSysFlag);
|
||||
|
||||
boolean changed = this.brokerController.getConsumerManager().registerConsumer(
|
||||
consumerData.getGroupName(),
|
||||
clientChannelInfo,
|
||||
consumerData.getConsumeType(),
|
||||
consumerData.getMessageModel(),
|
||||
consumerData.getConsumeFromWhere(),
|
||||
consumerData.getSubscriptionDataSet(),
|
||||
isNotifyConsumerIdsChangedEnable
|
||||
);
|
||||
if (changed) {
|
||||
LOGGER.info("ClientManageProcessor: registerConsumer info changed, SDK address={}, consumerData={}",
|
||||
RemotingHelper.parseChannelRemoteAddr(ctx.channel()), consumerData.toString());
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
for (ProducerData data : heartbeatData.getProducerDataSet()) {
|
||||
|
||||
Reference in New Issue
Block a user