From 549aa9acb056c0ffad08dcc96d31fa027fbc0057 Mon Sep 17 00:00:00 2001 From: root Date: Sun, 30 Aug 2026 01:33:08 +0800 Subject: [PATCH] [ISSUE #10988] Fix concurrent commitOffset race losing queue offsets for a new topic@group --- .../broker/offset/ConsumerOffsetManager.java | 17 ++++----- .../offset/ConsumerOffsetManagerTest.java | 36 +++++++++++++++++++ 2 files changed, 43 insertions(+), 10 deletions(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java b/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java index 1d3bf7bed09..71946701697 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java @@ -203,16 +203,13 @@ public void commitOffset(final String clientHost, final String group, final Stri } private void commitOffset(final String clientHost, final String key, final int queueId, final long offset) { - ConcurrentMap map = this.offsetTable.get(key); - if (null == map) { - map = new ConcurrentHashMap<>(2); - map.put(queueId, offset); - this.offsetTable.put(key, map); - } else { - Long storeOffset = map.put(queueId, offset); - if (storeOffset != null && offset < storeOffset) { - LOG.warn("[NOTIFYME]update consumer offset less than store. clientHost={}, key={}, queueId={}, requestOffset={}, storeOffset={}", clientHost, key, queueId, offset, storeOffset); - } + // computeIfAbsent, so concurrent commits for different queues of a new topic@group + // all land in the same map instead of overwriting each other's entries. + ConcurrentMap map = this.offsetTable.computeIfAbsent(key, + k -> new ConcurrentHashMap<>(2)); + Long storeOffset = map.put(queueId, offset); + if (storeOffset != null && offset < storeOffset) { + LOG.warn("[NOTIFYME]update consumer offset less than store. clientHost={}, key={}, queueId={}, requestOffset={}, storeOffset={}", clientHost, key, queueId, offset, storeOffset); } if (versionChangeCounter.incrementAndGet() % brokerController.getBrokerConfig().getConsumerOffsetUpdateVersionStep() == 0) { updateDataVersion(); diff --git a/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java b/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java index 7e4faa4e42f..29ce15a75ec 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManagerTest.java @@ -24,8 +24,13 @@ import org.junit.Before; import org.junit.Test; +import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import org.mockito.Mockito; import static org.apache.rocketmq.broker.offset.ConsumerOffsetManager.TOPIC_GROUP_SEPARATOR; @@ -121,4 +126,35 @@ public void testEraseResetOffset() { Assert.assertFalse(consumerOffsetManager.hasOffsetReset(topic, group, 1)); Assert.assertFalse(consumerOffsetManager.resetOffsetTable.containsKey(key)); } + + @Test + public void testConcurrentCommitOffsetDoesNotLoseQueues() throws InterruptedException { + Mockito.when(brokerController.getBrokerConfig()).thenReturn(new BrokerConfig()); + + final int queueCount = 8; + final int keyCount = 500; + ExecutorService executor = Executors.newFixedThreadPool(queueCount); + CountDownLatch latch = new CountDownLatch(queueCount); + for (int queueId = 0; queueId < queueCount; queueId++) { + final int qid = queueId; + executor.submit(() -> { + try { + // every thread commits a distinct queue of the same brand-new keys, + // so creation of each topic@group entry races + for (int i = 0; i < keyCount; i++) { + consumerOffsetManager.commitOffset("host", "group", "topic" + i, qid, i); + } + } finally { + latch.countDown(); + } + }); + } + Assert.assertTrue(latch.await(60, TimeUnit.SECONDS)); + executor.shutdownNow(); + + for (int i = 0; i < keyCount; i++) { + Map offsets = consumerOffsetManager.queryOffset("group", "topic" + i); + assertThat(offsets).as("offsets of topic" + i).hasSize(queueCount); + } + } }