From f4ebd9a9598e23651c8e85cbc9f22973fd38ae09 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 12 Nov 2021 16:03:32 +0800 Subject: [PATCH 1/2] Make produce and fetch purgatory shared by all connections --- .../handlers/kop/KafkaChannelInitializer.java | 18 ++++++++++ .../handlers/kop/KafkaProtocolHandler.java | 33 +++++++++++++++-- .../handlers/kop/KafkaRequestHandler.java | 19 ++++------ .../handlers/kop/EntryPublishTimeTest.java | 36 +------------------ .../pulsar/handlers/kop/KafkaApisTest.java | 33 +---------------- .../handlers/kop/KafkaRequestHandlerTest.java | 31 +--------------- ...kaRequestHandlerWithAuthorizationTest.java | 32 +---------------- .../kop/KafkaTopicConsumerManagerTest.java | 33 +---------------- .../kop/KopProtocolHandlerTestBase.java | 27 ++++++++++++++ 9 files changed, 86 insertions(+), 176 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java index e3fa64c890..e7c8efe0c0 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaChannelInitializer.java @@ -23,6 +23,8 @@ import io.netty.handler.ssl.SslHandler; import io.netty.handler.timeout.IdleStateHandler; import io.streamnative.pulsar.handlers.kop.stats.StatsLogger; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import io.streamnative.pulsar.handlers.kop.utils.ssl.SSLUtils; import java.util.concurrent.TimeUnit; import lombok.Getter; @@ -46,6 +48,8 @@ public class KafkaChannelInitializer extends ChannelInitializer { private final KopBrokerLookupManager kopBrokerLookupManager; private final AdminManager adminManager; + private DelayedOperationPurgatory producePurgatory; + private DelayedOperationPurgatory fetchPurgatory; @Getter private final boolean enableTls; @Getter @@ -60,6 +64,8 @@ public KafkaChannelInitializer(PulsarService pulsarService, TenantContextManager tenantContextManager, KopBrokerLookupManager kopBrokerLookupManager, AdminManager adminManager, + DelayedOperationPurgatory producePurgatory, + DelayedOperationPurgatory fetchPurgatory, boolean enableTLS, EndPoint advertisedEndPoint, StatsLogger statsLogger) { @@ -69,6 +75,8 @@ public KafkaChannelInitializer(PulsarService pulsarService, this.tenantContextManager = tenantContextManager; this.kopBrokerLookupManager = kopBrokerLookupManager; this.adminManager = adminManager; + this.producePurgatory = producePurgatory; + this.fetchPurgatory = fetchPurgatory; this.enableTls = enableTLS; this.advertisedEndPoint = advertisedEndPoint; this.statsLogger = statsLogger; @@ -100,6 +108,16 @@ protected void initChannel(SocketChannel ch) throws Exception { public KafkaRequestHandler newCnx() throws Exception { return new KafkaRequestHandler(pulsarService, kafkaConfig, tenantContextManager, kopBrokerLookupManager, adminManager, + producePurgatory, fetchPurgatory, + enableTls, advertisedEndPoint, statsLogger); + } + + @VisibleForTesting + public KafkaRequestHandler newCnx(final TenantContextManager tenantContextManager, + final StatsLogger statsLogger) throws Exception { + return new KafkaRequestHandler(pulsarService, kafkaConfig, + tenantContextManager, kopBrokerLookupManager, adminManager, + producePurgatory, fetchPurgatory, enableTls, advertisedEndPoint, statsLogger); } } 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 5a865804e1..ba807d7e59 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 @@ -31,6 +31,8 @@ import io.streamnative.pulsar.handlers.kop.utils.ConfigurationUtils; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; +import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; import io.streamnative.pulsar.handlers.kop.utils.timer.SystemTimer; import java.net.InetSocketAddress; import java.util.Map; @@ -77,6 +79,8 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag private KopBrokerLookupManager kopBrokerLookupManager; private AdminManager adminManager = null; private SystemTopicClient txnTopicClient; + private DelayedOperationPurgatory producePurgatory; + private DelayedOperationPurgatory fetchPurgatory; @VisibleForTesting @Getter private Map> channelInitializerMap; @@ -547,6 +551,8 @@ private KafkaChannelInitializer newKafkaChannelInitializer(final EndPoint endPoi this, kopBrokerLookupManager, adminManager, + producePurgatory, + fetchPurgatory, endPoint.isTlsEnabled(), endPoint, scopeStatsLogger); @@ -558,6 +564,15 @@ public Map> newChannelIniti checkState(kafkaConfig != null); checkState(brokerService != null); + producePurgatory = DelayedOperationPurgatory.builder() + .purgatoryName("produce") + .timeoutTimer(SystemTimer.builder().executorName("produce").build()) + .build(); + fetchPurgatory = DelayedOperationPurgatory.builder() + .purgatoryName("fetch") + .timeoutTimer(SystemTimer.builder().executorName("fetch").build()) + .build(); + try { ImmutableMap.Builder> builder = ImmutableMap.builder(); @@ -577,9 +592,21 @@ public Map> newChannelIniti @Override public void close() { Optional.ofNullable(LOOKUP_CLIENT_MAP.remove(brokerService.pulsar())).ifPresent(LookupClient::close); - offsetTopicClient.close(); - txnTopicClient.close(); - adminManager.shutdown(); + if (offsetTopicClient != null) { + offsetTopicClient.close(); + } + if (txnTopicClient != null) { + txnTopicClient.close(); + } + if (adminManager != null) { + adminManager.shutdown(); + } + if (producePurgatory != null) { + producePurgatory.shutdown(); + } + if (fetchPurgatory != null) { + fetchPurgatory.shutdown(); + } groupCoordinatorsByTenant.values().forEach(GroupCoordinator::shutdown); kopEventManager.close(); transactionCoordinatorByTenant.values().forEach(TransactionCoordinator::shutdown); 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 a7dd259a8f..d7778f5707 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 @@ -60,7 +60,6 @@ import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperation; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationKey; import io.streamnative.pulsar.handlers.kop.utils.delayed.DelayedOperationPurgatory; -import io.streamnative.pulsar.handlers.kop.utils.timer.SystemTimer; import java.net.InetSocketAddress; import java.nio.ByteBuffer; import java.util.ArrayList; @@ -219,16 +218,8 @@ public class KafkaRequestHandler extends KafkaCommandDecoder { // is found. private final Map pendingTopicFuturesMap = new ConcurrentHashMap<>(); // DelayedOperation for produce and fetch - private final DelayedOperationPurgatory producePurgatory = - DelayedOperationPurgatory.builder() - .purgatoryName("produce") - .timeoutTimer(SystemTimer.builder().executorName("produce").build()) - .build(); - private final DelayedOperationPurgatory fetchPurgatory = - DelayedOperationPurgatory.builder() - .purgatoryName("fetch") - .timeoutTimer(SystemTimer.builder().executorName("fetch").build()) - .build(); + private final DelayedOperationPurgatory producePurgatory; + private final DelayedOperationPurgatory fetchPurgatory; // Flag to manage throttling-publish-buffer by atomically enable/disable read-channel. private final long maxPendingBytes; @@ -285,6 +276,8 @@ public KafkaRequestHandler(PulsarService pulsarService, TenantContextManager tenantContextManager, KopBrokerLookupManager kopBrokerLookupManager, AdminManager adminManager, + DelayedOperationPurgatory producePurgatory, + DelayedOperationPurgatory fetchPurgatory, Boolean tlsEnabled, EndPoint advertisedEndPoint, StatsLogger statsLogger) throws Exception { @@ -306,6 +299,8 @@ public KafkaRequestHandler(PulsarService pulsarService, ? new SimpleAclAuthorizer(pulsarService) : null; this.adminManager = adminManager; + this.producePurgatory = producePurgatory; + this.fetchPurgatory = fetchPurgatory; this.tlsEnabled = tlsEnabled; this.advertisedEndPoint = advertisedEndPoint; this.topicManager = new KafkaTopicManager(this); @@ -355,8 +350,6 @@ protected void close() { log.info("currentConnectedGroup remove {}", clientHost); currentConnectedGroup.remove(clientHost); } - producePurgatory.shutdown(); - fetchPurgatory.shutdown(); // update alive channel count stat RequestStats.ALIVE_CHANNEL_COUNT_INSTANCE.decrementAndGet(); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java index 235b9ea517..227df6f62c 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/EntryPublishTimeTest.java @@ -18,14 +18,8 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; -import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; -import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import java.net.InetSocketAddress; import java.net.SocketAddress; -import org.apache.pulsar.broker.protocol.ProtocolHandler; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -33,11 +27,9 @@ * Test for entry publish time. */ public class EntryPublishTimeTest extends KopProtocolHandlerTestBase { - private static final Logger log = LoggerFactory.getLogger(EntryPublishTimeTest.class); KafkaRequestHandler kafkaRequestHandler; SocketAddress serviceAddress; - private AdminManager adminManager; public EntryPublishTimeTest(String format) { super(format); @@ -48,32 +40,7 @@ public EntryPublishTimeTest(String format) { protected void setup() throws Exception { super.internalSetup(); - ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka"); - GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler) - .getGroupCoordinator(conf.getKafkaMetadataTenant()); - TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler) - .getTransactionCoordinator(conf.getKafkaMetadataTenant()); - - adminManager = new AdminManager(pulsar.getAdminClient(), conf); - kafkaRequestHandler = new KafkaRequestHandler( - pulsar, - conf, - new TenantContextManager() { - @Override - public GroupCoordinator getGroupCoordinator(String tenant) { - return groupCoordinator; - } - - @Override - public TransactionCoordinator getTransactionCoordinator(String tenant) { - return transactionCoordinator; - } - }, - ((KafkaProtocolHandler) handler).getKopBrokerLookupManager(), - adminManager, - false, - getPlainEndPoint(), - NullStatsLogger.INSTANCE); + kafkaRequestHandler = newRequestHandler(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); doReturn(mockChannel).when(mockCtx).channel(); @@ -85,7 +52,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { @AfterMethod @Override protected void cleanup() throws Exception { - adminManager.shutdown(); super.internalCleanup(); } } \ No newline at end of file diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java index e770ae3bd9..530904a429 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaApisTest.java @@ -30,9 +30,6 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; -import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; -import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import java.net.InetSocketAddress; import java.net.SocketAddress; import java.nio.ByteBuffer; @@ -78,7 +75,6 @@ import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; -import org.apache.pulsar.broker.protocol.ProtocolHandler; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; @@ -93,7 +89,6 @@ public class KafkaApisTest extends KopProtocolHandlerTestBase { KafkaRequestHandler kafkaRequestHandler; SocketAddress serviceAddress; - private AdminManager adminManager; @Override protected void resetConfig() { @@ -123,32 +118,7 @@ protected void setup() throws Exception { log.info("created namespaces, init handler"); - ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka"); - GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler) - .getGroupCoordinator(conf.getKafkaMetadataTenant()); - TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler) - .getTransactionCoordinator(conf.getKafkaMetadataTenant()); - - adminManager = new AdminManager(pulsar.getAdminClient(), conf); - kafkaRequestHandler = new KafkaRequestHandler( - pulsar, - (KafkaServiceConfiguration) conf, - new TenantContextManager() { - @Override - public GroupCoordinator getGroupCoordinator(String tenant) { - return groupCoordinator; - } - - @Override - public TransactionCoordinator getTransactionCoordinator(String tenant) { - return transactionCoordinator; - } - }, - ((KafkaProtocolHandler) handler).getKopBrokerLookupManager(), - adminManager, - false, - getPlainEndPoint(), - NullStatsLogger.INSTANCE); + kafkaRequestHandler = newRequestHandler(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); doReturn(mockChannel).when(mockCtx).channel(); @@ -160,7 +130,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { @AfterMethod @Override protected void cleanup() throws Exception { - adminManager.shutdown(); super.internalCleanup(); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java index 8b05a5176d..195f316107 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java @@ -39,7 +39,6 @@ import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadataManager; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils; import java.net.InetSocketAddress; import java.nio.ByteBuffer; @@ -104,7 +103,6 @@ import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.utils.Time; -import org.apache.pulsar.broker.protocol.ProtocolHandler; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; import org.apache.pulsar.common.naming.TopicName; @@ -123,7 +121,6 @@ public class KafkaRequestHandlerTest extends KopProtocolHandlerTestBase { private KafkaRequestHandler handler; - private AdminManager adminManager; @DataProvider(name = "metadataVersions") public static Object[][] metadataVersions() { @@ -153,32 +150,7 @@ protected void setup() throws Exception { log.info("created namespaces, init handler"); - ProtocolHandler handler1 = pulsar.getProtocolHandlers().protocol("kafka"); - GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler1) - .getGroupCoordinator(conf.getKafkaMetadataTenant()); - TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler1) - .getTransactionCoordinator(conf.getKafkaMetadataTenant()); - - adminManager = new AdminManager(pulsar.getAdminClient(), conf); - handler = new KafkaRequestHandler( - pulsar, - conf, - new TenantContextManager() { - @Override - public GroupCoordinator getGroupCoordinator(String tenant) { - return groupCoordinator; - } - - @Override - public TransactionCoordinator getTransactionCoordinator(String tenant) { - return transactionCoordinator; - } - }, - ((KafkaProtocolHandler) handler1).getKopBrokerLookupManager(), - adminManager, - false, - getPlainEndPoint(), - NullStatsLogger.INSTANCE); + handler = newRequestHandler(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); doReturn(mockChannel).when(mockCtx).channel(); @@ -188,7 +160,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { @AfterClass @Override protected void cleanup() throws Exception { - adminManager.shutdown(); super.internalCleanup(); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerWithAuthorizationTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerWithAuthorizationTest.java index 7291ad79ca..d972ce1d83 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerWithAuthorizationTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerWithAuthorizationTest.java @@ -30,10 +30,8 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; -import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.security.auth.Resource; import io.streamnative.pulsar.handlers.kop.security.auth.ResourceType; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import java.net.InetSocketAddress; import java.net.SocketAddress; @@ -81,7 +79,6 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.authentication.AuthenticationProviderToken; import org.apache.pulsar.broker.authentication.utils.AuthTokenUtils; -import org.apache.pulsar.broker.protocol.ProtocolHandler; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.impl.auth.AuthenticationToken; @@ -111,7 +108,6 @@ public class KafkaRequestHandlerWithAuthorizationTest extends KopProtocolHandler private String adminToken; private KafkaRequestHandler handler; - private AdminManager adminManager; @BeforeClass @Override @@ -161,32 +157,7 @@ protected void setup() throws Exception { log.info("created namespaces, init handler"); - ProtocolHandler handler1 = pulsar.getProtocolHandlers().protocol("kafka"); - GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler1) - .getGroupCoordinator(conf.getKafkaMetadataTenant()); - TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler1) - .getTransactionCoordinator(conf.getKafkaMetadataTenant()); - - adminManager = new AdminManager(pulsar.getAdminClient(), conf); - handler = new KafkaRequestHandler( - pulsar, - conf, - new TenantContextManager() { - @Override - public GroupCoordinator getGroupCoordinator(String tenant) { - return groupCoordinator; - } - - @Override - public TransactionCoordinator getTransactionCoordinator(String tenant) { - return transactionCoordinator; - } - }, - ((KafkaProtocolHandler) handler1).getKopBrokerLookupManager(), - adminManager, - false, - getPlainEndPoint(), - NullStatsLogger.INSTANCE); + handler = newRequestHandler(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); doReturn(mockChannel).when(mockCtx).channel(); @@ -208,7 +179,6 @@ protected void createAdmin() throws Exception { @AfterClass @Override protected void cleanup() throws Exception { - adminManager.shutdown(); super.internalCleanup(); } 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 02839e4902..061cb3199e 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,9 +26,6 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; -import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; -import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; -import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.utils.KopTopic; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import java.time.Duration; @@ -60,7 +57,6 @@ 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; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.TopicName; @@ -78,7 +74,6 @@ public class KafkaTopicConsumerManagerTest extends KopProtocolHandlerTestBase { private KafkaTopicManager kafkaTopicManager; private KafkaRequestHandler kafkaRequestHandler; - private AdminManager adminManager; @BeforeMethod @Override @@ -86,32 +81,7 @@ protected void setup() throws Exception { super.internalSetup(); this.triggerTopicLookup(MetadataUtils.constructOffsetsTopicBaseName( TopicName.PUBLIC_TENANT, this.conf), this.conf.getOffsetsTopicNumPartitions()); - ProtocolHandler handler = pulsar.getProtocolHandlers().protocol("kafka"); - GroupCoordinator groupCoordinator = ((KafkaProtocolHandler) handler) - .getGroupCoordinator(conf.getKafkaMetadataTenant()); - TransactionCoordinator transactionCoordinator = ((KafkaProtocolHandler) handler) - .getTransactionCoordinator(conf.getKafkaMetadataTenant()); - - adminManager = new AdminManager(pulsar.getAdminClient(), conf); - kafkaRequestHandler = new KafkaRequestHandler( - pulsar, - conf, - new TenantContextManager() { - @Override - public GroupCoordinator getGroupCoordinator(String tenant) { - return groupCoordinator; - } - - @Override - public TransactionCoordinator getTransactionCoordinator(String tenant) { - return transactionCoordinator; - } - }, - ((KafkaProtocolHandler) handler).getKopBrokerLookupManager(), - adminManager, - false, - getPlainEndPoint(), - NullStatsLogger.INSTANCE); + kafkaRequestHandler = newRequestHandler(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); Channel mockChannel = mock(Channel.class); @@ -125,7 +95,6 @@ public TransactionCoordinator getTransactionCoordinator(String tenant) { @AfterMethod @Override protected void cleanup() throws Exception { - adminManager.shutdown(); super.internalCleanup(); } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java index 03036d9474..1004663078 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KopProtocolHandlerTestBase.java @@ -23,6 +23,9 @@ import io.confluent.kafka.schemaregistry.rest.SchemaRegistryConfig; import io.confluent.kafka.schemaregistry.rest.SchemaRegistryRestApplication; import io.netty.channel.EventLoopGroup; +import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; +import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; +import io.streamnative.pulsar.handlers.kop.stats.NullStatsLogger; import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import java.io.Closeable; import java.io.IOException; @@ -752,4 +755,28 @@ protected Properties newKafkaAdminClientProperties() { return adminProps; } + public KafkaChannelInitializer getFirstChannelInitializer() { + final KafkaProtocolHandler handler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka"); + return (KafkaChannelInitializer) handler.getChannelInitializerMap().entrySet().iterator().next().getValue(); + } + + public KafkaRequestHandler newRequestHandler() throws Exception { + final KafkaProtocolHandler handler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol("kafka"); + final GroupCoordinator groupCoordinator = handler.getGroupCoordinator(conf.getKafkaMetadataTenant()); + final TransactionCoordinator transactionCoordinator = + handler.getTransactionCoordinator(conf.getKafkaMetadataTenant()); + + return ((KafkaChannelInitializer) handler.getChannelInitializerMap().entrySet().iterator().next().getValue()) + .newCnx(new TenantContextManager() { + @Override + public GroupCoordinator getGroupCoordinator(String tenant) { + return groupCoordinator; + } + + @Override + public TransactionCoordinator getTransactionCoordinator(String tenant) { + return transactionCoordinator; + } + }, NullStatsLogger.INSTANCE); + } } From 426a904404af2fcd5672b4f9bce6b105075fd661 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 12 Nov 2021 19:28:19 +0800 Subject: [PATCH 2/2] Fix checkstyle --- .../pulsar/handlers/kop/KafkaRequestHandlerTest.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java index 195f316107..fac6e5c01d 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaRequestHandlerTest.java @@ -34,10 +34,8 @@ import io.netty.channel.ChannelHandlerContext; import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndResponse; -import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupCoordinator; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadata; import io.streamnative.pulsar.handlers.kop.coordinator.group.GroupMetadataManager; -import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionCoordinator; import io.streamnative.pulsar.handlers.kop.offset.OffsetAndMetadata; import io.streamnative.pulsar.handlers.kop.utils.TopicNameUtils; import java.net.InetSocketAddress;