From 134efd26a1265f18b805ef2dfa6f46a3ac3a0f56 Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:04:48 +0800 Subject: [PATCH 1/8] [ISSUE --- .../rocketmq/broker/pop/PopConsumerRocksdbStore.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java index dc68f9d9fe5..3e2fe8a36fe 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerRocksdbStore.java @@ -147,9 +147,11 @@ public List scanExpiredRecords(long lower, long upper, int ma // configure prefix indexing to improve the performance of scans. // However, in the current implementation, this is not the bottleneck. List consumerRecordList = new ArrayList<>(); - try (ReadOptions scanOptions = new ReadOptions() - .setIterateLowerBound(new Slice(ByteBuffer.allocate(Long.BYTES).putLong(lower).array())) - .setIterateUpperBound(new Slice(ByteBuffer.allocate(Long.BYTES).putLong(upper).array())); + try (Slice lowerBound = new Slice(ByteBuffer.allocate(Long.BYTES).putLong(lower).array()); + Slice upperBound = new Slice(ByteBuffer.allocate(Long.BYTES).putLong(upper).array()); + ReadOptions scanOptions = new ReadOptions() + .setIterateLowerBound(lowerBound) + .setIterateUpperBound(upperBound); RocksIterator iterator = db.newIterator(this.columnFamilyHandle, scanOptions)) { iterator.seek(ByteBuffer.allocate(Long.BYTES).putLong(lower).array()); while (iterator.isValid() && consumerRecordList.size() < maxCount) { From 3da362b6cae0d89d5a4d38601b9b347bdb27d18c Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:05:52 +0800 Subject: [PATCH 2/8] [ISSUE #10898] Wake NettyEventExecutor immediately on shutdown --- .../rocketmq/remoting/netty/NettyRemotingAbstract.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java index a735f8455d3..b9bf58ccb44 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingAbstract.java @@ -784,7 +784,7 @@ public void run() { while (!this.isStopped()) { try { NettyEvent event = this.eventQueue.poll(3000, TimeUnit.MILLISECONDS); - if (event != null && listener != null) { + if (event != null && event.getType() != null && listener != null) { switch (event.getType()) { case IDLE: listener.onChannelIdle(event.getRemoteAddr(), event.getChannel()); @@ -814,6 +814,14 @@ public void run() { log.info(this.getServiceName() + " service end"); } + @Override + public void shutdown(final boolean interrupt) { + // wake up the thread blocked on eventQueue.poll() so it observes the stopped flag + // immediately instead of waiting for the next poll timeout + this.eventQueue.offer(new NettyEvent(null, null, null)); + super.shutdown(interrupt); + } + @Override public String getServiceName() { return NettyEventExecutor.class.getSimpleName(); From 8e3d454941cb1067738964bf190f2aab52be48e6 Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:06:28 +0800 Subject: [PATCH 3/8] [ISSUE #10904] Guard malformed retry counters in proxy send path --- .../proxy/processor/ProducerProcessor.java | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java index 8c4907c588a..bf2e95717bc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java @@ -211,13 +211,23 @@ protected SendMessageRequestHeader buildSendMessageRequestHeader(List m if (requestHeader.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) { String reconsumeTimes = MessageAccessor.getReconsumeTime(message); if (reconsumeTimes != null) { - requestHeader.setReconsumeTimes(Integer.valueOf(reconsumeTimes)); + try { + requestHeader.setReconsumeTimes(Integer.parseInt(reconsumeTimes.trim())); + } catch (NumberFormatException e) { + log.warn("parse reconsumeTimes error, with value:{}", reconsumeTimes); + requestHeader.setReconsumeTimes(0); + } MessageAccessor.clearProperty(message, MessageConst.PROPERTY_RECONSUME_TIME); } String maxReconsumeTimes = MessageAccessor.getMaxReconsumeTimes(message); if (maxReconsumeTimes != null) { - requestHeader.setMaxReconsumeTimes(Integer.valueOf(maxReconsumeTimes)); + try { + requestHeader.setMaxReconsumeTimes(Integer.parseInt(maxReconsumeTimes.trim())); + } catch (NumberFormatException e) { + log.warn("parse maxReconsumeTimes error, with value:{}", maxReconsumeTimes); + requestHeader.setMaxReconsumeTimes(0); + } MessageAccessor.clearProperty(message, MessageConst.PROPERTY_MAX_RECONSUME_TIMES); } } From 2496d51c73e2de85b40df69ccbfd7bd9c54bf533 Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:07:02 +0800 Subject: [PATCH 4/8] [ISSUE #10901] Fix OTLP header parsing and avoid leaking header config --- .../rocketmq/proxy/metrics/ProxyMetricsManager.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java index 81db576e3d2..f617c1de0cc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/metrics/ProxyMetricsManager.java @@ -148,9 +148,9 @@ public void start() throws Exception { if (StringUtils.isNotBlank(labels)) { List kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(labels); for (String item : kvPairs) { - String[] split = item.split(":"); + String[] split = item.split(":", 2); if (split.length != 2) { - log.warn("metricsLabel is not valid: {}", labels); + log.warn("metricsLabel is not valid: {}", item); continue; } LABEL_MAP.put(split[0], split[1]); @@ -188,9 +188,9 @@ public void start() throws Exception { Map headerMap = new HashMap<>(); List kvPairs = Splitter.on(',').omitEmptyStrings().splitToList(headers); for (String item : kvPairs) { - String[] split = item.split(":"); + String[] split = item.split(":", 2); if (split.length != 2) { - log.warn("metricsGrpcExporterHeader is not valid: {}", headers); + log.warn("metricsGrpcExporterHeader is not valid: {}", item); continue; } headerMap.put(split[0], split[1]); From ffbb2833198e02ea80e98d975e620b494548cdea Mon Sep 17 00:00:00 2001 From: 123123213weqw <1939455790@qq.com> Date: Thu, 13 Aug 2026 00:07:59 +0800 Subject: [PATCH 5/8] [ISSUE #10895] Validate every Lite consumer subscription against its bound topic --- .../apache/rocketmq/proxy/processor/ClientProcessor.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java index c73e66416da..f0a1b36c79d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ClientProcessor.java @@ -197,8 +197,10 @@ protected void validateLiteSubTopic(ProxyContext ctx, String group, Set Date: Thu, 13 Aug 2026 00:09:25 +0800 Subject: [PATCH 6/8] [ISSUE #10894] Enforce Lite subscription quota for complete and batch updates --- .../broker/lite/LiteSubscriptionRegistryImpl.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java index b487b8757f4..119c5952b0c 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/lite/LiteSubscriptionRegistryImpl.java @@ -97,6 +97,9 @@ public void addPartialSubscription(String clientId, String group, String topic, if (!liteLifecycleManager.isSubscriptionActive(topic, lmqName)) { continue; } + if (!isLiteTopicSubscribed(clientGroup, lmqName) && getActiveSubscriptionNum() >= maxCount) { + throw new LiteQuotaException("lite subscription quota exceeded " + maxCount); + } thisSub.addLiteTopic(lmqName); // First remove the old subscription if (LiteMetadataUtil.isSubLiteExclusive(group, brokerController)) { @@ -147,6 +150,10 @@ public void addCompleteSubscription(String clientId, String group, String topic, removeTopicGroup(clientGroup, lmqName, false); }); lmqNameNew.forEach(lmqName -> { + long maxCount = brokerController.getBrokerConfig().getMaxLiteSubscriptionCount(); + if (!isLiteTopicSubscribed(clientGroup, lmqName) && getActiveSubscriptionNum() >= maxCount) { + throw new LiteQuotaException("lite subscription quota exceeded " + maxCount); + } thisSub.addLiteTopic(lmqName); addTopicGroup(clientGroup, lmqName); }); @@ -269,6 +276,11 @@ public void cleanSubscription(String lmqName, boolean notifyClient) { } } + protected boolean isLiteTopicSubscribed(ClientGroup clientGroup, String lmqName) { + Set topicGroupSet = liteTopic2Group.get(lmqName); + return topicGroupSet != null && topicGroupSet.contains(clientGroup); + } + protected void addTopicGroup(ClientGroup clientGroup, String lmqName) { Set topicGroupSet = liteTopic2Group .computeIfAbsent(lmqName, k -> ConcurrentHashMap.newKeySet()); From aac534b48e496862b0e377531f120830cccf3ba2 Mon Sep 17 00:00:00 2001 From: btlqql <2977859784@qq.com> Date: Thu, 13 Aug 2026 00:10:46 +0800 Subject: [PATCH 7/8] [ISSUE #10873] Guard BrokerData.selectBrokerAddr against empty address table --- .../apache/rocketmq/remoting/protocol/route/BrokerData.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/remoting/src/main/java/org/apache/rocketmq/remoting/protocol/route/BrokerData.java b/remoting/src/main/java/org/apache/rocketmq/remoting/protocol/route/BrokerData.java index de911d17c26..766cae3e421 100644 --- a/remoting/src/main/java/org/apache/rocketmq/remoting/protocol/route/BrokerData.java +++ b/remoting/src/main/java/org/apache/rocketmq/remoting/protocol/route/BrokerData.java @@ -89,6 +89,9 @@ public BrokerData(String cluster, String brokerName, HashMap broke * @return Broker address. */ public String selectBrokerAddr() { + if (this.brokerAddrs == null || this.brokerAddrs.isEmpty()) { + return null; + } String masterAddress = this.brokerAddrs.get(MixAll.MASTER_ID); if (masterAddress == null) { From c311521121372665b8711880df6c3d3cd796000f Mon Sep 17 00:00:00 2001 From: btlqql <2977859784@qq.com> Date: Thu, 13 Aug 2026 00:11:24 +0800 Subject: [PATCH 8/8] [ISSUE #10876] Fix BrokerIdentityInfo.equals for partial identities --- .../controller/impl/heartbeat/BrokerIdentityInfo.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/heartbeat/BrokerIdentityInfo.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/heartbeat/BrokerIdentityInfo.java index 8fc04957e68..a0c0dccc724 100644 --- a/controller/src/main/java/org/apache/rocketmq/controller/impl/heartbeat/BrokerIdentityInfo.java +++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/heartbeat/BrokerIdentityInfo.java @@ -63,7 +63,9 @@ public boolean equals(Object obj) { if (obj instanceof BrokerIdentityInfo) { BrokerIdentityInfo addr = (BrokerIdentityInfo) obj; - return clusterName.equals(addr.clusterName) && brokerName.equals(addr.brokerName) && brokerId.equals(addr.brokerId); + return Objects.equals(clusterName, addr.clusterName) + && Objects.equals(brokerName, addr.brokerName) + && Objects.equals(brokerId, addr.brokerId); } return false; }