From a9998637808ee497e5d6edc9367185647c7d8aba Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Tue, 17 Aug 2021 22:01:43 +0800 Subject: [PATCH] Remove closed KafkaTopicConsumerManager from cache --- .../pulsar/handlers/kop/MessageFetchContext.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) 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 090bffdf06..767683f58a 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 @@ -250,7 +250,7 @@ public void handleFetch() { statsLogger.getPrepareMetadataStats().registerFailedEvent( MathUtils.elapsedNanos(startPrepareMetadataNanos), TimeUnit.NANOSECONDS); // remove null future cache - KafkaTopicManager.removeKafkaTopicConsumerManager(KopTopic.toString(topicPartition)); + KafkaTopicManager.removeKafkaTopicConsumerManager(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NOT_LEADER_FOR_PARTITION); return; } @@ -279,7 +279,9 @@ public void handleFetch() { 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); + log.warn("[{}] KafkaTopicConsumerManager is closed, remove TCM of {}", + requestHandler.ctx, fullTopicName); + KafkaTopicManager.removeKafkaTopicConsumerManager(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NONE); return; }