diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalServerCnx.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalServerCnx.java index 9233a961a1..195f256d84 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalServerCnx.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/InternalServerCnx.java @@ -14,6 +14,7 @@ package io.streamnative.pulsar.handlers.kop; import java.net.InetSocketAddress; +import java.net.SocketAddress; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.service.Producer; @@ -27,6 +28,9 @@ */ @Slf4j public class InternalServerCnx extends ServerCnx { + + public static final SocketAddress MOCKED_REMOTE_ADDRESS = new InetSocketAddress("localhost", 9999); + @Getter KafkaRequestHandler kafkaRequestHandler; @@ -39,7 +43,7 @@ public InternalServerCnx(KafkaRequestHandler kafkaRequestHandler) { // mock some values, or Producer create will meet NPE. // used in test, which will not call channel.active, and not call updateCtx. if (this.remoteAddress == null) { - this.remoteAddress = new InetSocketAddress("localhost", 9999); + this.remoteAddress = MOCKED_REMOTE_ADDRESS; } } @@ -56,8 +60,8 @@ public void closeProducer(Producer producer) { } // called after channel active - public void updateCtx() { - this.remoteAddress = kafkaRequestHandler.remoteAddress; + public void updateCtx(final SocketAddress remoteAddress) { + this.remoteAddress = remoteAddress; } @Override diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java index 5c1a7d4e72..074f1e57ad 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandler.java @@ -294,7 +294,7 @@ public KafkaRequestHandler(PulsarService pulsarService, @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { super.channelActive(ctx); - getTopicManager().updateCtx(); + topicManager.setRemoteAddress(ctx.channel().remoteAddress()); if (authenticator != null) { authenticator.reset(); } 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 5a56d9348d..cc1cec1eb0 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,6 +15,7 @@ import static com.google.common.base.Preconditions.checkArgument; +import com.google.common.annotations.VisibleForTesting; import io.streamnative.pulsar.handlers.kop.utils.OffsetSearchPredicate; import java.io.Closeable; import java.util.ArrayList; @@ -24,6 +25,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCursorCallback; @@ -42,6 +44,10 @@ */ @Slf4j public class KafkaTopicConsumerManager implements Closeable { + + private static final AtomicIntegerFieldUpdater NUM_CREATED_CURSORS_UPDATER = + AtomicIntegerFieldUpdater.newUpdater(KafkaTopicConsumerManager.class, "numCreatedCursors"); + private final PersistentTopic topic; private final KafkaRequestHandler requestHandler; @@ -60,6 +66,9 @@ public class KafkaTopicConsumerManager implements Closeable { @Getter private final Map lastAccessTimes; + // Record the number of created cursor + private volatile int numCreatedCursors = 0; + KafkaTopicConsumerManager(KafkaRequestHandler requestHandler, PersistentTopic topic) { this.topic = topic; this.cursors = new ConcurrentHashMap<>(); @@ -191,6 +200,7 @@ public void close() { log.debug("[{}] Close TCM for topic {}.", requestHandler.ctx.channel(), topic.getName()); } + NUM_CREATED_CURSORS_UPDATER.set(this, 0); final List>> cursorFuturesToClose = new ArrayList<>(); cursors.forEach((ignored, cursorFuture) -> cursorFuturesToClose.add(cursorFuture)); cursors.clear(); @@ -235,6 +245,7 @@ private CompletableFuture> asyncGetCursorByOffset(long } try { final ManagedCursor newCursor = ledger.newNonDurableCursor(previous, cursorName); + NUM_CREATED_CURSORS_UPDATER.incrementAndGet(this); createdCursors.putIfAbsent(newCursor.getName(), newCursor); lastAccessTimes.put(offset, System.currentTimeMillis()); return Pair.of(newCursor, offset); @@ -249,4 +260,9 @@ private CompletableFuture> asyncGetCursorByOffset(long public ManagedLedger getManagedLedger() { return topic.getManagedLedger(); } + + @VisibleForTesting + public int getNumCreatedCursors() { + return numCreatedCursors; + } } 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 new file mode 100644 index 0000000000..3e76f62530 --- /dev/null +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaTopicConsumerManagerCache.java @@ -0,0 +1,112 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.kop; + +import com.google.common.annotations.VisibleForTesting; +import java.net.SocketAddress; +import java.util.Collections; +import java.util.List; +import java.util.Map; +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; +import java.util.stream.Collectors; +import lombok.NonNull; +import lombok.extern.slf4j.Slf4j; + +/** + * The cache for {@link KafkaTopicConsumerManager}, aka TCM. + */ +@Slf4j +public class KafkaTopicConsumerManagerCache { + + private static final KafkaTopicConsumerManagerCache TCM_CACHE = new KafkaTopicConsumerManagerCache(); + + // The 1st key is the full topic name, the 2nd key is the remote address of Kafka client. + // Because a topic could have multiple connected consumers, for different consumers we should maintain different + // KafkaTopicConsumerManagers, which are responsible for maintaining the cursors. + private final Map>> + cache = new ConcurrentHashMap<>(); + + public static KafkaTopicConsumerManagerCache getInstance() { + return TCM_CACHE; + } + + private KafkaTopicConsumerManagerCache() { + // No ops + } + + public CompletableFuture computeIfAbsent( + final String fullTopicName, + final SocketAddress remoteAddress, + final Supplier> mappingFunction) { + return cache.computeIfAbsent(fullTopicName, ignored -> new ConcurrentHashMap<>()) + .computeIfAbsent(remoteAddress, ignored -> mappingFunction.get()); + } + + public void forEach(final Consumer> action) { + cache.values().forEach(internalMap -> { + internalMap.values().forEach(action); + }); + } + + public void removeAndClose(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(); + } + }); + })); + } + + 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); + } + }); + }); + } + + @VisibleForTesting + public int getCount() { + final AtomicInteger count = new AtomicInteger(0); + forEach(ignored -> count.incrementAndGet()); + return count.get(); + } + + @VisibleForTesting + public @NonNull List getTopicConsumerManagers(final String fullTopicName) { + return cache.getOrDefault(fullTopicName, Collections.emptyMap()).values().stream() + .map(CompletableFuture::join) + .collect(Collectors.toList()); + } +} 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 c2bd35b923..14ada50cce 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 @@ -13,8 +13,8 @@ */ package io.streamnative.pulsar.handlers.kop; -import com.google.common.annotations.VisibleForTesting; import java.net.InetSocketAddress; +import java.net.SocketAddress; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; @@ -22,7 +22,6 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicBoolean; import lombok.Getter; import lombok.extern.slf4j.Slf4j; @@ -41,14 +40,12 @@ @Slf4j public class KafkaTopicManager { + private static final KafkaTopicConsumerManagerCache TCM_CACHE = KafkaTopicConsumerManagerCache.getInstance(); + private final KafkaRequestHandler requestHandler; - private final PulsarService pulsarService; private final BrokerService brokerService; private final LookupClient lookupClient; - - // consumerTopicManagers for consumers cache. - private static final ConcurrentHashMap> - consumerTopicManagers = new ConcurrentHashMap<>(); + private volatile SocketAddress remoteAddress; // cache for topics: , for removing producer @Getter @@ -77,7 +74,7 @@ public class KafkaTopicManager { KafkaTopicManager(KafkaRequestHandler kafkaRequestHandler) { this.requestHandler = kafkaRequestHandler; - this.pulsarService = kafkaRequestHandler.getPulsarService(); + PulsarService pulsarService = kafkaRequestHandler.getPulsarService(); this.brokerService = pulsarService.getBrokerService(); this.internalServerCnx = new InternalServerCnx(requestHandler); this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); @@ -92,7 +89,7 @@ private static void initializeCursorExpireTask(final ScheduledExecutorService ex // check expired cursor every 1 min. cursorExpireTask = executor.scheduleWithFixedDelay(() -> { long current = System.currentTimeMillis(); - consumerTopicManagers.values().forEach(future -> { + TCM_CACHE.forEach(future -> { if (future != null && future.isDone() && !future.isCompletedExceptionally()) { future.join().deleteExpiredCursor(current, expirePeriodMillis); } @@ -104,8 +101,9 @@ private static void initializeCursorExpireTask(final ScheduledExecutorService ex } // update Ctx information, since at internalServerCnx create time there is no ctx passed into kafkaRequestHandler. - public void updateCtx() { - internalServerCnx.updateCtx(); + public void setRemoteAddress(SocketAddress remoteAddress) { + internalServerCnx.updateCtx(remoteAddress); + this.remoteAddress = remoteAddress; } // topicName is in pulsar format. e.g. persistent://public/default/topic-partition-0 @@ -118,11 +116,17 @@ public CompletableFuture getTopicConsumerManager(Stri } return CompletableFuture.completedFuture(null); } - return consumerTopicManagers.computeIfAbsent( + if (remoteAddress == null) { + log.error("[{}] Try to getTopicConsumerManager({}) while remoteAddress is not set", + requestHandler.ctx.channel(), topicName); + return CompletableFuture.completedFuture(null); + } + return TCM_CACHE.computeIfAbsent( topicName, - t -> { + remoteAddress, + () -> { final CompletableFuture tcmFuture = new CompletableFuture<>(); - getTopic(t).whenComplete((persistentTopic, throwable) -> { + getTopic(topicName).whenComplete((persistentTopic, throwable) -> { if (persistentTopic.isPresent() && throwable == null) { if (log.isDebugEnabled()) { log.debug("[{}] Call getTopicConsumerManager for {}, and create TCM for {}.", @@ -321,25 +325,13 @@ public static void deReference(String topicName) { try { removeTopicManagerCache(topicName); - Optional.ofNullable(consumerTopicManagers.remove(topicName)).ifPresent( - // Use thenAccept to avoid blocking - tcmFuture -> tcmFuture.thenAccept(tcm -> { - if (tcm != null) { - tcm.close(); - } - }) - ); - + TCM_CACHE.removeAndClose(topicName); removePersistentTopicAndReferenceProducer(topicName); } catch (Exception e) { log.error("Failed to close reference for individual topic {}. exception:", topicName, e); } } - public static void removeKafkaTopicConsumerManager(String topicName) { - consumerTopicManagers.remove(topicName); - } - public static void closeKafkaTopicConsumerManagers() { synchronized (KafkaTopicManager.class) { if (cursorExpireTask != null) { @@ -347,19 +339,6 @@ public static void closeKafkaTopicConsumerManagers() { cursorExpireTask = null; } } - 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(); - } - - @VisibleForTesting - public static int getNumberOfKafkaTopicConsumerManagers() { - return consumerTopicManagers.size(); + 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 3561ec65f8..436d7cba62 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 @@ -281,7 +281,7 @@ private void handlePartitionData(final TopicPartition topicPartition, statsLogger.getPrepareMetadataStats().registerFailedEvent( MathUtils.elapsedNanos(startPrepareMetadataNanos), TimeUnit.NANOSECONDS); // remove null future cache - KafkaTopicManager.removeKafkaTopicConsumerManager(fullTopicName); + KafkaTopicConsumerManagerCache.getInstance().removeAndClose(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NOT_LEADER_FOR_PARTITION); return; } @@ -315,7 +315,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); - KafkaTopicManager.removeKafkaTopicConsumerManager(fullTopicName); + KafkaTopicConsumerManagerCache.getInstance().removeAndClose(fullTopicName); addErrorPartitionResponse(topicPartition, Errors.NONE); return; } 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 be6c6589a9..d76997b45e 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 @@ -27,14 +27,23 @@ import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.time.Duration; import java.util.Collections; +import java.util.List; import java.util.Properties; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; +import java.util.stream.Collectors; +import java.util.stream.IntStream; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.commons.lang3.tuple.Pair; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; @@ -46,6 +55,7 @@ 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; @@ -87,6 +97,7 @@ protected void setup() throws Exception { kafkaRequestHandler.ctx = mockCtx; kafkaTopicManager = new KafkaTopicManager(kafkaRequestHandler); + kafkaTopicManager.setRemoteAddress(InternalServerCnx.MOCKED_REMOTE_ADDRESS); } @AfterMethod @@ -113,7 +124,7 @@ public void testGetTopicConsumerManager() throws Exception { KafkaTopicConsumerManager topicConsumerManager2 = tcm.get(); assertTrue(topicConsumerManager == topicConsumerManager2); - assertEquals(KafkaTopicManager.getNumberOfKafkaTopicConsumerManagers(), 1); + assertEquals(KafkaTopicConsumerManagerCache.getInstance().getCount(), 1); // 2. verify another get with different topic will return different tcm String topicName2 = "persistent://public/default/testGetTopicConsumerManager2"; @@ -121,7 +132,7 @@ public void testGetTopicConsumerManager() throws Exception { tcm = kafkaTopicManager.getTopicConsumerManager(topicName2); topicConsumerManager2 = tcm.get(); assertTrue(topicConsumerManager != topicConsumerManager2); - assertEquals(KafkaTopicManager.getNumberOfKafkaTopicConsumerManagers(), 2); + assertEquals(KafkaTopicConsumerManagerCache.getInstance().getCount(), 2); } @@ -325,6 +336,7 @@ private void verifyBacklogAndNumCursor(PersistentTopic persistentTopic, @Test(timeOut = 20000) public void testOnlyOneCursorCreated() throws Exception { final String topic = "testOnlyOneCursorCreated"; + final String partitionName = new KopTopic(topic).getPartitionName(0); admin.topics().createPartitionedTopic(topic, 1); final int numMessages = 100; @@ -344,10 +356,71 @@ public void testOnlyOneCursorCreated() throws Exception { numReceived += consumer.poll(Duration.ofSeconds(1)).count(); } - final KafkaTopicConsumerManager tcm = - kafkaTopicManager.getTopicConsumerManager(new KopTopic(topic).getPartitionName(0)).get(); + final List tcmList = + KafkaTopicConsumerManagerCache.getInstance().getTopicConsumerManagers(partitionName); + Assert.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(tcm.getCreatedCursors().size(), 1); + assertEquals(tcmList.get(0).getCreatedCursors().size(), 1); + assertEquals(tcmList.get(0).getNumCreatedCursors(), 1); + } + + @Test(timeOut = 20000) + public void testCursorCountForMultiGroups() throws Exception { + final String topic = "test-cursor-count-for-multi-groups"; + final String partitionName = new KopTopic(topic).getPartitionName(0); + final int numMessages = 100; + final int numConsumers = 5; + + final KafkaProducer producer = new KafkaProducer<>(newKafkaProducerProperties()); + for (int i = 0; i < numMessages; i++) { + producer.send(new ProducerRecord<>(topic, "msg-" + i)).get(); + } + producer.close(); + + final List> consumers = IntStream.range(0, numConsumers) + .mapToObj(i -> { + final Properties props = newKafkaConsumerProperties(); + props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "group-" + i); + final KafkaConsumer consumer = new KafkaConsumer<>(props); + consumer.subscribe(Collections.singleton(topic)); + return consumer; + }).collect(Collectors.toList()); + + final CountDownLatch latch = new CountDownLatch(numConsumers); + final ExecutorService executor = Executors.newFixedThreadPool(numConsumers); + for (int i = 0; i < numConsumers; i++) { + final int index = i; + final KafkaConsumer consumer = consumers.get(i); + executor.execute(() -> { + int numReceived = 0; + while (numReceived < numMessages) { + final ConsumerRecords records = consumer.poll(Duration.ofSeconds(1)); + records.forEach(record -> { + if (log.isDebugEnabled()) { + log.debug("Group {} received message {}", index, record.value()); + } + }); + numReceived += records.count(); + } + latch.countDown(); + }); + } + latch.await(10, TimeUnit.SECONDS); + + final List tcmList = + KafkaTopicConsumerManagerCache.getInstance().getTopicConsumerManagers(partitionName); + assertEquals(tcmList.size(), numConsumers); + + for (int i = 0; i < numConsumers; i++) { + assertEquals(tcmList.get(i).getNumCreatedCursors(), 1); + } + + // Since consumer close will make connection disconnected and all TCMs will be cleared, we should call it after + // the test is verified. + consumers.forEach(KafkaConsumer::close); + for (int i = 0; i < numConsumers; i++) { + assertEquals(tcmList.get(i).getNumCreatedCursors(), 0); + } } }