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; } }