From b5f76bb01365a9ce5895297c726fb37d78022715 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 9 Jun 2021 00:48:58 +0800 Subject: [PATCH 1/2] Fix deadlock caused by synchronous asyncFindPosition --- .../kop/KafkaTopicConsumerManager.java | 155 +++++++-------- .../handlers/kop/MessageFetchContext.java | 184 +++++++++--------- .../handlers/kop/utils/MessageIdUtils.java | 13 -- .../kop/KafkaTopicConsumerManagerTest.java | 26 +-- 4 files changed, 182 insertions(+), 196 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java index 60672b8fbe..764d03760a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java @@ -15,9 +15,10 @@ import static com.google.common.base.Preconditions.checkArgument; -import io.streamnative.pulsar.handlers.kop.utils.MessageIdUtils; +import io.streamnative.pulsar.handlers.kop.utils.OffsetSearchPredicate; import java.io.Closeable; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.locks.ReentrantReadWriteLock; @@ -25,6 +26,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCursorCallback; import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; @@ -47,9 +49,10 @@ public class KafkaTopicConsumerManager implements Closeable { private final ReentrantReadWriteLock rwLock; private boolean closed; - // keep fetch offset and related cursor. keep cursor and its last offset in Pair. + // key is the offset, value is the future of (cursor, offset), whose offset is the last offset in pair. @Getter - private final ConcurrentLongHashMap> consumers; + private final ConcurrentLongHashMap>> cursors; + // used to track all created cursor, since above consumers may be remove and in fly, // use this map will not leak cursor when close. private final ConcurrentMap createdCursors; @@ -60,7 +63,7 @@ public class KafkaTopicConsumerManager implements Closeable { KafkaTopicConsumerManager(KafkaRequestHandler requestHandler, PersistentTopic topic) { this.topic = topic; - this.consumers = new ConcurrentLongHashMap<>(); + this.cursors = new ConcurrentLongHashMap<>(); this.createdCursors = new ConcurrentHashMap<>(); this.lastAccessTimes = new ConcurrentLongHashMap<>(); this.requestHandler = requestHandler; @@ -78,7 +81,7 @@ void deleteExpiredCursor(long current, long expirePeriodMillis) { } void deleteOneExpiredCursor(long offset) { - Pair pair; + CompletableFuture> cursorFuture; // need not do anything, since this tcm already in closing state. and close() will delete every thing. rwLock.readLock().lock(); @@ -86,20 +89,26 @@ void deleteOneExpiredCursor(long offset) { if (closed) { return; } - pair = consumers.remove(offset); + cursorFuture = cursors.remove(offset); lastAccessTimes.remove(offset); } finally { rwLock.readLock().unlock(); } - if (pair != null) { + if (cursorFuture != null) { if (log.isDebugEnabled()) { log.debug("[{}] Cursor timed out for offset: {}, cursors cache size: {}", - requestHandler.ctx.channel(), offset, consumers.size()); + requestHandler.ctx.channel(), offset, cursors.size()); } - ManagedCursor managedCursor = pair.getKey(); - deleteOneCursorAsync(managedCursor, "cursor expired"); + // TODO: Should we just cancel this future? + cursorFuture.whenComplete((pair, e) -> { + if (e != null || pair == null) { + return; + } + ManagedCursor managedCursor = pair.getKey(); + deleteOneCursorAsync(managedCursor, "cursor expired"); + }); } } @@ -128,67 +137,49 @@ public void deleteCursorFailed(ManagedLedgerException exception, Object ctx) { // get one cursor offset pair. // remove from cache, so another same offset read could happen. // each success remove should have a following add. - public Pair remove(long offset) { - Pair cursor; - - // should not return cursor for Fetch to read, since this tcm already in closing state. + public CompletableFuture> removeCursorFuture(long offset) { rwLock.readLock().lock(); try { if (closed) { return null; } - cursor = consumers.remove(offset); lastAccessTimes.remove(offset); - } finally { - rwLock.readLock().unlock(); - } + final CompletableFuture> cursorFuture = cursors.remove(offset); + if (cursorFuture == null) { + return asyncCreateCursorIfNotExists(offset); + } - if (cursor != null) { if (log.isDebugEnabled()) { log.debug("[{}] Get cursor for offset: {} in cache. cache size: {}", - requestHandler.ctx.channel(), offset, consumers.size()); + requestHandler.ctx.channel(), offset, cursors.size()); } - return cursor; + return cursorFuture; + } finally { + rwLock.readLock().unlock(); } - - return createCursorIfNotExists(offset); } - private Pair createCursorIfNotExists(long offset) { - - Pair cursor; - + private CompletableFuture> asyncCreateCursorIfNotExists(long offset) { rwLock.readLock().lock(); try { if (closed) { return null; } - // handle offset not exist in consumers, need create cursor. - Pair managedCursorLongPair = getCursorByOffset(offset); - if (null == managedCursorLongPair) { - log.error("[{}] Failed to get cursor by offset {}", requestHandler.ctx.channel(), offset); - return null; - } - - consumers.putIfAbsent(offset, managedCursorLongPair); + cursors.putIfAbsent(offset, asyncGetCursorByOffset(offset)); // notice: above would add a - cursor = consumers.remove(offset); lastAccessTimes.remove(offset); + return cursors.remove(offset); } finally { rwLock.readLock().unlock(); } - - return cursor; } - // once entry read complete, add new offset back. public void add(long offset, Pair pair) { checkArgument(offset == pair.getRight(), - "offset not equal. key: " + offset + " value: " + pair.getRight()); + "offset not equal. key: " + offset + " value: " + pair.getRight()); rwLock.readLock().lock(); - // should delete the cursor since this tcm already in closing state. try { if (closed) { ManagedCursor managedCursor = pair.getLeft(); @@ -199,22 +190,23 @@ public void add(long offset, Pair pair) { rwLock.readLock().unlock(); } - Pair oldPair = consumers.putIfAbsent(offset, pair); - if (oldPair != null) { + final CompletableFuture> cursorFuture = CompletableFuture.completedFuture(pair); + if (cursors.putIfAbsent(offset, cursorFuture) != null) { deleteOneCursorAsync(pair.getLeft(), "reason: A race - same cursor already cached"); } lastAccessTimes.put(offset, System.currentTimeMillis()); if (log.isDebugEnabled()) { log.debug("[{}] Add cursor back {} for offset: {}", - requestHandler.ctx.channel(), pair.getLeft().getName(), offset); + requestHandler.ctx.channel(), pair.getLeft().getName(), offset); } } // called when channel closed. @Override public void close() { - ConcurrentLongHashMap> consumersToClose; + final ConcurrentLongHashMap>> cursorFuturesToClose = + new ConcurrentLongHashMap<>(); ConcurrentMap cursorsToClose; rwLock.writeLock().lock(); try { @@ -226,25 +218,28 @@ public void close() { log.debug("[{}] Close TCM for topic {}.", requestHandler.ctx.channel(), topic.getName()); } - consumersToClose = new ConcurrentLongHashMap<>(); - consumers.forEach((k, v) -> consumersToClose.put(k, v)); - consumers.clear(); + cursors.forEach(cursorFuturesToClose::put); + cursors.clear(); lastAccessTimes.clear(); cursorsToClose = new ConcurrentHashMap<>(); - createdCursors.forEach((k, v) -> cursorsToClose.put(k, v)); + createdCursors.forEach(cursorsToClose::put); createdCursors.clear(); } finally { rwLock.writeLock().unlock(); } - consumersToClose.values() - .forEach(pair -> { - ManagedCursor cursor = pair.getLeft(); - deleteOneCursorAsync(cursor, "TopicConsumerManager close"); - if (null != cursor) { - cursorsToClose.remove(cursor.getName()); - } + cursorFuturesToClose.values().forEach(cursorFuture -> { + cursorFuture.whenComplete((pair, e) -> { + if (e != null || pair == null) { + return; + } + ManagedCursor cursor = pair.getLeft(); + deleteOneCursorAsync(cursor, "TopicConsumerManager close"); + if (cursor != null) { + cursorsToClose.remove(cursor.getName()); + } }); + }); // delete dangling createdCursors cursorsToClose.values().forEach(cursor -> @@ -252,31 +247,29 @@ public void close() { cursorsToClose.clear(); } - private Pair getCursorByOffset(Long offset) { - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topic.getManagedLedger(); - PositionImpl position = MessageIdUtils.getPositionForOffset(ledger, offset); + private CompletableFuture> asyncGetCursorByOffset(long offset) { + final ManagedLedger ledger = topic.getManagedLedger(); + return ledger.asyncFindPosition(new OffsetSearchPredicate(offset)).thenApply(position -> { + final String cursorName = "kop-consumer-cursor-" + topic.getName() + + "-" + position.getLedgerId() + "-" + position.getEntryId() + + "-" + DigestUtils.sha1Hex(UUID.randomUUID().toString()).substring(0, 10); - String cursorName = "kop-consumer-cursor-" + topic.getName() - + "-" + position.getLedgerId() + "-" + position.getEntryId() - + "-" + DigestUtils.sha1Hex(UUID.randomUUID().toString()).substring(0, 10); - - // get previous position, because NonDurableCursor is read from next position. - PositionImpl previous = ledger.getPreviousPosition(position); - if (log.isDebugEnabled()) { - log.debug("[{}] Create cursor {} for offset: {}. position: {}, previousPosition: {}", - requestHandler.ctx.channel(), cursorName, offset, position, previous); - } - ManagedCursor newCursor; - try { - newCursor = ledger.newNonDurableCursor(previous, cursorName); - createdCursors.put(newCursor.getName(), newCursor); - } catch (ManagedLedgerException e) { - log.error("[{}] Error new cursor for topic {} at offset {} - {}. will cause fetch data error.", - requestHandler.ctx.channel(), topic.getName(), offset, previous, e); - return null; - } - - lastAccessTimes.put(offset, System.currentTimeMillis()); - return Pair.of(newCursor, offset); + // get previous position, because NonDurableCursor is read from next position. + final PositionImpl previous = ((ManagedLedgerImpl) ledger).getPreviousPosition((PositionImpl) position); + if (log.isDebugEnabled()) { + log.debug("[{}] Create cursor {} for offset: {}. position: {}, previousPosition: {}", + requestHandler.ctx.channel(), cursorName, offset, position, previous); + } + try { + final ManagedCursor newCursor = ledger.newNonDurableCursor(previous, cursorName); + createdCursors.putIfAbsent(newCursor.getName(), newCursor); + lastAccessTimes.put(offset, System.currentTimeMillis()); + return Pair.of(newCursor, offset); + } catch (ManagedLedgerException e) { + log.error("[{}] Error new cursor for topic {} at offset {} - {}. will cause fetch data error.", + requestHandler.ctx.channel(), topic.getName(), offset, previous, e); + return null; + } + }); } } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java index 618691386f..445fc834b6 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java @@ -222,109 +222,115 @@ public void handleFetch() { topicPartition, offset); } - final Pair cursorLongPair = tcm.remove(offset); - if (cursorLongPair == null) { - log.warn("KafkaTopicConsumerManager.remove({}) return null for topic {}. " - + "Fetch for topic return error.", - offset, topicPartition); - addErrorPartitionResponse(topicPartition, Errors.NOT_LEADER_FOR_PARTITION); + final CompletableFuture> cursorFuture = tcm.removeCursorFuture(offset); + if (cursorFuture == null) { + // tcm is closed, just return a NONE error because the channel may be still active + log.warn("[{}] KafkaTopicConsumerManager is closed", requestHandler.ctx); + addErrorPartitionResponse(topicPartition, Errors.NONE); return; } - final ManagedCursor cursor = cursorLongPair.getLeft(); - final AtomicLong cursorOffset = new AtomicLong(cursorLongPair.getRight()); - final long highWatermark = MessageIdUtils.getHighWatermark( - cursorLongPair.getLeft().getManagedLedger()); - statsLogger.getPrepareMetadataStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(startPrepareMetadataNanos), TimeUnit.NANOSECONDS); - readEntries(cursor, topicPartition, cursorOffset).whenComplete((entries, throwable) -> { - if (throwable != null) { - tcm.deleteOneCursorAsync(cursorLongPair.getLeft(), "cursor.readEntry fail. deleteCursor"); - addErrorPartitionResponse(topicPartition, Errors.forException(throwable)); - return; - } - if (entries == null) { - addErrorPartitionResponse(topicPartition, - Errors.forException(new ApiException("Cursor is null"))); + // cursorFuture is never completed exceptionally because ManagedLedgerImpl#asyncFindPosition is never + // completed exceptionally. + cursorFuture.thenAccept(cursorLongPair -> { + if (cursorLongPair == null) { + log.warn("KafkaTopicConsumerManager.remove({}) return null for topic {}. " + + "Fetch for topic return error.", + offset, topicPartition); + addErrorPartitionResponse(topicPartition, Errors.NOT_LEADER_FOR_PARTITION); return; } - // Add new offset back to TCM after entries are read successfully - tcm.add(cursorOffset.get(), Pair.of(cursor, cursorOffset.get())); - - final long lso = (readCommitted - ? tc.getLastStableOffset(TopicName.get(fullTopicName), highWatermark) - : highWatermark); - List committedEntries = entries; - if (readCommitted) { - committedEntries = new ArrayList<>(); - for (Entry entry : entries) { - if (lso >= MessageIdUtils.peekBaseOffsetFromEntry(entry)) { - committedEntries.add(entry); - } else { - break; - } + final ManagedCursor cursor = cursorLongPair.getLeft(); + final AtomicLong cursorOffset = new AtomicLong(cursorLongPair.getRight()); + final long highWatermark = MessageIdUtils.getHighWatermark( + cursorLongPair.getLeft().getManagedLedger()); + statsLogger.getPrepareMetadataStats().registerSuccessfulEvent( + MathUtils.elapsedNanos(startPrepareMetadataNanos), TimeUnit.NANOSECONDS); + readEntries(cursor, topicPartition, cursorOffset).whenComplete((entries, throwable) -> { + if (throwable != null) { + tcm.deleteOneCursorAsync(cursorLongPair.getLeft(), "cursor.readEntry fail. deleteCursor"); + addErrorPartitionResponse(topicPartition, Errors.forException(throwable)); + return; } - if (log.isDebugEnabled()) { - log.debug("Request {}: read {} entries but only {} entries are committed", - header, entries.size(), committedEntries.size()); + if (entries == null) { + addErrorPartitionResponse(topicPartition, + Errors.forException(new ApiException("Cursor is null"))); + return; } - } else { - if (log.isDebugEnabled()) { - log.debug("Request {}: read {} entries", header, entries.size()); + + // Add new offset back to TCM after entries are read successfully + tcm.add(cursorOffset.get(), Pair.of(cursor, cursorOffset.get())); + + final long lso = (readCommitted + ? tc.getLastStableOffset(TopicName.get(fullTopicName), highWatermark) + : highWatermark); + List committedEntries = entries; + if (readCommitted) { + committedEntries = new ArrayList<>(); + for (Entry entry : entries) { + if (lso >= MessageIdUtils.peekBaseOffsetFromEntry(entry)) { + committedEntries.add(entry); + } else { + break; + } + } + if (log.isDebugEnabled()) { + log.debug("Request {}: read {} entries but only {} entries are committed", + header, entries.size(), committedEntries.size()); + } + } else { + if (log.isDebugEnabled()) { + log.debug("Request {}: read {} entries", header, entries.size()); + } + } + if (committedEntries.isEmpty()) { + addErrorPartitionResponse(topicPartition, Errors.NONE); + return; } - } - if (committedEntries.isEmpty()) { - addErrorPartitionResponse(topicPartition, Errors.NONE); - return; - } - // use compatible magic value by apiVersion - short apiVersion = header.apiVersion(); - byte magic = RecordBatch.CURRENT_MAGIC_VALUE; - if (apiVersion <= 1) { - magic = RecordBatch.MAGIC_VALUE_V0; - } else if (apiVersion <= 3) { - magic = RecordBatch.MAGIC_VALUE_V1; - } + // use compatible magic value by apiVersion + short apiVersion = header.apiVersion(); + byte magic = RecordBatch.CURRENT_MAGIC_VALUE; + if (apiVersion <= 1) { + magic = RecordBatch.MAGIC_VALUE_V0; + } else if (apiVersion <= 3) { + magic = RecordBatch.MAGIC_VALUE_V1; + } - // get group and consumer - final String groupName = requestHandler - .getCurrentConnectedGroup().computeIfAbsent(clientHost, ignored -> { - String zkSubPath = ZooKeeperUtils.groupIdPathFormat(clientHost, - header.clientId()); - String groupId = ZooKeeperUtils.getData(requestHandler.getPulsarService().getZkClient(), - requestHandler.getGroupIdStoredPath(), zkSubPath); - if (groupId.isEmpty()) { - log.error("get empty group name from zk for current connection: {}, consumer stats" - + "won't be updated", clientHost); - } else { - log.info("get group name from zk for current connection: {} groupId: {}", + // get group and consumer + final String groupName = requestHandler + .getCurrentConnectedGroup().computeIfAbsent(clientHost, ignored -> { + String zkSubPath = ZooKeeperUtils.groupIdPathFormat(clientHost, + header.clientId()); + String groupId = ZooKeeperUtils.getData( + requestHandler.getPulsarService().getZkClient(), + requestHandler.getGroupIdStoredPath(), + zkSubPath); + log.info("get group name from zk for current connection:{} groupId:{}", clientHost, groupId); - } - return groupId; - }); - final long startDecodingEntriesNanos = MathUtils.nowInNano(); - final DecodeResult decodeResult = requestHandler.getEntryFormatter().decode(entries, magic); - requestHandler.requestStats.getFetchDecodeStats().registerSuccessfulEvent( - MathUtils.elapsedNanos(startDecodingEntriesNanos), TimeUnit.NANOSECONDS); - decodeResults.add(decodeResult); - - // collect consumer metrics - if (!groupName.isEmpty()) { + return groupId; + }); + final long startDecodingEntriesNanos = MathUtils.nowInNano(); + final DecodeResult decodeResult = requestHandler.getEntryFormatter().decode(entries, magic); + requestHandler.requestStats.getFetchDecodeStats().registerSuccessfulEvent( + MathUtils.elapsedNanos(startDecodingEntriesNanos), TimeUnit.NANOSECONDS); + decodeResults.add(decodeResult); + + // collect consumer metrics updateConsumerStats(topicPartition, decodeResult.getRecords(), entries.size(), groupName); - } - final List abortedTransactions = - (readCommitted ? tc.getAbortedIndexList(partitionData.fetchOffset) : null); - responseData.put(topicPartition, new PartitionData<>( - Errors.NONE, - highWatermark, - lso, - highWatermark, // TODO: should it be changed to the logStartOffset? - abortedTransactions, - decodeResult.getRecords())); - tryComplete(); + final List abortedTransactions = + (readCommitted ? tc.getAbortedIndexList(partitionData.fetchOffset) : null); + responseData.put(topicPartition, new PartitionData<>( + Errors.NONE, + highWatermark, + lso, + highWatermark, // TODO: should it be changed to the logStartOffset? + abortedTransactions, + decodeResult.getRecords())); + tryComplete(); + }); }); }); }); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java index 751ecd1557..1e34bd7276 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/utils/MessageIdUtils.java @@ -15,7 +15,6 @@ import io.netty.buffer.ByteBuf; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ExecutionException; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedger; @@ -25,14 +24,11 @@ import org.apache.pulsar.broker.intercept.ManagedLedgerInterceptorImpl; import org.apache.pulsar.common.api.proto.MessageMetadata; import org.apache.pulsar.common.protocol.Commands; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; /** * Utils for Pulsar MessageId. */ public class MessageIdUtils { - private static final Logger log = LoggerFactory.getLogger(MessageIdUtils.class); public static long getCurrentOffset(ManagedLedger managedLedger) { return ((ManagedLedgerInterceptorImpl) managedLedger.getManagedLedgerInterceptor()).getIndex(); @@ -93,15 +89,6 @@ public void readEntryComplete(Entry entry, Object ctx) { return future; } - public static PositionImpl getPositionForOffset(ManagedLedger managedLedger, Long offset) { - try { - return (PositionImpl) managedLedger.asyncFindPosition(new OffsetSearchPredicate(offset)).get(); - } catch (InterruptedException | ExecutionException e) { - log.error("[{}] Failed to find position for offset {}", managedLedger.getName(), offset); - throw new RuntimeException(managedLedger.getName() + " failed to find position for offset " + offset); - } - } - public static long peekOffsetFromEntry(Entry entry) { return Commands.peekBrokerEntryMetadataIfExist(entry.getDataBuffer()).getIndex(); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java index f174d331e0..25b94cca7b 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java @@ -144,8 +144,8 @@ public void testTopicConsumerManagerRemoveAndAdd() throws Exception { KafkaTopicConsumerManager topicConsumerManager = tcm.get(); // before a read, first get cursor of offset. - Pair cursorPair = topicConsumerManager.remove(offset); - assertEquals(topicConsumerManager.getConsumers().size(), 0); + Pair cursorPair = topicConsumerManager.removeCursorFuture(offset).get(); + assertEquals(topicConsumerManager.getCursors().size(), 0); ManagedCursor cursor = cursorPair.getLeft(); assertEquals(cursorPair.getRight(), Long.valueOf(offset)); @@ -156,11 +156,11 @@ public void testTopicConsumerManagerRemoveAndAdd() throws Exception { // simulate a read complete; offset++; topicConsumerManager.add(offset, Pair.of(cursor, offset)); - assertEquals(topicConsumerManager.getConsumers().size(), 1); + assertEquals(topicConsumerManager.getCursors().size(), 1); // another read, cache hit. - cursorPair = topicConsumerManager.remove(offset); - assertEquals(topicConsumerManager.getConsumers().size(), 0); + cursorPair = topicConsumerManager.removeCursorFuture(offset).get(); + assertEquals(topicConsumerManager.getCursors().size(), 0); ManagedCursor cursor2 = cursorPair.getLeft(); assertEquals(cursor2, cursor); @@ -178,9 +178,9 @@ public void testTopicConsumerManagerRemoveAndAdd() throws Exception { } // try read last messages, so read not continuous - cursorPair = topicConsumerManager.remove(offset); + cursorPair = topicConsumerManager.removeCursorFuture(offset).get(); // since above remove will use a new cursor. there should be one in the map. - assertEquals(topicConsumerManager.getConsumers().size(), 1); + assertEquals(topicConsumerManager.getCursors().size(), 1); cursor2 = cursorPair.getLeft(); assertNotEquals(cursor2.getName(), cursor.getName()); assertEquals(cursorPair.getRight(), Long.valueOf(offset)); @@ -236,10 +236,10 @@ public void testTopicConsumerManagerRemoveCursorAndBacklog() throws Exception { KafkaTopicConsumerManager topicConsumerManager = tcm.get(); // before a read, first get cursor of offset. - Pair cursorPair1 = topicConsumerManager.remove(offset1); - Pair cursorPair2 = topicConsumerManager.remove(offset2); - Pair cursorPair3 = topicConsumerManager.remove(offset3); - assertEquals(topicConsumerManager.getConsumers().size(), 0); + Pair cursorPair1 = topicConsumerManager.removeCursorFuture(offset1).get(); + Pair cursorPair2 = topicConsumerManager.removeCursorFuture(offset2).get(); + Pair cursorPair3 = topicConsumerManager.removeCursorFuture(offset3).get(); + assertEquals(topicConsumerManager.getCursors().size(), 0); ManagedCursor cursor1 = cursorPair1.getLeft(); ManagedCursor cursor2 = cursorPair2.getLeft(); @@ -259,7 +259,7 @@ public void testTopicConsumerManagerRemoveCursorAndBacklog() throws Exception { topicConsumerManager.add(offset1, Pair.of(cursor1, offset1)); topicConsumerManager.add(offset2, Pair.of(cursor2, offset2)); topicConsumerManager.add(offset3, Pair.of(cursor3, offset3)); - assertEquals(topicConsumerManager.getConsumers().size(), 3); + assertEquals(topicConsumerManager.getCursors().size(), 3); // simulate cursor deleted, and backlog cleared. topicConsumerManager.deleteOneExpiredCursor(offset3); @@ -269,7 +269,7 @@ public void testTopicConsumerManagerRemoveCursorAndBacklog() throws Exception { topicConsumerManager.deleteOneExpiredCursor(offset1); verifyBacklogAndNumCursor(persistentTopic, 0, 0); - assertEquals(topicConsumerManager.getConsumers().size(), 0); + assertEquals(topicConsumerManager.getCursors().size(), 0); } // dump Topic Stats, mainly want to get and verify backlogSize. From bc5cb24ee595521f1bd9decf9d7421a35db1216f Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 9 Jun 2021 11:30:58 +0800 Subject: [PATCH 2/2] Add test to ensure cursor is only created once for a consumer --- .../kop/KafkaTopicConsumerManager.java | 1 + .../kop/KafkaTopicConsumerManagerTest.java | 33 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java index 764d03760a..99f28e94e8 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManager.java @@ -55,6 +55,7 @@ public class KafkaTopicConsumerManager implements Closeable { // used to track all created cursor, since above consumers may be remove and in fly, // use this map will not leak cursor when close. + @Getter private final ConcurrentMap createdCursors; // track last access time(millis) for offsets diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java index 25b94cca7b..4a7b9a6c1c 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerTest.java @@ -24,6 +24,9 @@ import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; +import io.streamnative.pulsar.handlers.kop.utils.KopTopic; +import java.time.Duration; +import java.util.Collections; import java.util.Properties; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicInteger; @@ -32,6 +35,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.commons.lang3.tuple.Pair; +import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; @@ -315,4 +319,33 @@ private void verifyBacklogAndNumCursor(PersistentTopic persistentTopic, assertEquals(backlog.get(), expectedBacklog); assertEquals(cursorCount.get(), numCursor); } + + @Test(timeOut = 20000) + public void testOnlyOneCursorCreated() throws Exception { + final String topic = "testOnlyOneCursorCreated"; + admin.topics().createPartitionedTopic(topic, 1); + + final int numMessages = 100; + + @Cleanup + final KafkaProducer producer = new KafkaProducer<>(newKafkaProducerProperties()); + for (int i = 0; i < numMessages; i++) { + producer.send(new ProducerRecord<>(topic, "msg-" + i)).get(); + } + + @Cleanup + final KafkaConsumer consumer = new KafkaConsumer<>(newKafkaConsumerProperties()); + consumer.subscribe(Collections.singleton(topic)); + + int numReceived = 0; + while (numReceived < numMessages) { + numReceived += consumer.poll(Duration.ofSeconds(1)).count(); + } + + final KafkaTopicConsumerManager tcm = + kafkaTopicManager.getTopicConsumerManager(new KopTopic(topic).getPartitionName(0)).get(); + // Only 1 cursor should be created for a consumer even if there were a lot of FETCH requests + // This check is to ensure that KafkaTopicConsumerManager#add is called in FETCH request handler + assertEquals(tcm.getCreatedCursors().size(), 1); + } }