From cb356bbc7613e984b2080db3e55e5ae61426c310 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 18 Aug 2021 21:04:20 +0800 Subject: [PATCH] Avoid creating a lot of MetadataCache instances --- .../pulsar/handlers/kop/KafkaChannelInitializer.java | 9 +++++++-- .../pulsar/handlers/kop/KafkaProtocolHandler.java | 11 +++++++++-- .../pulsar/handlers/kop/KafkaRequestHandler.java | 7 ++++--- .../pulsar/handlers/kop/EntryPublishTimeTest.java | 2 ++ .../pulsar/handlers/kop/KafkaApisTest.java | 2 ++ .../pulsar/handlers/kop/KafkaRequestHandlerTest.java | 2 ++ .../handlers/kop/KafkaTopicConsumerManagerTest.java | 2 ++ 7 files changed, 28 insertions(+), 7 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 4f2c10727f..ab86d9ea2c 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 @@ -26,6 +26,8 @@ import io.streamnative.pulsar.handlers.kop.utils.ssl.SSLUtils; import lombok.Getter; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; import org.eclipse.jetty.util.ssl.SslContextFactory; /** @@ -52,6 +54,7 @@ public class KafkaChannelInitializer extends ChannelInitializer { private final SslContextFactory.Server sslContextFactory; @Getter private final StatsLogger statsLogger; + private final MetadataCache localBrokerDataCache; public KafkaChannelInitializer(PulsarService pulsarService, KafkaServiceConfiguration kafkaConfig, @@ -60,7 +63,8 @@ public KafkaChannelInitializer(PulsarService pulsarService, AdminManager adminManager, boolean enableTLS, EndPoint advertisedEndPoint, - StatsLogger statsLogger) { + StatsLogger statsLogger, + MetadataCache localBrokerDataCache) { super(); this.pulsarService = pulsarService; this.kafkaConfig = kafkaConfig; @@ -70,6 +74,7 @@ public KafkaChannelInitializer(PulsarService pulsarService, this.enableTls = enableTLS; this.advertisedEndPoint = advertisedEndPoint; this.statsLogger = statsLogger; + this.localBrokerDataCache = localBrokerDataCache; if (enableTls) { sslContextFactory = SSLUtils.createSslContextFactory(kafkaConfig); @@ -88,7 +93,7 @@ protected void initChannel(SocketChannel ch) throws Exception { new LengthFieldBasedFrameDecoder(MAX_FRAME_LENGTH, 0, 4, 0, 4)); ch.pipeline().addLast("handler", new KafkaRequestHandler(pulsarService, kafkaConfig, - groupCoordinator, transactionCoordinator, adminManager, + groupCoordinator, transactionCoordinator, adminManager, localBrokerDataCache, 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 461d61f405..cb0b585901 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 @@ -68,6 +68,8 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.util.FutureUtil; +import org.apache.pulsar.metadata.api.MetadataCache; +import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; /** * Kafka Protocol Handler load and run by Pulsar Service. @@ -83,6 +85,7 @@ public class KafkaProtocolHandler implements ProtocolHandler { private PrometheusMetricsProvider statsProvider; private KopBrokerLookupManager kopBrokerLookupManager; private AdminManager adminManager = null; + private MetadataCache localBrokerDataCache; @Getter private KafkaServiceConfiguration kafkaConfig; @@ -259,6 +262,10 @@ public void start(BrokerService service) { KopVersion.getBuildHost(), KopVersion.getBuildTime()); + // Currently each time getMetadataCache() is called, a new MetadataCache instance will be created, even for + // the same type. So we must reuse the same MetadataCache to avoid creating a lot of instances. + localBrokerDataCache = brokerService.pulsar().getLocalMetadataStore().getMetadataCache(LocalBrokerData.class); + ZooKeeperUtils.tryCreatePath(brokerService.pulsar().getZkClient(), kafkaConfig.getGroupIdZooKeeperPath(), new byte[0]); @@ -353,13 +360,13 @@ public Map> newChannelIniti case SASL_PLAINTEXT: builder.put(endPoint.getInetAddress(), new KafkaChannelInitializer(brokerService.getPulsar(), kafkaConfig, groupCoordinator, transactionCoordinator, adminManager, false, - advertisedEndPoint, rootStatsLogger.scope(SERVER_SCOPE))); + advertisedEndPoint, rootStatsLogger.scope(SERVER_SCOPE), localBrokerDataCache)); break; case SSL: case SASL_SSL: builder.put(endPoint.getInetAddress(), new KafkaChannelInitializer(brokerService.getPulsar(), kafkaConfig, groupCoordinator, transactionCoordinator, adminManager, true, - advertisedEndPoint, rootStatsLogger.scope(SERVER_SCOPE))); + advertisedEndPoint, rootStatsLogger.scope(SERVER_SCOPE), localBrokerDataCache)); break; } }); 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 69ed843a66..606170e8b4 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 @@ -194,6 +194,7 @@ public class KafkaRequestHandler extends KafkaCommandDecoder { private final SaslAuthenticator authenticator; private final Authorizer authorizer; private final AdminManager adminManager; + private final MetadataCache localBrokerDataCache; private final Boolean tlsEnabled; private final EndPoint advertisedEndPoint; @@ -236,6 +237,7 @@ public KafkaRequestHandler(PulsarService pulsarService, GroupCoordinator groupCoordinator, TransactionCoordinator transactionCoordinator, AdminManager adminManager, + MetadataCache localBrokerDataCache, Boolean tlsEnabled, EndPoint advertisedEndPoint, StatsLogger statsLogger) throws Exception { @@ -256,6 +258,7 @@ public KafkaRequestHandler(PulsarService pulsarService, ? new SimpleAclAuthorizer(pulsarService) : null; this.adminManager = adminManager; + this.localBrokerDataCache = localBrokerDataCache; this.tlsEnabled = tlsEnabled; this.advertisedEndPoint = advertisedEndPoint; this.advertisedListeners = kafkaConfig.getKafkaAdvertisedListeners(); @@ -1970,10 +1973,8 @@ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { } // Get a list of ServiceLookupData for each matchBroker. - final MetadataCache metadataCache = pulsarService.getLocalMetadataStore() - .getMetadataCache(LocalBrokerData.class); List>> list = matchBrokers.stream() - .map(matchBroker -> metadataCache.get( + .map(matchBroker -> localBrokerDataCache.get( String.format("%s/%s", LoadManager.LOADBALANCE_BROKERS_ROOT, matchBroker))) .collect(Collectors.toList()); 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 99aaedd603..40f5d05ad0 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 @@ -24,6 +24,7 @@ import java.net.InetSocketAddress; import java.net.SocketAddress; import org.apache.pulsar.broker.protocol.ProtocolHandler; +import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.annotations.AfterMethod; @@ -59,6 +60,7 @@ protected void setup() throws Exception { groupCoordinator, transactionCoordinator, adminManager, + pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class), false, getPlainEndPoint(), NullStatsLogger.INSTANCE); 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 372bdc1dc4..e51b4f23c4 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 @@ -72,6 +72,7 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.broker.protocol.ProtocolHandler; import org.apache.pulsar.common.policies.data.RetentionPolicies; +import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Ignore; @@ -126,6 +127,7 @@ protected void setup() throws Exception { groupCoordinator, transactionCoordinator, adminManager, + pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class), false, getPlainEndPoint(), NullStatsLogger.INSTANCE); 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 117b24e3fe..6eb0d15f30 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 @@ -93,6 +93,7 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; +import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; import org.testng.Assert; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; @@ -147,6 +148,7 @@ protected void setup() throws Exception { groupCoordinator, transactionCoordinator, adminManager, + pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class), false, getPlainEndPoint(), NullStatsLogger.INSTANCE); 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 ad151b70c2..be6c6589a9 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 @@ -45,6 +45,7 @@ import org.apache.pulsar.broker.service.persistent.PersistentTopic; 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.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -75,6 +76,7 @@ protected void setup() throws Exception { groupCoordinator, transactionCoordinator, adminManager, + pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class), false, getPlainEndPoint(), NullStatsLogger.INSTANCE);