From 1c2301aa8226202151994a16f1a49912ac55e320 Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 24 Dec 2023 09:41:03 +0800 Subject: [PATCH 1/8] build: Restrict the Snapshot Daily Release Automation action to be triggered only in the official repository --- .github/workflows/snapshot-automation.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/snapshot-automation.yml b/.github/workflows/snapshot-automation.yml index 88f5f4e0ccb..99855d3aa0d 100644 --- a/.github/workflows/snapshot-automation.yml +++ b/.github/workflows/snapshot-automation.yml @@ -36,6 +36,7 @@ env: jobs: dist-tar: + if: github.repository == 'apache/rocketmq' name: Build dist tar runs-on: ubuntu-latest timeout-minutes: 30 From e31ae9adbe3099c12ffb9f64a3e58f8e81ac10f4 Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 22 Jun 2025 01:52:38 +0800 Subject: [PATCH 2/8] [#ISSUE 9422]: change defaultInstanceName to pid + groupName --- .../org/apache/rocketmq/client/ClientConfig.java | 4 ++-- .../impl/consumer/DefaultLitePullConsumerImpl.java | 2 +- .../impl/consumer/DefaultMQPullConsumerImpl.java | 2 +- .../impl/consumer/DefaultMQPushConsumerImpl.java | 2 +- .../client/impl/producer/DefaultMQProducerImpl.java | 2 +- .../consumer/DefaultLitePullConsumerTest.java | 2 +- .../client/consumer/DefaultMQPushConsumerTest.java | 2 +- .../ConsumeMessageConcurrentlyServiceTest.java | 2 +- .../consumer/ConsumeMessageOrderlyServiceTest.java | 2 +- .../trace/DefaultMQConsumerWithOpenTracingTest.java | 2 +- .../trace/DefaultMQConsumerWithTraceTest.java | 2 +- .../DefaultMQLitePullConsumerWithTraceTest.java | 2 +- .../rocketmq/example/quickstart/Consumer.java | 13 +++++++++++-- .../rocketmq/example/quickstart/Producer.java | 4 ++-- .../rocketmq/tools/admin/DefaultMQAdminExtImpl.java | 2 +- 15 files changed, 27 insertions(+), 18 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java index 696b073b373..3b357d5a308 100644 --- a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java +++ b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java @@ -154,9 +154,9 @@ public void setInstanceName(String instanceName) { this.instanceName = instanceName; } - public void changeInstanceNameToPID() { + public void changeInstanceNameToPIDWithGroupInfo(String groupInfo) { if (this.instanceName.equals("DEFAULT")) { - this.instanceName = UtilAll.getPid() + "#" + System.nanoTime(); + this.instanceName = UtilAll.getPid() + "#" + groupInfo; } } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java index f85dcc7b459..cb555bb67ce 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java @@ -286,7 +286,7 @@ public synchronized void start() throws MQClientException { this.checkConfig(); if (this.defaultLitePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultLitePullConsumer.changeInstanceNameToPID(); + this.defaultLitePullConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultLitePullConsumer.getConsumerGroup()); } initScheduledThreadPoolExecutor(); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java index 9d46e28f5d4..30cc8e1bd02 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java @@ -710,7 +710,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPullConsumer.changeInstanceNameToPID(); + this.defaultMQPullConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQPullConsumer.getConsumerGroup()); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPullConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java index 0ae779971c8..4ff42c2fa9e 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java @@ -932,7 +932,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPushConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPushConsumer.changeInstanceNameToPID(); + this.defaultMQPushConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQPushConsumer.getConsumerGroup()); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPushConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index 4aa605821f0..3f166efd711 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -248,7 +248,7 @@ public void start(final boolean startFactory) throws MQClientException { this.checkConfig(); if (!this.defaultMQProducer.getProducerGroup().equals(MixAll.CLIENT_INNER_PRODUCER_GROUP)) { - this.defaultMQProducer.changeInstanceNameToPID(); + this.defaultMQProducer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQProducer.getProducerGroup()); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQProducer, rpcHook); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java index 592c247057b..d87b2edaee2 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java @@ -892,7 +892,7 @@ private PullResultExt createPullResult(PullMessageRequestHeader requestHeader, P private void suppressUpdateTopicRouteInfoFromNameServer( DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPID(); + litePullConsumer.changeInstanceNameToPIDWithGroupInfo(litePullConsumer.getConsumerGroup()); } ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java index 834be5cf16f..301dfb115fa 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java @@ -152,7 +152,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java index 395c0ff2335..f998b79a17c 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java @@ -118,7 +118,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo field.set(pushConsumerImpl, rebalancePushImpl); pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java index 5fa78b70090..d99c988e45a 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java @@ -109,7 +109,7 @@ public ConsumeOrderlyStatus consumeMessage(List msgs, ConsumeOrderly pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java index 028445ef2d7..1edbbd69c50 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java @@ -153,7 +153,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java index fc63cce1ce4..d3dad70641c 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java @@ -142,7 +142,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true))); factoryTable.put(pushConsumer.buildMQClientId(), mQClientFactory); doReturn(false).when(mQClientFactory).updateTopicRouteInfoFromNameServer(anyString()); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java index e0573bdfb0b..e398fe7b7a8 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java @@ -308,7 +308,7 @@ private SendResult createSendResult(SendStatus sendStatus) { private static void suppressUpdateTopicRouteInfoFromNameServer(DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { DefaultLitePullConsumerImpl defaultLitePullConsumerImpl = (DefaultLitePullConsumerImpl) FieldUtils.readDeclaredField(litePullConsumer, "defaultLitePullConsumerImpl", true); if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPID(); + litePullConsumer.changeInstanceNameToPIDWithGroupInfo(litePullConsumer.getConsumerGroup()); } MQClientInstance mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(litePullConsumer, (RPCHook) FieldUtils.readDeclaredField(defaultLitePullConsumerImpl, "rpcHook", true))); ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java index 3a101bf664f..cbfc59776d2 100644 --- a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java +++ b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java @@ -27,7 +27,7 @@ */ public class Consumer { - public static final String CONSUMER_GROUP = "please_rename_unique_group_name_4"; + public static final String CONSUMER_GROUP = "c_please_rename_unique_group_name_4"; public static final String DEFAULT_NAMESRVADDR = "127.0.0.1:9876"; public static final String TOPIC = "TopicTest"; @@ -37,6 +37,7 @@ public static void main(String[] args) throws MQClientException { * Instantiate with specified consumer group name. */ DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP); + DefaultMQPushConsumer consumer2 = new DefaultMQPushConsumer(CONSUMER_GROUP); /* * Specify name server addresses. @@ -50,17 +51,20 @@ public static void main(String[] args) throws MQClientException { * */ // Uncomment the following line while debugging, namesrvAddr should be set to your local address - // consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + consumer2.setNamesrvAddr(DEFAULT_NAMESRVADDR); /* * Specify where to start in case the specific consumer group is a brand-new one. */ consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); + consumer2.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); /* * Subscribe one more topic to consume. */ consumer.subscribe(TOPIC, "*"); + consumer2.subscribe(TOPIC, "*"); /* * Register callback to execute on arrival of messages fetched from brokers. @@ -69,11 +73,16 @@ public static void main(String[] args) throws MQClientException { System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); + consumer2.registerMessageListener((MessageListenerConcurrently) (msg, context) -> { + System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); + return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; + }); /* * Launch the consumer instance. */ consumer.start(); + consumer2.start(); System.out.printf("Consumer Started.%n"); } diff --git a/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java b/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java index cae16cec52d..7c62efff6a2 100644 --- a/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java +++ b/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java @@ -31,7 +31,7 @@ public class Producer { * The number of produced messages. */ public static final int MESSAGE_COUNT = 1000; - public static final String PRODUCER_GROUP = "please_rename_unique_group_name"; + public static final String PRODUCER_GROUP = "p_please_rename_unique_group_name"; public static final String DEFAULT_NAMESRVADDR = "127.0.0.1:9876"; public static final String TOPIC = "TopicTest"; public static final String TAG = "TagA"; @@ -54,7 +54,7 @@ public static void main(String[] args) throws MQClientException, InterruptedExce * */ // Uncomment the following line while debugging, namesrvAddr should be set to your local address - // producer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + producer.setNamesrvAddr(DEFAULT_NAMESRVADDR); /* * Launch the instance. diff --git a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java index 1bdcc765d61..79064030936 100644 --- a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java +++ b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java @@ -155,7 +155,7 @@ public void start() throws MQClientException { case CREATE_JUST: this.serviceState = ServiceState.START_FAILED; - this.defaultMQAdminExt.changeInstanceNameToPID(); + this.defaultMQAdminExt.changeInstanceNameToPIDWithGroupInfo(this.defaultMQAdminExt.getAdminExtGroup()); if ("{}".equals(this.defaultMQAdminExt.getSocksProxyConfig())) { String proxyConfig = System.getenv(SOCKS_PROXY_JSON); From d95228b2d5bf70c649442e170bd32c753c79bc04 Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 22 Jun 2025 01:55:37 +0800 Subject: [PATCH 3/8] Revert unwanted changes --- .../apache/rocketmq/example/quickstart/Consumer.java | 11 +---------- .../apache/rocketmq/example/quickstart/Producer.java | 2 +- 2 files changed, 2 insertions(+), 11 deletions(-) diff --git a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java index cbfc59776d2..fa1af9d33b9 100644 --- a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java +++ b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java @@ -37,7 +37,6 @@ public static void main(String[] args) throws MQClientException { * Instantiate with specified consumer group name. */ DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP); - DefaultMQPushConsumer consumer2 = new DefaultMQPushConsumer(CONSUMER_GROUP); /* * Specify name server addresses. @@ -51,20 +50,17 @@ public static void main(String[] args) throws MQClientException { * */ // Uncomment the following line while debugging, namesrvAddr should be set to your local address - consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); - consumer2.setNamesrvAddr(DEFAULT_NAMESRVADDR); + // consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); /* * Specify where to start in case the specific consumer group is a brand-new one. */ consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); - consumer2.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); /* * Subscribe one more topic to consume. */ consumer.subscribe(TOPIC, "*"); - consumer2.subscribe(TOPIC, "*"); /* * Register callback to execute on arrival of messages fetched from brokers. @@ -73,16 +69,11 @@ public static void main(String[] args) throws MQClientException { System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); - consumer2.registerMessageListener((MessageListenerConcurrently) (msg, context) -> { - System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); - return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; - }); /* * Launch the consumer instance. */ consumer.start(); - consumer2.start(); System.out.printf("Consumer Started.%n"); } diff --git a/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java b/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java index 7c62efff6a2..5665517715e 100644 --- a/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java +++ b/example/src/main/java/org/apache/rocketmq/example/quickstart/Producer.java @@ -54,7 +54,7 @@ public static void main(String[] args) throws MQClientException, InterruptedExce * */ // Uncomment the following line while debugging, namesrvAddr should be set to your local address - producer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + // producer.setNamesrvAddr(DEFAULT_NAMESRVADDR); /* * Launch the instance. From eeec75bbcbeb6ba414f771f4a9029e9c0b3f843b Mon Sep 17 00:00:00 2001 From: totalo Date: Sat, 28 Jun 2025 22:25:35 +0800 Subject: [PATCH 4/8] Revert unwanted changes --- .../main/java/org/apache/rocketmq/client/ClientConfig.java | 4 ++-- .../client/impl/consumer/DefaultLitePullConsumerImpl.java | 2 +- .../client/impl/consumer/DefaultMQPullConsumerImpl.java | 2 +- .../client/impl/consumer/DefaultMQPushConsumerImpl.java | 2 +- .../rocketmq/client/impl/producer/DefaultMQProducerImpl.java | 2 +- .../rocketmq/client/consumer/DefaultLitePullConsumerTest.java | 2 +- .../rocketmq/client/consumer/DefaultMQPushConsumerTest.java | 2 +- .../impl/consumer/ConsumeMessageConcurrentlyServiceTest.java | 2 +- .../impl/consumer/ConsumeMessageOrderlyServiceTest.java | 2 +- .../client/trace/DefaultMQConsumerWithOpenTracingTest.java | 2 +- .../rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java | 2 +- .../client/trace/DefaultMQLitePullConsumerWithTraceTest.java | 2 +- .../src/test/java/org/apache/rocketmq/test/base/BaseConf.java | 1 + .../apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java | 2 +- 14 files changed, 15 insertions(+), 14 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java index 3b357d5a308..696b073b373 100644 --- a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java +++ b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java @@ -154,9 +154,9 @@ public void setInstanceName(String instanceName) { this.instanceName = instanceName; } - public void changeInstanceNameToPIDWithGroupInfo(String groupInfo) { + public void changeInstanceNameToPID() { if (this.instanceName.equals("DEFAULT")) { - this.instanceName = UtilAll.getPid() + "#" + groupInfo; + this.instanceName = UtilAll.getPid() + "#" + System.nanoTime(); } } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java index cb555bb67ce..f85dcc7b459 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java @@ -286,7 +286,7 @@ public synchronized void start() throws MQClientException { this.checkConfig(); if (this.defaultLitePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultLitePullConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultLitePullConsumer.getConsumerGroup()); + this.defaultLitePullConsumer.changeInstanceNameToPID(); } initScheduledThreadPoolExecutor(); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java index 30cc8e1bd02..9d46e28f5d4 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java @@ -710,7 +710,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPullConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQPullConsumer.getConsumerGroup()); + this.defaultMQPullConsumer.changeInstanceNameToPID(); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPullConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java index 4ff42c2fa9e..0ae779971c8 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java @@ -932,7 +932,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPushConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPushConsumer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQPushConsumer.getConsumerGroup()); + this.defaultMQPushConsumer.changeInstanceNameToPID(); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPushConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index 3f166efd711..4aa605821f0 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -248,7 +248,7 @@ public void start(final boolean startFactory) throws MQClientException { this.checkConfig(); if (!this.defaultMQProducer.getProducerGroup().equals(MixAll.CLIENT_INNER_PRODUCER_GROUP)) { - this.defaultMQProducer.changeInstanceNameToPIDWithGroupInfo(this.defaultMQProducer.getProducerGroup()); + this.defaultMQProducer.changeInstanceNameToPID(); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQProducer, rpcHook); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java index d87b2edaee2..592c247057b 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java @@ -892,7 +892,7 @@ private PullResultExt createPullResult(PullMessageRequestHeader requestHeader, P private void suppressUpdateTopicRouteInfoFromNameServer( DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPIDWithGroupInfo(litePullConsumer.getConsumerGroup()); + litePullConsumer.changeInstanceNameToPID(); } ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java index 301dfb115fa..834be5cf16f 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java @@ -152,7 +152,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); + pushConsumer.changeInstanceNameToPID(); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java index f998b79a17c..395c0ff2335 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java @@ -118,7 +118,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo field.set(pushConsumerImpl, rebalancePushImpl); pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); + pushConsumer.changeInstanceNameToPID(); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java index d99c988e45a..5fa78b70090 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java @@ -109,7 +109,7 @@ public ConsumeOrderlyStatus consumeMessage(List msgs, ConsumeOrderly pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); + pushConsumer.changeInstanceNameToPID(); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java index 1edbbd69c50..028445ef2d7 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java @@ -153,7 +153,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); + pushConsumer.changeInstanceNameToPID(); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java index d3dad70641c..fc63cce1ce4 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java @@ -142,7 +142,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPIDWithGroupInfo(pushConsumer.getConsumerGroup()); + pushConsumer.changeInstanceNameToPID(); mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true))); factoryTable.put(pushConsumer.buildMQClientId(), mQClientFactory); doReturn(false).when(mQClientFactory).updateTopicRouteInfoFromNameServer(anyString()); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java index e398fe7b7a8..e0573bdfb0b 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java @@ -308,7 +308,7 @@ private SendResult createSendResult(SendStatus sendStatus) { private static void suppressUpdateTopicRouteInfoFromNameServer(DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { DefaultLitePullConsumerImpl defaultLitePullConsumerImpl = (DefaultLitePullConsumerImpl) FieldUtils.readDeclaredField(litePullConsumer, "defaultLitePullConsumerImpl", true); if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPIDWithGroupInfo(litePullConsumer.getConsumerGroup()); + litePullConsumer.changeInstanceNameToPID(); } MQClientInstance mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(litePullConsumer, (RPCHook) FieldUtils.readDeclaredField(defaultLitePullConsumerImpl, "rpcHook", true))); ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java index 472e106ce35..6b6fa512bd9 100644 --- a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java +++ b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java @@ -54,6 +54,7 @@ import org.apache.rocketmq.tools.admin.DefaultMQAdminExt; import org.apache.rocketmq.tools.admin.MQAdminExt; import org.junit.Assert; +import org.junit.jupiter.api.BeforeEach; import static org.apache.rocketmq.test.base.IntegrationTestBase.initMQAdmin; import static org.awaitility.Awaitility.await; diff --git a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java index 79064030936..1bdcc765d61 100644 --- a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java +++ b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java @@ -155,7 +155,7 @@ public void start() throws MQClientException { case CREATE_JUST: this.serviceState = ServiceState.START_FAILED; - this.defaultMQAdminExt.changeInstanceNameToPIDWithGroupInfo(this.defaultMQAdminExt.getAdminExtGroup()); + this.defaultMQAdminExt.changeInstanceNameToPID(); if ("{}".equals(this.defaultMQAdminExt.getSocksProxyConfig())) { String proxyConfig = System.getenv(SOCKS_PROXY_JSON); From c06eb2ea409207ff890621788f223274b1d46271 Mon Sep 17 00:00:00 2001 From: totalo Date: Sat, 28 Jun 2025 23:57:38 +0800 Subject: [PATCH 5/8] fix checkstyle --- test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java | 1 - 1 file changed, 1 deletion(-) diff --git a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java index 6b6fa512bd9..472e106ce35 100644 --- a/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java +++ b/test/src/test/java/org/apache/rocketmq/test/base/BaseConf.java @@ -54,7 +54,6 @@ import org.apache.rocketmq.tools.admin.DefaultMQAdminExt; import org.apache.rocketmq.tools.admin.MQAdminExt; import org.junit.Assert; -import org.junit.jupiter.api.BeforeEach; import static org.apache.rocketmq.test.base.IntegrationTestBase.initMQAdmin; import static org.awaitility.Awaitility.await; From e186a35e37a5c91ff13ed253bf35f74b59fd31ac Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 29 Jun 2025 11:05:47 +0800 Subject: [PATCH 6/8] Change default client instance name --- .../org/apache/rocketmq/client/ClientConfig.java | 4 ++-- .../impl/consumer/DefaultLitePullConsumerImpl.java | 2 +- .../impl/consumer/DefaultMQPullConsumerImpl.java | 2 +- .../impl/consumer/DefaultMQPushConsumerImpl.java | 2 +- .../client/impl/producer/DefaultMQProducerImpl.java | 3 +-- .../client/consumer/DefaultLitePullConsumerTest.java | 2 +- .../client/consumer/DefaultMQPushConsumerTest.java | 2 +- .../ConsumeMessageConcurrentlyServiceTest.java | 2 +- .../consumer/ConsumeMessageOrderlyServiceTest.java | 2 +- .../trace/DefaultMQConsumerWithOpenTracingTest.java | 2 +- .../client/trace/DefaultMQConsumerWithTraceTest.java | 2 +- .../DefaultMQLitePullConsumerWithTraceTest.java | 2 +- .../apache/rocketmq/example/quickstart/Consumer.java | 12 +++++++++++- .../client.consumer.DefaultLitePullConsumer.schema | 2 +- .../api/client.consumer.DefaultMQPullConsumer.schema | 2 +- .../api/client.consumer.DefaultMQPushConsumer.schema | 2 +- .../api/client.producer.DefaultMQProducer.schema | 2 +- .../schema/api/tools.admin.DefaultMQAdminExt.schema | 2 +- .../rocketmq/tools/admin/DefaultMQAdminExtImpl.java | 2 +- 19 files changed, 30 insertions(+), 21 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java index 696b073b373..18012fa8044 100644 --- a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java +++ b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java @@ -154,9 +154,9 @@ public void setInstanceName(String instanceName) { this.instanceName = instanceName; } - public void changeInstanceNameToPID() { + public void changeInstanceNameToIpWithPidAndGroupInfo(String groupInfo) { if (this.instanceName.equals("DEFAULT")) { - this.instanceName = UtilAll.getPid() + "#" + System.nanoTime(); + this.instanceName = this.clientIP + "#" + UtilAll.getPid() + "#" + groupInfo; } } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java index f85dcc7b459..d76051b0319 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultLitePullConsumerImpl.java @@ -286,7 +286,7 @@ public synchronized void start() throws MQClientException { this.checkConfig(); if (this.defaultLitePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultLitePullConsumer.changeInstanceNameToPID(); + this.defaultLitePullConsumer.changeInstanceNameToIpWithPidAndGroupInfo(this.defaultLitePullConsumer.getConsumerGroup()); } initScheduledThreadPoolExecutor(); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java index 9d46e28f5d4..29af4350680 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPullConsumerImpl.java @@ -710,7 +710,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPullConsumer.changeInstanceNameToPID(); + this.defaultMQPullConsumer.changeInstanceNameToIpWithPidAndGroupInfo(this.defaultMQPullConsumer.getConsumerGroup()); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPullConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java index 0ae779971c8..7662074052a 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java @@ -932,7 +932,7 @@ public synchronized void start() throws MQClientException { this.copySubscription(); if (this.defaultMQPushConsumer.getMessageModel() == MessageModel.CLUSTERING) { - this.defaultMQPushConsumer.changeInstanceNameToPID(); + this.defaultMQPushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(this.defaultMQPushConsumer.getConsumerGroup()); } this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQPushConsumer, this.rpcHook); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index 4aa605821f0..2baea0766be 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -248,9 +248,8 @@ public void start(final boolean startFactory) throws MQClientException { this.checkConfig(); if (!this.defaultMQProducer.getProducerGroup().equals(MixAll.CLIENT_INNER_PRODUCER_GROUP)) { - this.defaultMQProducer.changeInstanceNameToPID(); + this.defaultMQProducer.changeInstanceNameToIpWithPidAndGroupInfo(this.defaultMQProducer.getProducerGroup()); } - this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultMQProducer, rpcHook); defaultMQProducer.initProduceAccumulator(); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java index 592c247057b..cc66e7ee050 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultLitePullConsumerTest.java @@ -892,7 +892,7 @@ private PullResultExt createPullResult(PullMessageRequestHeader requestHeader, P private void suppressUpdateTopicRouteInfoFromNameServer( DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPID(); + litePullConsumer.changeInstanceNameToIpWithPidAndGroupInfo(litePullConsumer.getConsumerGroup()); } ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java index 834be5cf16f..83cf8ccf315 100644 --- a/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumerTest.java @@ -152,7 +152,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java index 395c0ff2335..e2ae29dd69b 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyServiceTest.java @@ -118,7 +118,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo field.set(pushConsumerImpl, rebalancePushImpl); pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java index 5fa78b70090..d70da83c2d7 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageOrderlyServiceTest.java @@ -109,7 +109,7 @@ public ConsumeOrderlyStatus consumeMessage(List msgs, ConsumeOrderly pushConsumer.subscribe(topic, "*"); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); mQClientFactory = spy(mQClientFactory); field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory"); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java index 028445ef2d7..17d447fa191 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithOpenTracingTest.java @@ -153,7 +153,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true)); FieldUtils.writeDeclaredField(mQClientFactory, "mQClientAPIImpl", mQClientAPIImpl, true); mQClientFactory = spy(mQClientFactory); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java index fc63cce1ce4..969b3098a09 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQConsumerWithTraceTest.java @@ -142,7 +142,7 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, DefaultMQPushConsumerImpl pushConsumerImpl = pushConsumer.getDefaultMQPushConsumerImpl(); // suppress updateTopicRouteInfoFromNameServer - pushConsumer.changeInstanceNameToPID(); + pushConsumer.changeInstanceNameToIpWithPidAndGroupInfo(pushConsumer.getConsumerGroup()); mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true))); factoryTable.put(pushConsumer.buildMQClientId(), mQClientFactory); doReturn(false).when(mQClientFactory).updateTopicRouteInfoFromNameServer(anyString()); diff --git a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java index e0573bdfb0b..710698efc0c 100644 --- a/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/trace/DefaultMQLitePullConsumerWithTraceTest.java @@ -308,7 +308,7 @@ private SendResult createSendResult(SendStatus sendStatus) { private static void suppressUpdateTopicRouteInfoFromNameServer(DefaultLitePullConsumer litePullConsumer) throws IllegalAccessException { DefaultLitePullConsumerImpl defaultLitePullConsumerImpl = (DefaultLitePullConsumerImpl) FieldUtils.readDeclaredField(litePullConsumer, "defaultLitePullConsumerImpl", true); if (litePullConsumer.getMessageModel() == MessageModel.CLUSTERING) { - litePullConsumer.changeInstanceNameToPID(); + litePullConsumer.changeInstanceNameToIpWithPidAndGroupInfo(litePullConsumer.getConsumerGroup()); } MQClientInstance mQClientFactory = spy(MQClientManager.getInstance().getOrCreateMQClientInstance(litePullConsumer, (RPCHook) FieldUtils.readDeclaredField(defaultLitePullConsumerImpl, "rpcHook", true))); ConcurrentMap factoryTable = (ConcurrentMap) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true); diff --git a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java index fa1af9d33b9..86389e20a30 100644 --- a/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java +++ b/example/src/main/java/org/apache/rocketmq/example/quickstart/Consumer.java @@ -37,6 +37,7 @@ public static void main(String[] args) throws MQClientException { * Instantiate with specified consumer group name. */ DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(CONSUMER_GROUP); + DefaultMQPushConsumer consumer2 = new DefaultMQPushConsumer(CONSUMER_GROUP); /* * Specify name server addresses. @@ -50,17 +51,21 @@ public static void main(String[] args) throws MQClientException { * */ // Uncomment the following line while debugging, namesrvAddr should be set to your local address - // consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + consumer.setNamesrvAddr(DEFAULT_NAMESRVADDR); + consumer2.setNamesrvAddr(DEFAULT_NAMESRVADDR); /* * Specify where to start in case the specific consumer group is a brand-new one. */ consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); + consumer2.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET); /* * Subscribe one more topic to consume. */ consumer.subscribe(TOPIC, "*"); + consumer.subscribe(TOPIC + "1", "*"); + consumer2.subscribe(TOPIC, "*"); /* * Register callback to execute on arrival of messages fetched from brokers. @@ -69,11 +74,16 @@ public static void main(String[] args) throws MQClientException { System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); + consumer2.registerMessageListener((MessageListenerConcurrently) (msg, context) -> { + System.out.printf("%s Receive New Messages: %s %n", Thread.currentThread().getName(), msg); + return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; + }); /* * Launch the consumer instance. */ consumer.start(); + consumer2.start(); System.out.printf("Consumer Started.%n"); } diff --git a/test/src/test/resources/schema/api/client.consumer.DefaultLitePullConsumer.schema b/test/src/test/resources/schema/api/client.consumer.DefaultLitePullConsumer.schema index 79cd6c12595..4fe0d05b170 100644 --- a/test/src/test/resources/schema/api/client.consumer.DefaultLitePullConsumer.schema +++ b/test/src/test/resources/schema/api/client.consumer.DefaultLitePullConsumer.schema @@ -63,7 +63,7 @@ Field useTLS : private boolean false Field vipChannelEnabled : private boolean false Method assign(java.util.Collection) : public throws (void) Method buildMQClientId() : public throws (java.lang.String) -Method changeInstanceNameToPID() : public throws (void) +Method changeInstanceNameToIpWithPidAndGroupInfo(java.lang.String) : public throws (void) Method cloneClientConfig() : public throws (org.apache.rocketmq.client.ClientConfig) Method commit(boolean,java.util.Set) : public throws (void) Method commitSync() : public throws (void) diff --git a/test/src/test/resources/schema/api/client.consumer.DefaultMQPullConsumer.schema b/test/src/test/resources/schema/api/client.consumer.DefaultMQPullConsumer.schema index 92d4cbadd13..89839d93de4 100644 --- a/test/src/test/resources/schema/api/client.consumer.DefaultMQPullConsumer.schema +++ b/test/src/test/resources/schema/api/client.consumer.DefaultMQPullConsumer.schema @@ -47,7 +47,7 @@ Field unitName : private java.lang.String null Field useTLS : private boolean false Field vipChannelEnabled : private boolean false Method buildMQClientId() : public throws (java.lang.String) -Method changeInstanceNameToPID() : public throws (void) +Method changeInstanceNameToIpWithPidAndGroupInfo(java.lang.String) : public throws (void) Method cloneClientConfig() : public throws (org.apache.rocketmq.client.ClientConfig) Method createTopic(int,int,java.lang.String,java.lang.String) : public throws (void) Method createTopic(int,java.lang.String,java.lang.String) : public throws (void) diff --git a/test/src/test/resources/schema/api/client.consumer.DefaultMQPushConsumer.schema b/test/src/test/resources/schema/api/client.consumer.DefaultMQPushConsumer.schema index 3aed1d66149..1e5ccab97a3 100644 --- a/test/src/test/resources/schema/api/client.consumer.DefaultMQPushConsumer.schema +++ b/test/src/test/resources/schema/api/client.consumer.DefaultMQPushConsumer.schema @@ -63,7 +63,7 @@ Field unitName : private java.lang.String null Field useTLS : private boolean false Field vipChannelEnabled : private boolean false Method buildMQClientId() : public throws (java.lang.String) -Method changeInstanceNameToPID() : public throws (void) +Method changeInstanceNameToIpWithPidAndGroupInfo(java.lang.String) : public throws (void) Method cloneClientConfig() : public throws (org.apache.rocketmq.client.ClientConfig) Method createTopic(int,int,java.lang.String,java.lang.String) : public throws (void) Method createTopic(int,java.lang.String,java.lang.String) : public throws (void) diff --git a/test/src/test/resources/schema/api/client.producer.DefaultMQProducer.schema b/test/src/test/resources/schema/api/client.producer.DefaultMQProducer.schema index d1111fb4572..1a0bfe46ae0 100644 --- a/test/src/test/resources/schema/api/client.producer.DefaultMQProducer.schema +++ b/test/src/test/resources/schema/api/client.producer.DefaultMQProducer.schema @@ -50,7 +50,7 @@ Field useTLS : private boolean false Field vipChannelEnabled : private boolean false Method addRetryResponseCode(int) : public throws (void) Method buildMQClientId() : public throws (java.lang.String) -Method changeInstanceNameToPID() : public throws (void) +Method changeInstanceNameToIpWithPidAndGroupInfo(java.lang.String) : public throws (void) Method cloneClientConfig() : public throws (org.apache.rocketmq.client.ClientConfig) Method createTopic(int,int,java.lang.String,java.lang.String) : public throws (void) Method createTopic(int,java.lang.String,java.lang.String) : public throws (void) diff --git a/test/src/test/resources/schema/api/tools.admin.DefaultMQAdminExt.schema b/test/src/test/resources/schema/api/tools.admin.DefaultMQAdminExt.schema index 026c1975dfe..10cc6b6a5fc 100644 --- a/test/src/test/resources/schema/api/tools.admin.DefaultMQAdminExt.schema +++ b/test/src/test/resources/schema/api/tools.admin.DefaultMQAdminExt.schema @@ -41,7 +41,7 @@ Field useTLS : private boolean false Field vipChannelEnabled : private boolean false Method addWritePermOfBroker(java.lang.String,java.lang.String) : public throws (int) Method buildMQClientId() : public throws (java.lang.String) -Method changeInstanceNameToPID() : public throws (void) +Method changeInstanceNameToIpWithPidAndGroupInfo(java.lang.String) : public throws (void) Method cleanExpiredConsumerQueue(java.lang.String) : public throws (boolean) Method cleanExpiredConsumerQueueByAddr(java.lang.String) : public throws (boolean) Method cleanUnusedTopic(java.lang.String) : public throws (boolean) diff --git a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java index 1bdcc765d61..6706a46102e 100644 --- a/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java +++ b/tools/src/main/java/org/apache/rocketmq/tools/admin/DefaultMQAdminExtImpl.java @@ -155,7 +155,7 @@ public void start() throws MQClientException { case CREATE_JUST: this.serviceState = ServiceState.START_FAILED; - this.defaultMQAdminExt.changeInstanceNameToPID(); + this.defaultMQAdminExt.changeInstanceNameToIpWithPidAndGroupInfo(this.defaultMQAdminExt.getAdminExtGroup()); if ("{}".equals(this.defaultMQAdminExt.getSocksProxyConfig())) { String proxyConfig = System.getenv(SOCKS_PROXY_JSON); From 25690d0faedbb9f87c6f235bc06af612d28fd8a6 Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 13 Jul 2025 22:33:05 +0800 Subject: [PATCH 7/8] fix test --- .../java/org/apache/rocketmq/test/util/MQAdminTestUtils.java | 1 + 1 file changed, 1 insertion(+) diff --git a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java index 276d08d8061..e57377469e6 100644 --- a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java +++ b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java @@ -64,6 +64,7 @@ public class MQAdminTestUtils { public static void startAdmin(String nameSrvAddr) throws MQClientException { mqAdminExt = new DefaultMQAdminExt(); mqAdminExt.setNamesrvAddr(nameSrvAddr); + mqAdminExt.setInstanceName(UUID.randomUUID().toString()); mqAdminExt.start(); } From 7da1f404a211528b0f071539dcb3635b75b88a07 Mon Sep 17 00:00:00 2001 From: totalo Date: Sun, 13 Jul 2025 22:40:49 +0800 Subject: [PATCH 8/8] fix compile --- .../apache/rocketmq/store/queue/CombineConsumeQueueStore.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/queue/CombineConsumeQueueStore.java b/store/src/main/java/org/apache/rocketmq/store/queue/CombineConsumeQueueStore.java index cbcc941fe98..6e9bd2f2bba 100644 --- a/store/src/main/java/org/apache/rocketmq/store/queue/CombineConsumeQueueStore.java +++ b/store/src/main/java/org/apache/rocketmq/store/queue/CombineConsumeQueueStore.java @@ -31,6 +31,8 @@ import org.apache.rocketmq.common.Pair; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.common.message.MessageExtBrokerInner; +import org.apache.rocketmq.logging.org.slf4j.Logger; +import org.apache.rocketmq.logging.org.slf4j.LoggerFactory; import org.apache.rocketmq.store.DefaultMessageStore; import org.apache.rocketmq.store.DispatchRequest; import org.apache.rocketmq.store.StoreType; @@ -38,8 +40,6 @@ import org.apache.rocketmq.store.exception.ConsumeQueueException; import org.apache.rocketmq.store.exception.StoreException; import org.rocksdb.RocksDBException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; public class CombineConsumeQueueStore implements ConsumeQueueStoreInterface { private static final Logger log = LoggerFactory.getLogger(LoggerName.STORE_LOGGER_NAME);