Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<Integer, Long> 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<Integer, Long> 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Integer, Long> offsets = consumerOffsetManager.queryOffset("group", "topic" + i);
assertThat(offsets).as("offsets of topic" + i).hasSize(queueCount);
}
}
}