From 76d6053fc4bfc56105f6290a681c7146450ef138 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Sat, 2 Oct 2021 20:53:25 +0800 Subject: [PATCH 1/4] Remove TCM cache by address --- .../handlers/kop/KafkaProtocolHandler.java | 2 +- .../kop/KafkaTopicConsumerManagerCache.java | 22 +++++++++--------- .../handlers/kop/KafkaTopicManager.java | 23 +++++++++---------- .../handlers/kop/MessageFetchContext.java | 4 ++-- 4 files changed, 25 insertions(+), 26 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 abb1f9281a..867239e0cd 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 @@ -472,7 +472,7 @@ public void close() { kopEventManager.close(); KafkaTopicManager.LOOKUP_CACHE.clear(); KopBrokerLookupManager.clear(); - KafkaTopicManager.closeKafkaTopicConsumerManagers(); + KafkaTopicManager.cancelCursorExpireTask(); KafkaTopicManager.getReferences().clear(); KafkaTopicManager.getTopics().clear(); statsProvider.stop(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java index 3e76f62530..94a8396c64 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java @@ -21,9 +21,6 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; import java.util.function.Supplier; @@ -67,7 +64,7 @@ public void forEach(final Consumer> }); } - public void removeAndClose(final String fullTopicName) { + public void removeAndCloseByTopic(final String fullTopicName) { // The TCM future could be completed with null, so we should process this case Optional.ofNullable(cache.remove(fullTopicName)).ifPresent(map -> map.forEach((remoteAddress, future) -> { @@ -83,15 +80,18 @@ public void removeAndClose(final String fullTopicName) { })); } - public void close() { + public void removeAndCloseByAddress(final SocketAddress remoteAddress) { cache.forEach((fullTopicName, internalMap) -> { - internalMap.forEach((remoteAddress, future) -> { - try { - Optional.ofNullable(future.get(100, TimeUnit.MILLISECONDS)) - .ifPresent(KafkaTopicConsumerManager::close); - } catch (InterruptedException | ExecutionException | TimeoutException e) { - log.warn("[{}][{}] Failed to get TCM future when trying to close it", fullTopicName, remoteAddress); + Optional.ofNullable(internalMap.remove(remoteAddress)).ifPresent(future -> { + if (log.isDebugEnabled()) { + log.debug("[{}][{}] Remove and close TCM", fullTopicName, remoteAddress); } + // Use thenAccept to avoid blocking + future.thenAccept(tcm -> { + if (tcm != null) { + tcm.close();; + } + }); }); }); } 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 df8c4db836..74d90613d7 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 @@ -101,6 +101,15 @@ private static void initializeCursorExpireTask(final ScheduledExecutorService ex } } + public static void cancelCursorExpireTask() { + synchronized (KafkaTopicManager.class) { + if (cursorExpireTask != null) { + cursorExpireTask.cancel(true); + cursorExpireTask = null; + } + } + } + // update Ctx information, since at internalServerCnx create time there is no ctx passed into kafkaRequestHandler. public void setRemoteAddress(SocketAddress remoteAddress) { internalServerCnx.updateCtx(remoteAddress); @@ -326,7 +335,7 @@ public void close() { } try { - closeKafkaTopicConsumerManagers(); + TCM_CACHE.removeAndCloseByAddress(remoteAddress); topics.keySet().forEach(topicName -> { if (log.isDebugEnabled()) { @@ -377,20 +386,10 @@ public static void deReference(String topicName) { try { removeTopicManagerCache(topicName); - TCM_CACHE.removeAndClose(topicName); + TCM_CACHE.removeAndCloseByTopic(topicName); removePersistentTopicAndReferenceProducer(topicName); } catch (Exception e) { log.error("Failed to close reference for individual topic {}. exception:", topicName, e); } } - - public static void closeKafkaTopicConsumerManagers() { - synchronized (KafkaTopicManager.class) { - if (cursorExpireTask != null) { - cursorExpireTask.cancel(true); - cursorExpireTask = null; - } - } - TCM_CACHE.close(); - } } 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 f7a67e2d60..788b409f64 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 @@ -301,7 +301,7 @@ private void handlePartitionData(final TopicPartition topicPartition, statsLogger.getPrepareMetadataStats().registerFailedEvent( MathUtils.elapsedNanos(startPrepareMetadataNanos), TimeUnit.NANOSECONDS); // remove null future cache - KafkaTopicConsumerManagerCache.getInstance().removeAndClose(fullTopicName); + KafkaTopicConsumerManagerCache.getInstance().removeAndCloseByTopic(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NOT_LEADER_FOR_PARTITION); return; } @@ -335,7 +335,7 @@ private void handlePartitionData(final TopicPartition topicPartition, // tcm is closed, just return a NONE error because the channel may be still active log.warn("[{}] KafkaTopicConsumerManager is closed, remove TCM of {}", requestHandler.ctx, fullTopicName); - KafkaTopicConsumerManagerCache.getInstance().removeAndClose(fullTopicName); + KafkaTopicConsumerManagerCache.getInstance().removeAndCloseByTopic(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NONE); return; } From 1edfb53b6e0645159a1dbac39dcb7a1a3b2e8db6 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 8 Oct 2021 11:49:35 +0800 Subject: [PATCH 2/4] Add tests for closing any producer or consumer --- .../kop/KafkaTopicConsumerManager.java | 5 ++ .../kop/KafkaTopicConsumerManagerTest.java | 59 +++++++++++++++++++ 2 files changed, 64 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 cc1cec1eb0..e5884ca6e3 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 @@ -265,4 +265,9 @@ public ManagedLedger getManagedLedger() { public int getNumCreatedCursors() { return numCreatedCursors; } + + @VisibleForTesting + public boolean isClosed() { + return closed.get(); + } } 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 d811c3fc62..f5b9fb6219 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 @@ -26,6 +26,7 @@ import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.time.Duration; +import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Properties; @@ -36,6 +37,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Function; import java.util.stream.Collectors; import java.util.stream.IntStream; import lombok.Cleanup; @@ -48,6 +50,7 @@ import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.broker.protocol.ProtocolHandler; @@ -434,4 +437,60 @@ public void testCursorCountForMultiGroups() throws Exception { assertEquals(tcmList.get(i).getNumCreatedCursors(), 0); } } + + // KafkaTopicManager#close should only remove TCM cache for the specific address + @Test(timeOut = 20000) + public void testTopicManagerClose() throws Exception { + final String topic = "test-topic-manager-close"; + final int numPartitions = 2; + admin.topics().createPartitionedTopic(topic, numPartitions); + + final List> consumers = new ArrayList<>(); + for (int i = 0; i < numPartitions; i++) { + consumers.add(new KafkaConsumer<>(newKafkaConsumerProperties())); + consumers.get(i).assign(Collections.singleton(new TopicPartition(topic, i))); + } + + final KafkaProducer producer = new KafkaProducer<>(newKafkaProducerProperties()); + for (int i = 0; i < numPartitions; i++) { + producer.send(new ProducerRecord<>(topic, i, null, "msg-" + i)).get(); + final ConsumerRecords records = consumers.get(i).poll(Duration.ofSeconds(1)); + Assert.assertEquals(records.count(), 1); + Assert.assertEquals(records.iterator().next().value(), "msg-" + i); + } + + final Function getTcmForPartition = partition -> { + final String fullTopicName = new KopTopic(topic).getPartitionName(partition); + final List tcmList = + KafkaTopicConsumerManagerCache.getInstance().getTopicConsumerManagers(fullTopicName); + return tcmList.isEmpty() ? null : tcmList.get(0); + }; + + final List originalTcmList = new ArrayList<>(); + for (int i = 0; i < numPartitions; i++) { + final KafkaTopicConsumerManager tcm = getTcmForPartition.apply(i); + Assert.assertNotNull(tcm); + Assert.assertFalse(tcm.isClosed()); + originalTcmList.add(tcm); + } + + producer.close(); // trigger KafkaTopicManager#close but the TCM cache was not affected + Assert.assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); + Assert.assertFalse(originalTcmList.get(0).isClosed()); + Assert.assertSame(getTcmForPartition.apply(1), originalTcmList.get(1)); + Assert.assertFalse(originalTcmList.get(1).isClosed()); + + consumers.get(1).close(); // trigger KafkaTopicManager#close, only the partition 1 related cache was removed + Assert.assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); + Assert.assertFalse(originalTcmList.get(0).isClosed()); + // The tcm of partition 1 was closed and it was removed from cache + Assert.assertNull(getTcmForPartition.apply(1)); + Assert.assertTrue(originalTcmList.get(1).isClosed()); + + consumers.get(0).close(); // Now all TCM cache was cleared + Assert.assertNull(getTcmForPartition.apply(0)); + Assert.assertNull(getTcmForPartition.apply(1)); + Assert.assertTrue(originalTcmList.get(0).isClosed()); + Assert.assertTrue(originalTcmList.get(1).isClosed()); + } } From 36b5e38f55168b7831186452b4f9ba1a3a44af43 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 8 Oct 2021 12:07:33 +0800 Subject: [PATCH 3/4] Close TCM cache in protocol handler's close --- .../handlers/kop/KafkaProtocolHandler.java | 1 + .../kop/KafkaTopicConsumerManagerCache.java | 40 +++++++++++++------ 2 files changed, 28 insertions(+), 13 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 867239e0cd..6739bb19e2 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 @@ -473,6 +473,7 @@ public void close() { KafkaTopicManager.LOOKUP_CACHE.clear(); KopBrokerLookupManager.clear(); KafkaTopicManager.cancelCursorExpireTask(); + KafkaTopicConsumerManagerCache.getInstance().close(); KafkaTopicManager.getReferences().clear(); KafkaTopicManager.getTopics().clear(); statsProvider.stop(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java index 94a8396c64..478e76db9d 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java @@ -21,6 +21,9 @@ import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; import java.util.function.Supplier; @@ -64,19 +67,22 @@ public void forEach(final Consumer> }); } + private static void closeTcmFuture(final CompletableFuture tcmFuture) { + // Use thenAccept to avoid blocking + tcmFuture.thenAccept(tcm -> { + if (tcm != null) { + tcm.close(); + } + }); + } + public void removeAndCloseByTopic(final String fullTopicName) { - // The TCM future could be completed with null, so we should process this case Optional.ofNullable(cache.remove(fullTopicName)).ifPresent(map -> map.forEach((remoteAddress, future) -> { if (log.isDebugEnabled()) { log.debug("[{}][{}] Remove and close TCM", fullTopicName, remoteAddress); } - // Use thenAccept to avoid blocking - future.thenAccept(tcm -> { - if (tcm != null) { - tcm.close(); - } - }); + closeTcmFuture(future); })); } @@ -86,12 +92,20 @@ public void removeAndCloseByAddress(final SocketAddress remoteAddress) { if (log.isDebugEnabled()) { log.debug("[{}][{}] Remove and close TCM", fullTopicName, remoteAddress); } - // Use thenAccept to avoid blocking - future.thenAccept(tcm -> { - if (tcm != null) { - tcm.close();; - } - }); + closeTcmFuture(future); + }); + }); + } + + public void close() { + cache.forEach((fullTopicName, internalMap) -> { + internalMap.forEach((remoteAddress, future) -> { + try { + Optional.ofNullable(future.get(100, TimeUnit.MILLISECONDS)) + .ifPresent(KafkaTopicConsumerManager::close); + } catch (InterruptedException | ExecutionException | TimeoutException e) { + log.warn("[{}][{}] Failed to get TCM future when trying to close it", fullTopicName, remoteAddress); + } }); }); } From d94aff3818bb9aec16f9408b031d706ac0ec0954 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 8 Oct 2021 12:38:13 +0800 Subject: [PATCH 4/4] Fix codacy check --- .../kop/KafkaTopicConsumerManagerTest.java | 39 ++++++++++--------- 1 file changed, 21 insertions(+), 18 deletions(-) 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 f5b9fb6219..0084c20985 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 @@ -16,7 +16,11 @@ import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotEquals; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertSame; import static org.testng.Assert.assertTrue; import io.netty.channel.Channel; @@ -58,7 +62,6 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; -import org.testng.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -372,7 +375,7 @@ public void testOnlyOneCursorCreated() throws Exception { final List tcmList = KafkaTopicConsumerManagerCache.getInstance().getTopicConsumerManagers(partitionName); - Assert.assertFalse(tcmList.isEmpty()); + assertFalse(tcmList.isEmpty()); // 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(tcmList.get(0).getCreatedCursors().size(), 1); @@ -455,8 +458,8 @@ public void testTopicManagerClose() throws Exception { for (int i = 0; i < numPartitions; i++) { producer.send(new ProducerRecord<>(topic, i, null, "msg-" + i)).get(); final ConsumerRecords records = consumers.get(i).poll(Duration.ofSeconds(1)); - Assert.assertEquals(records.count(), 1); - Assert.assertEquals(records.iterator().next().value(), "msg-" + i); + assertEquals(records.count(), 1); + assertEquals(records.iterator().next().value(), "msg-" + i); } final Function getTcmForPartition = partition -> { @@ -469,28 +472,28 @@ public void testTopicManagerClose() throws Exception { final List originalTcmList = new ArrayList<>(); for (int i = 0; i < numPartitions; i++) { final KafkaTopicConsumerManager tcm = getTcmForPartition.apply(i); - Assert.assertNotNull(tcm); - Assert.assertFalse(tcm.isClosed()); + assertNotNull(tcm); + assertFalse(tcm.isClosed()); originalTcmList.add(tcm); } producer.close(); // trigger KafkaTopicManager#close but the TCM cache was not affected - Assert.assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); - Assert.assertFalse(originalTcmList.get(0).isClosed()); - Assert.assertSame(getTcmForPartition.apply(1), originalTcmList.get(1)); - Assert.assertFalse(originalTcmList.get(1).isClosed()); + assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); + assertFalse(originalTcmList.get(0).isClosed()); + assertSame(getTcmForPartition.apply(1), originalTcmList.get(1)); + assertFalse(originalTcmList.get(1).isClosed()); consumers.get(1).close(); // trigger KafkaTopicManager#close, only the partition 1 related cache was removed - Assert.assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); - Assert.assertFalse(originalTcmList.get(0).isClosed()); + assertSame(getTcmForPartition.apply(0), originalTcmList.get(0)); + assertFalse(originalTcmList.get(0).isClosed()); // The tcm of partition 1 was closed and it was removed from cache - Assert.assertNull(getTcmForPartition.apply(1)); - Assert.assertTrue(originalTcmList.get(1).isClosed()); + assertNull(getTcmForPartition.apply(1)); + assertTrue(originalTcmList.get(1).isClosed()); consumers.get(0).close(); // Now all TCM cache was cleared - Assert.assertNull(getTcmForPartition.apply(0)); - Assert.assertNull(getTcmForPartition.apply(1)); - Assert.assertTrue(originalTcmList.get(0).isClosed()); - Assert.assertTrue(originalTcmList.get(1).isClosed()); + assertNull(getTcmForPartition.apply(0)); + assertNull(getTcmForPartition.apply(1)); + assertTrue(originalTcmList.get(0).isClosed()); + assertTrue(originalTcmList.get(1).isClosed()); } }