From 09cb22dfa3838cf6ec4d14336bff7984c228a91e Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Mon, 25 Oct 2021 17:35:56 +0800 Subject: [PATCH] Fix flaky test --- .../handlers/kop/KafkaTopicConsumerManagerTest.java | 8 ++++++++ 1 file changed, 8 insertions(+) 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 c680bc738a..f39c77dc6d 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 @@ -486,6 +486,11 @@ public void testTopicManagerClose() throws Exception { assertFalse(originalTcmList.get(1).isClosed()); consumers.get(1).close(); // trigger KafkaTopicManager#close, only the partition 1 related cache was removed + // Because the KafkaRequestHandler.close() is called by channelInActive, when channelInActive called, + // the tcp connect already closed. We need ensure topicManager.close() is called. + Awaitility.await() + .atMost(Duration.ofSeconds(3)) + .until(() -> originalTcmList.get(1).getNumCreatedCursors() == 0); 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 @@ -493,6 +498,9 @@ public void testTopicManagerClose() throws Exception { assertTrue(originalTcmList.get(1).isClosed()); consumers.get(0).close(); // Now all TCM cache was cleared + Awaitility.await() + .atMost(Duration.ofSeconds(3)) + .until(() -> originalTcmList.get(0).getNumCreatedCursors() == 0); assertNull(getTcmForPartition.apply(0)); assertNull(getTcmForPartition.apply(1)); assertTrue(originalTcmList.get(0).isClosed());