Skip to content
This repository was archived by the owner on Jan 24, 2024. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand All @@ -52,6 +54,7 @@ public class KafkaChannelInitializer extends ChannelInitializer<SocketChannel> {
private final SslContextFactory.Server sslContextFactory;
@Getter
private final StatsLogger statsLogger;
private final MetadataCache<LocalBrokerData> localBrokerDataCache;

public KafkaChannelInitializer(PulsarService pulsarService,
KafkaServiceConfiguration kafkaConfig,
Expand All @@ -60,7 +63,8 @@ public KafkaChannelInitializer(PulsarService pulsarService,
AdminManager adminManager,
boolean enableTLS,
EndPoint advertisedEndPoint,
StatsLogger statsLogger) {
StatsLogger statsLogger,
MetadataCache<LocalBrokerData> localBrokerDataCache) {
super();
this.pulsarService = pulsarService;
this.kafkaConfig = kafkaConfig;
Expand All @@ -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);
Expand All @@ -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));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -83,6 +85,7 @@ public class KafkaProtocolHandler implements ProtocolHandler {
private PrometheusMetricsProvider statsProvider;
private KopBrokerLookupManager kopBrokerLookupManager;
private AdminManager adminManager = null;
private MetadataCache<LocalBrokerData> localBrokerDataCache;

@Getter
private KafkaServiceConfiguration kafkaConfig;
Expand Down Expand Up @@ -259,6 +262,10 @@ public void start(BrokerService service) {
KopVersion.getBuildHost(),
KopVersion.getBuildTime());

// Currently each time getMetadataCache() is called, a new MetadataCache<T> instance will be created, even for
// the same type. So we must reuse the same MetadataCache<LocalBrokerData> to avoid creating a lot of instances.
localBrokerDataCache = brokerService.pulsar().getLocalMetadataStore().getMetadataCache(LocalBrokerData.class);

ZooKeeperUtils.tryCreatePath(brokerService.pulsar().getZkClient(),
kafkaConfig.getGroupIdZooKeeperPath(), new byte[0]);

Expand Down Expand Up @@ -353,13 +360,13 @@ public Map<InetSocketAddress, ChannelInitializer<SocketChannel>> 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;
}
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,7 @@ public class KafkaRequestHandler extends KafkaCommandDecoder {
private final SaslAuthenticator authenticator;
private final Authorizer authorizer;
private final AdminManager adminManager;
private final MetadataCache<LocalBrokerData> localBrokerDataCache;

private final Boolean tlsEnabled;
private final EndPoint advertisedEndPoint;
Expand Down Expand Up @@ -236,6 +237,7 @@ public KafkaRequestHandler(PulsarService pulsarService,
GroupCoordinator groupCoordinator,
TransactionCoordinator transactionCoordinator,
AdminManager adminManager,
MetadataCache<LocalBrokerData> localBrokerDataCache,
Boolean tlsEnabled,
EndPoint advertisedEndPoint,
StatsLogger statsLogger) throws Exception {
Expand All @@ -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();
Expand Down Expand Up @@ -1970,10 +1973,8 @@ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) {
}

// Get a list of ServiceLookupData for each matchBroker.
final MetadataCache<LocalBrokerData> metadataCache = pulsarService.getLocalMetadataStore()
.getMetadataCache(LocalBrokerData.class);
List<CompletableFuture<Optional<LocalBrokerData>>> list = matchBrokers.stream()
.map(matchBroker -> metadataCache.get(
.map(matchBroker -> localBrokerDataCache.get(
String.format("%s/%s", LoadManager.LOADBALANCE_BROKERS_ROOT, matchBroker)))
.collect(Collectors.toList());

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -59,6 +60,7 @@ protected void setup() throws Exception {
groupCoordinator,
transactionCoordinator,
adminManager,
pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class),
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -126,6 +127,7 @@ protected void setup() throws Exception {
groupCoordinator,
transactionCoordinator,
adminManager,
pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class),
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -147,6 +148,7 @@ protected void setup() throws Exception {
groupCoordinator,
transactionCoordinator,
adminManager,
pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class),
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -75,6 +76,7 @@ protected void setup() throws Exception {
groupCoordinator,
transactionCoordinator,
adminManager,
pulsar.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class),
false,
getPlainEndPoint(),
NullStatsLogger.INSTANCE);
Expand Down