From e2f212aa6d658509a67db0601093a79eb098858a Mon Sep 17 00:00:00 2001 From: "terrance.lzm" Date: Wed, 20 Aug 2025 17:43:45 +0800 Subject: [PATCH 1/2] [ISSUE #9626] Prevent premature offset commit before consumer record flush --- .../rocketmq/broker/pop/PopConsumerCache.java | 106 ++++++++---------- .../broker/pop/PopConsumerCacheTest.java | 10 +- 2 files changed, 53 insertions(+), 63 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java index e7ce68e0193..43c4b0a8b83 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java @@ -20,13 +20,11 @@ import java.util.Iterator; import java.util.List; import java.util.Map; -import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.locks.Lock; -import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.offset.ConsumerOffsetManager; @@ -128,30 +126,29 @@ public int cleanupRecords(Consumer consumer) { records.getGroupId(), records.getTopicId()); if (timeout) { - List removeExpiredRecords = - records.removeExpiredRecords(Long.MAX_VALUE); - if (removeExpiredRecords != null) { - consumerRecordStore.writeRecords(removeExpiredRecords); + records.stageExpiredRecords(Long.MAX_VALUE); + List writeConsumerRecords = + new ArrayList<>(records.getRemoveTreeMap().values()); + if (!writeConsumerRecords.isEmpty()) { + consumerRecordStore.writeRecords(writeConsumerRecords); } + records.clearStagedRecords(); log.info("PopConsumerOffline, so clean expire records, groupId={}, topic={}, queueId={}, records={}", - records.getGroupId(), records.getTopicId(), records.getQueueId(), - removeExpiredRecords != null ? removeExpiredRecords.size() : 0); + records.getGroupId(), records.getTopicId(), records.getQueueId(), records.getInFlightRecordCount()); iterator.remove(); continue; } long currentTime = System.currentTimeMillis(); + records.stageExpiredRecords(currentTime); List writeConsumerRecords = new ArrayList<>(); - List consumerRecords = records.removeExpiredRecords(currentTime); - if (consumerRecords != null) { - consumerRecords.forEach(consumerRecord -> { - if (consumerRecord.getVisibilityTimeout() <= currentTime) { - consumer.accept(consumerRecord); - } else { - writeConsumerRecords.add(consumerRecord); - } - }); - } + records.getRemoveTreeMap().values().forEach(record -> { + if (record.getVisibilityTimeout() <= currentTime) { + consumer.accept(record); + } else { + writeConsumerRecords.add(record); + } + }); // write to store and handle it later consumerRecordStore.writeRecords(writeConsumerRecords); @@ -209,72 +206,64 @@ public void run() { protected static class ConsumerRecords { - private final Lock lock; private final String groupId; private final String topicId; private final int queueId; private final BrokerConfig brokerConfig; - private final TreeMap recordTreeMap; + private final ConcurrentSkipListMap removeTreeMap; + private final ConcurrentSkipListMap recordTreeMap; public ConsumerRecords(BrokerConfig brokerConfig, String groupId, String topicId, int queueId) { this.groupId = groupId; this.topicId = topicId; this.queueId = queueId; - this.lock = new ReentrantLock(); this.brokerConfig = brokerConfig; - this.recordTreeMap = new TreeMap<>(); + this.removeTreeMap = new ConcurrentSkipListMap<>(); + this.recordTreeMap = new ConcurrentSkipListMap<>(); } public void write(PopConsumerRecord record) { - lock.lock(); - try { - recordTreeMap.put(record.getOffset(), record); - } finally { - lock.unlock(); - } + recordTreeMap.put(record.getOffset(), record); } public boolean delete(PopConsumerRecord record) { - PopConsumerRecord popConsumerRecord; - lock.lock(); - try { - popConsumerRecord = recordTreeMap.remove(record.getOffset()); - } finally { - lock.unlock(); - } - return popConsumerRecord != null; + return recordTreeMap.remove(record.getOffset()) != null; } public long getMinOffsetInBuffer() { - Map.Entry entry = recordTreeMap.firstEntry(); + Map.Entry entry = removeTreeMap.firstEntry(); + if (entry != null) { + return entry.getKey(); + } + entry = recordTreeMap.firstEntry(); return entry != null ? entry.getKey() : OFFSET_NOT_EXIST; } public int getInFlightRecordCount() { - return recordTreeMap.size(); + return removeTreeMap.size() + recordTreeMap.size(); } - public List removeExpiredRecords(long currentTime) { - List result = null; - lock.lock(); - try { - Iterator> iterator = recordTreeMap.entrySet().iterator(); - while (iterator.hasNext()) { - Map.Entry entry = iterator.next(); - // org.apache.rocketmq.broker.processor.PopBufferMergeService.scan - if (entry.getValue().getVisibilityTimeout() <= currentTime || - entry.getValue().getPopTime() + brokerConfig.getPopCkStayBufferTime() <= currentTime) { - if (result == null) { - result = new ArrayList<>(); - } - result.add(entry.getValue()); - iterator.remove(); - } + public void stageExpiredRecords(long currentTime) { + Iterator> + iterator = recordTreeMap.entrySet().iterator(); + + // refer: org.apache.rocketmq.broker.processor.PopBufferMergeService.scan + while (iterator.hasNext()) { + Map.Entry entry = iterator.next(); + if (entry.getValue().getVisibilityTimeout() <= currentTime || + entry.getValue().getPopTime() + brokerConfig.getPopCkStayBufferTime() <= currentTime) { + removeTreeMap.put(entry.getKey(), entry.getValue()); + iterator.remove(); } - } finally { - lock.unlock(); } - return result; + } + + public void clearStagedRecords() { + removeTreeMap.clear(); + } + + public ConcurrentSkipListMap getRemoveTreeMap() { + return removeTreeMap; } public String getGroupId() { @@ -292,7 +281,6 @@ public int getQueueId() { @Override public String toString() { return "ConsumerRecords{" + - "lock=" + lock + ", topicId=" + topicId + ", groupId=" + groupId + ", queueId=" + queueId + diff --git a/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerCacheTest.java b/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerCacheTest.java index 3f6e893a527..28045ca26e7 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerCacheTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerCacheTest.java @@ -61,10 +61,12 @@ public void consumerRecordsTest() { Assert.assertEquals(3, consumerRecords.getInFlightRecordCount()); long bufferTimeout = brokerConfig.getPopCkStayBufferTime(); - Assert.assertEquals(1, consumerRecords.removeExpiredRecords(bufferTimeout + 2).size()); - Assert.assertNull(consumerRecords.removeExpiredRecords(bufferTimeout + 2)); - Assert.assertEquals(2, consumerRecords.removeExpiredRecords(bufferTimeout + 4).size()); - Assert.assertNull(consumerRecords.removeExpiredRecords(bufferTimeout + 4)); + consumerRecords.stageExpiredRecords(bufferTimeout + 2); + Assert.assertEquals(1, consumerRecords.getRemoveTreeMap().size()); + consumerRecords.clearStagedRecords(); + consumerRecords.stageExpiredRecords(bufferTimeout + 4); + Assert.assertEquals(2, consumerRecords.getRemoveTreeMap().size()); + consumerRecords.clearStagedRecords(); } @Test From e819344465ef6717eb52255030cd3510dc7eb77d Mon Sep 17 00:00:00 2001 From: "terrance.lzm" Date: Thu, 21 Aug 2025 14:41:38 +0800 Subject: [PATCH 2/2] [ISSUE #9626] Prevent premature offset commit before consumer record flush --- .../java/org/apache/rocketmq/broker/pop/PopConsumerCache.java | 1 + 1 file changed, 1 insertion(+) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java index 43c4b0a8b83..7f518171676 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerCache.java @@ -152,6 +152,7 @@ public int cleanupRecords(Consumer consumer) { // write to store and handle it later consumerRecordStore.writeRecords(writeConsumerRecords); + records.clearStagedRecords(); // commit min offset in buffer to offset store long offset = records.getMinOffsetInBuffer();