From 664713383b490143ff2369c6b52b81b58368b353 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Sun, 25 Jul 2021 20:44:53 +0800 Subject: [PATCH 1/2] [PR #620] Fix deadlock caused by KafkaTopicConsumerManager --- .../kop/KafkaTopicConsumerManager.java | 186 ++++++++---------- .../handlers/kop/KafkaTopicManager.java | 38 ++-- 2 files changed, 99 insertions(+), 125 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 2b99798abd..842a8d1c6a 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 @@ -18,10 +18,12 @@ import io.streamnative.pulsar.handlers.kop.utils.MessageIdUtils; import java.io.Closeable; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.concurrent.atomic.AtomicBoolean; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCursorCallback; @@ -32,7 +34,6 @@ import org.apache.commons.codec.digest.DigestUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.service.persistent.PersistentTopic; -import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; /** * KafkaTopicConsumerManager manages a topic and its related offset cursor. @@ -43,30 +44,25 @@ public class KafkaTopicConsumerManager implements Closeable { private final PersistentTopic topic; private final KafkaRequestHandler requestHandler; - // the lock for closed status change. - // once closed, should not add new cursor back, since consumers are cleared. - private final ReentrantReadWriteLock rwLock; - private boolean closed; + private final AtomicBoolean closed = new AtomicBoolean(false); // keep fetch offset and related cursor. keep cursor and its last offset in Pair. @Getter - private final ConcurrentLongHashMap> consumers; + private final Map> consumers; // 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; + private final Map createdCursors; // track last access time(millis) for offsets @Getter - private final ConcurrentLongHashMap lastAccessTimes; + private final Map lastAccessTimes; KafkaTopicConsumerManager(KafkaRequestHandler requestHandler, PersistentTopic topic) { this.topic = topic; - this.consumers = new ConcurrentLongHashMap<>(); + this.consumers = new ConcurrentHashMap<>(); this.createdCursors = new ConcurrentHashMap<>(); - this.lastAccessTimes = new ConcurrentLongHashMap<>(); + this.lastAccessTimes = new ConcurrentHashMap<>(); this.requestHandler = requestHandler; - this.rwLock = new ReentrantReadWriteLock(); - this.closed = false; } // delete expired cursors, so backlog can be cleared. @@ -79,19 +75,13 @@ void deleteExpiredCursor(long current, long expirePeriodMillis) { } void deleteOneExpiredCursor(long offset) { - Pair pair; + if (closed.get()) { + return; + } // need not do anything, since this tcm already in closing state. and close() will delete every thing. - rwLock.readLock().lock(); - try { - if (closed) { - return; - } - pair = consumers.remove(offset); - lastAccessTimes.remove(offset); - } finally { - rwLock.readLock().unlock(); - } + final Pair pair = consumers.remove(offset); + lastAccessTimes.remove(offset); if (pair != null) { if (log.isDebugEnabled()) { @@ -106,6 +96,9 @@ void deleteOneExpiredCursor(long offset) { // delete passed in cursor. void deleteOneCursorAsync(ManagedCursor cursor, String reason) { + if (closed.get()) { + return; + } if (cursor != null) { topic.getManagedLedger().asyncDeleteCursor(cursor.getName(), new DeleteCursorCallback() { @Override @@ -130,19 +123,12 @@ public void deleteCursorFailed(ManagedLedgerException exception, Object ctx) { // 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; + if (closed.get()) { + return null; + } // should not return cursor for Fetch to read, since this tcm already in closing state. - rwLock.readLock().lock(); - try { - if (closed) { - return null; - } - cursor = consumers.remove(offset); - lastAccessTimes.remove(offset); - } finally { - rwLock.readLock().unlock(); - } + final Pair cursor = consumers.remove(offset); if (cursor != null) { if (log.isDebugEnabled()) { @@ -156,53 +142,48 @@ public Pair remove(long offset) { } private Pair createCursorIfNotExists(long offset) { + if (closed.get()) { + return null; + } // This is for read a new entry, first check if offset is from a batched message request. offset = offsetAfterBatchIndex(offset); Pair cursor; - rwLock.readLock().lock(); - try { - if (closed) { - return null; - } - // handle offset not exist in consumers, need create cursor. - consumers.computeIfAbsent( - offset, - off -> { - PositionImpl position = MessageIdUtils.getPosition(off); - - 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. - ManagedLedgerImpl ledger = (ManagedLedgerImpl) topic.getManagedLedger(); - PositionImpl previous = ledger.getPreviousPosition(position); - if (log.isDebugEnabled()) { - log.debug("[{}] Create cursor {} for offset: {}. position: {}, previousPosition: {}", - requestHandler.ctx.channel(), cursorName, off, 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(), off, previous, e); - return null; - } + // handle offset not exist in consumers, need create cursor. + consumers.computeIfAbsent( + offset, + off -> { + PositionImpl position = MessageIdUtils.getPosition(off); + + 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. + ManagedLedgerImpl ledger = (ManagedLedgerImpl) topic.getManagedLedger(); + PositionImpl previous = ledger.getPreviousPosition(position); + if (log.isDebugEnabled()) { + log.debug("[{}] Create cursor {} for offset: {}. position: {}, previousPosition: {}", + requestHandler.ctx.channel(), cursorName, off, 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(), off, previous, e); + return null; + } - lastAccessTimes.put(off, System.currentTimeMillis()); - return Pair.of(newCursor, off); - }); + lastAccessTimes.put(off, System.currentTimeMillis()); + return Pair.of(newCursor, off); + }); - // notice: above would add a - cursor = consumers.remove(offset); - lastAccessTimes.remove(offset); - } finally { - rwLock.readLock().unlock(); - } + // notice: above would add a + cursor = consumers.remove(offset); + lastAccessTimes.remove(offset); return cursor; } @@ -212,16 +193,10 @@ public void add(long offset, Pair pair) { checkArgument(offset == 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(); - deleteOneCursorAsync(managedCursor, "A race - add cursor back but tcm already closed"); - return; - } - } finally { - rwLock.readLock().unlock(); + if (closed.get()) { + ManagedCursor managedCursor = pair.getLeft(); + deleteOneCursorAsync(managedCursor, "A race - add cursor back but tcm already closed"); + return; } Pair oldPair = consumers.putIfAbsent(offset, pair); @@ -239,30 +214,23 @@ public void add(long offset, Pair pair) { // called when channel closed. @Override public void close() { - ConcurrentLongHashMap> consumersToClose; - ConcurrentMap cursorsToClose; - rwLock.writeLock().lock(); - try { - if (closed) { - return; - } - closed = true; - if (log.isDebugEnabled()) { - log.debug("[{}] Close TCM for topic {}.", - requestHandler.ctx.channel(), topic.getName()); - } - consumersToClose = new ConcurrentLongHashMap<>(); - consumers.forEach((k, v) -> consumersToClose.put(k, v)); - consumers.clear(); - lastAccessTimes.clear(); - cursorsToClose = new ConcurrentHashMap<>(); - createdCursors.forEach((k, v) -> cursorsToClose.put(k, v)); - createdCursors.clear(); - } finally { - rwLock.writeLock().unlock(); + if (!closed.compareAndSet(false, true)) { + return; } - consumersToClose.values() + if (log.isDebugEnabled()) { + log.debug("[{}] Close TCM for topic {}.", + requestHandler.ctx.channel(), topic.getName()); + } + final List> consumersToClose = new ArrayList<>(); + consumers.forEach((k, v) -> consumersToClose.add(v)); + consumers.clear(); + lastAccessTimes.clear(); + final List cursorsToClose = new ArrayList<>(); + createdCursors.forEach((k, v) -> cursorsToClose.add(v)); + createdCursors.clear(); + + consumersToClose .forEach(pair -> { ManagedCursor cursor = pair.getLeft(); deleteOneCursorAsync(cursor, "TopicConsumerManager close"); @@ -272,7 +240,7 @@ public void close() { }); // delete dangling createdCursors - cursorsToClose.values().forEach(cursor -> + cursorsToClose.forEach(cursor -> deleteOneCursorAsync(cursor, "TopicConsumerManager close but cursor is still outstanding")); cursorsToClose.clear(); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java index 9bfc1e4a90..d1af9b9202 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java @@ -21,10 +21,10 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.ReentrantReadWriteLock; - import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.PulsarServerException; @@ -65,9 +65,9 @@ public class KafkaTopicManager { // every 1 min, check if the KafkaTopicConsumerManagers have expired cursors. // remove expired cursors, so backlog can be cleared. - private long checkPeriodMillis = 1 * 60 * 1000; - private long expirePeriodMillis = 2 * 60 * 1000; - private final ScheduledFuture cursorExpireTask; + private static final long checkPeriodMillis = 1 * 60 * 1000; + private static final long expirePeriodMillis = 2 * 60 * 1000; + private static ScheduledFuture cursorExpireTask = null; // the lock for closed status change. private final ReentrantReadWriteLock rwLock; @@ -88,19 +88,25 @@ public class KafkaTopicManager { this.rwLock = new ReentrantReadWriteLock(); this.closed = false; - // check expired cursor every 1 min. - this.cursorExpireTask = brokerService.executor().scheduleWithFixedDelay(() -> { - long current = System.currentTimeMillis(); - if (log.isDebugEnabled()) { - log.debug("[{}] Schedule a check of expired cursor", - requestHandler.ctx.channel()); - } - consumerTopicManagers.values().forEach(future -> { - if (future != null && future.isDone() && !future.isCompletedExceptionally()) { - future.join().deleteExpiredCursor(current, expirePeriodMillis); + initializeCursorExpireTask(brokerService.executor()); + } + + private static void initializeCursorExpireTask(final ScheduledExecutorService executor) { + if (cursorExpireTask == null) { + synchronized (KafkaTopicManager.class) { + if (cursorExpireTask == null) { + // check expired cursor every 1 min. + cursorExpireTask = executor.scheduleWithFixedDelay(() -> { + long current = System.currentTimeMillis(); + consumerTopicManagers.values().forEach(future -> { + if (future != null && future.isDone() && !future.isCompletedExceptionally()) { + future.join().deleteExpiredCursor(current, expirePeriodMillis); + } + }); + }, checkPeriodMillis, checkPeriodMillis, TimeUnit.MILLISECONDS); } - }); - }, checkPeriodMillis, checkPeriodMillis, TimeUnit.MILLISECONDS); + } + } } // update Ctx information, since at internalServerCnx create time there is no ctx passed into kafkaRequestHandler. From bfe3e7c5045a1bbff8fa3dae1df536c0e99b41d8 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Sun, 25 Jul 2021 21:28:55 +0800 Subject: [PATCH 2/2] Fix spotbugs check and correct the time to cancel cursorExpireTask --- .../handlers/kop/KafkaProtocolHandler.java | 1 + .../handlers/kop/KafkaTopicManager.java | 28 ++++++++++++++----- 2 files changed, 22 insertions(+), 7 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index 6a3ed488b3..964158e53a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -320,6 +320,7 @@ public void close() { KafkaTopicManager.getConsumerTopicManagers().clear(); KafkaTopicManager.getReferences().clear(); KafkaTopicManager.getTopics().clear(); + KafkaTopicManager.closeKafkaTopicConsumerManagers(); OffsetAcker.CONSUMERS.clear(); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java index d1af9b9202..893a3dd637 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicManager.java @@ -21,9 +21,11 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.locks.ReentrantReadWriteLock; import lombok.Getter; import lombok.extern.slf4j.Slf4j; @@ -67,7 +69,7 @@ public class KafkaTopicManager { // remove expired cursors, so backlog can be cleared. private static final long checkPeriodMillis = 1 * 60 * 1000; private static final long expirePeriodMillis = 2 * 60 * 1000; - private static ScheduledFuture cursorExpireTask = null; + private static volatile ScheduledFuture cursorExpireTask = null; // the lock for closed status change. private final ReentrantReadWriteLock rwLock; @@ -319,6 +321,23 @@ public void registerProducerInPersistentTopic (String topicName, PersistentTopic } } + public static void closeKafkaTopicConsumerManagers() { + synchronized (KafkaTopicManager.class) { + if (cursorExpireTask != null) { + cursorExpireTask.cancel(true); + } + } + consumerTopicManagers.forEach((topic, tcmFuture) -> { + try { + Optional.ofNullable(tcmFuture.get(300, TimeUnit.SECONDS)) + .ifPresent(KafkaTopicConsumerManager::close); + } catch (InterruptedException | ExecutionException | TimeoutException e) { + log.warn("Failed to get TCM future of {} when trying to close it", topic); + } + }); + consumerTopicManagers.clear(); + } + // when channel close, release all the topics reference in persistentTopic public synchronized void close() { rwLock.writeLock().lock(); @@ -336,12 +355,7 @@ public synchronized void close() { } try { - this.cursorExpireTask.cancel(true); - - for (CompletableFuture manager : consumerTopicManagers.values()) { - manager.get().close(); - } - consumerTopicManagers.clear(); + closeKafkaTopicConsumerManagers(); for (Map.Entry> entry : topics.entrySet()) { String topicName = entry.getKey();