From e30b59b6c54120dc5f1a86b7bf85f3424a3309c2 Mon Sep 17 00:00:00 2001 From: yyqdbngt <300715189+yyqdbngt@users.noreply.github.com> Date: Fri, 31 Jul 2026 23:16:07 +0800 Subject: [PATCH] fix: infinite loop in getBulkData, CME in PopReviveService, and NPE in Broker2Client - CommitLog.getBulkData: when findMappedFileByOffset returns null the loop body does nothing, so remainSize/startOffset never advance and the thread spins forever; break out with a warning log instead - PopReviveService: removed entries from the fail-fast TreeMap-backed map via map.remove() while iterating, throwing ConcurrentModificationException and aborting offset commits; use the iterator's remove() - Broker2Client.getConsumeStatus: getConsumerGroupInfo(group) is null when the group is not online, dereferencing it throws NPE; guard it Compiled and verified on the build server (mvn -pl broker -am compile). --- .../apache/rocketmq/broker/client/net/Broker2Client.java | 4 ++-- .../rocketmq/broker/processor/PopReviveService.java | 9 +++++++-- .../main/java/org/apache/rocketmq/store/CommitLog.java | 3 +++ 3 files changed, 12 insertions(+), 4 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/client/net/Broker2Client.java b/broker/src/main/java/org/apache/rocketmq/broker/client/net/Broker2Client.java index 5a6c4c94c47..f1b73bdccfb 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/client/net/Broker2Client.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/client/net/Broker2Client.java @@ -253,8 +253,8 @@ public RemotingCommand getConsumeStatus(String topic, String group, String origi requestHeader); Map> consumerStatusTable = new HashMap<>(); - ConcurrentMap channelInfoTable = - this.brokerController.getConsumerManager().getConsumerGroupInfo(group).getChannelInfoTable(); + ConsumerGroupInfo consumerGroupInfo = this.brokerController.getConsumerManager().getConsumerGroupInfo(group); + ConcurrentMap channelInfoTable = consumerGroupInfo == null ? null : consumerGroupInfo.getChannelInfoTable(); if (null == channelInfoTable || channelInfoTable.isEmpty()) { result.setCode(ResponseCode.SYSTEM_ERROR); result.setRemark(String.format("No Any Consumer online in the consumer group: [%s]", group)); diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java index 07f16e98965..238c663fab6 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopReviveService.java @@ -54,6 +54,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.NavigableMap; @@ -592,12 +593,16 @@ private void reviveMsgFromCk(PopCheckPoint popCheckPoint) { if (inflightReviveRequestMap.containsKey(popCheckPoint)) { inflightReviveRequestMap.get(popCheckPoint).setObject2(true); } - for (Map.Entry> entry : inflightReviveRequestMap.entrySet()) { + // Use the iterator to remove entries: directly calling map.remove() while + // iterating over the fail-fast TreeMap iterator throws ConcurrentModificationException. + Iterator>> iterator = inflightReviveRequestMap.entrySet().iterator(); + while (iterator.hasNext()) { + Map.Entry> entry = iterator.next(); PopCheckPoint oldCK = entry.getKey(); Pair pair = entry.getValue(); if (pair.getObject2()) { brokerController.getConsumerOffsetManager().commitOffset(PopAckConstants.LOCAL_HOST, PopAckConstants.REVIVE_GROUP, reviveTopic, queueId, oldCK.getReviveOffset()); - inflightReviveRequestMap.remove(oldCK); + iterator.remove(); } else { break; } diff --git a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java index 260f021aed7..1b068c11757 100644 --- a/store/src/main/java/org/apache/rocketmq/store/CommitLog.java +++ b/store/src/main/java/org/apache/rocketmq/store/CommitLog.java @@ -311,6 +311,9 @@ public List getBulkData(final long offset, final int s bufferResultList.add(bufferResult); remainSize -= readSize; startOffset += readSize; + } else { + log.warn("getBulkData: can not find mapped file by offset, break to avoid infinite loop. offset: {}, size: {}", startOffset, remainSize); + break; } }