From ececb04b3cb48ad2e48fc0622e73332a51f277a6 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 20 Oct 2021 16:34:22 +0800 Subject: [PATCH 1/4] Use MetadataStoreCacheLoader replace MetadataStoreCache in KopBrokerLookupManager --- .../handlers/kop/KafkaProtocolHandler.java | 10 +- .../handlers/kop/KafkaRequestHandler.java | 9 -- .../kop/KafkaServiceConfiguration.java | 6 + .../handlers/kop/KopBrokerLookupManager.java | 129 ++++++------------ 4 files changed, 59 insertions(+), 95 deletions(-) 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 c9a0b26bb5..8ab151044f 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 @@ -448,8 +448,14 @@ public void start(BrokerService service) { offsetTopicClient = new SystemTopicClient(brokerService.pulsar(), kafkaConfig); txnTopicClient = new SystemTopicClient(brokerService.pulsar(), kafkaConfig); - kopBrokerLookupManager = new KopBrokerLookupManager( - brokerService.getPulsar(), kafkaConfig.getKafkaAdvertisedListeners()); + try { + kopBrokerLookupManager = new KopBrokerLookupManager( + brokerService.getPulsar(), kafkaConfig.getKafkaAdvertisedListeners(), + kafkaConfig.getBrokerLookupTimeoutSeconds()); + } catch (Exception ex) { + log.error("Failed to get kopBrokerLookupManager", ex); + throw new IllegalStateException(ex); + } brokerService.pulsar() .getNamespaceService() 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 cad8b94254..7995157f84 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 @@ -186,7 +186,6 @@ import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.Murmur3_32Hash; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; -import org.apache.pulsar.policies.data.loadbalancer.ServiceLookupData; /** * This class contains all the request handling methods. @@ -2658,14 +2657,6 @@ static AbstractResponse failedResponse(KafkaHeaderAndRequest requestHar, Throwab return requestHar.getRequest().getErrorResponse(((Integer) THROTTLE_TIME_MS.defaultValue), e); } - // whether a ServiceLookupData contains wanted address. - static boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { - return (data.getPulsarServiceUrl() != null && data.getPulsarServiceUrl().contains(hostAndPort)) - || (data.getPulsarServiceUrlTls() != null && data.getPulsarServiceUrlTls().contains(hostAndPort)) - || (data.getWebServiceUrl() != null && data.getWebServiceUrl().contains(hostAndPort)) - || (data.getWebServiceUrlTls() != null && data.getWebServiceUrlTls().contains(hostAndPort)); - } - private static MemoryRecords validateRecords(short version, TopicPartition topicPartition, MemoryRecords records) { if (version >= 3) { Iterator iterator = records.batches().iterator(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java index 94da33e965..c055f3c2f3 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java @@ -231,6 +231,12 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { ) private int failedAuthenticationDelayMs = 300; + @FieldContext( + category = CATEGORY_KOP, + doc = "The timeout for broker lookups (in seconds)" + ) + private int brokerLookupTimeoutSeconds = 30; + // Kafka SSL configs @FieldContext( category = CATEGORY_KOP_SSL, diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java index e5d5296948..811a3324ca 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java @@ -13,25 +13,21 @@ */ package io.streamnative.pulsar.handlers.kop; -import static io.streamnative.pulsar.handlers.kop.KafkaRequestHandler.lookupDataContainsAddress; -import com.google.common.collect.Lists; +import java.io.IOException; import java.net.InetSocketAddress; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.stream.Collectors; import javax.annotation.Nullable; import lombok.NonNull; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.broker.PulsarService; -import org.apache.pulsar.broker.loadbalance.LoadManager; +import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; import org.apache.pulsar.common.naming.TopicName; -import org.apache.pulsar.common.util.FutureUtil; -import org.apache.pulsar.metadata.api.MetadataCache; -import org.apache.pulsar.policies.data.loadbalancer.LocalBrokerData; +import org.apache.pulsar.policies.data.loadbalancer.LoadManagerReport; import org.apache.pulsar.policies.data.loadbalancer.ServiceLookupData; @@ -41,10 +37,9 @@ @Slf4j public class KopBrokerLookupManager { - private final PulsarService pulsarService; private final String advertisedListeners; private final LookupClient lookupClient; - private final MetadataCache localBrokerDataCache; + private final MetadataStoreCacheLoader metadataStoreCacheLoader; private final AtomicBoolean closed = new AtomicBoolean(false); @@ -54,11 +49,12 @@ public class KopBrokerLookupManager { public static final ConcurrentHashMap>> KOP_ADDRESS_CACHE = new ConcurrentHashMap<>(); - public KopBrokerLookupManager(PulsarService pulsarService, String advertisedListeners) { - this.pulsarService = pulsarService; + public KopBrokerLookupManager( + PulsarService pulsarService, String advertisedListeners, int brokerLookupTimeoutSeconds) throws Exception { this.advertisedListeners = advertisedListeners; this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); - this.localBrokerDataCache = pulsarService.getLocalMetadataStore().getMetadataCache(LocalBrokerData.class); + this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), + brokerLookupTimeoutSeconds); } public CompletableFuture> findBroker(@NonNull TopicName topic, @@ -184,83 +180,43 @@ private CompletableFuture> getProtocolDataToAdvertise( return returnFuture; } - // advertised data is write in /loadbalance/brokers/advertisedAddress:webServicePort + // advertised data is written in /loadbalance/brokers/advertisedAddress:webServicePort // here we get the broker url, need to find related webServiceUrl. - pulsarService.getPulsarResources() - .getDynamicConfigResources() - .getChildrenAsync(LoadManager.LOADBALANCE_BROKERS_ROOT) - .whenComplete((set, throwable) -> { - if (throwable != null) { - log.error("Error in getChildrenAsync(zk://loadbalance) for {}", pulsarAddress, throwable); - returnFuture.complete(Optional.empty()); - return; - } - - String hostAndPort = pulsarAddress.getHostName() + ":" + pulsarAddress.getPort(); - List matchBrokers = Lists.newArrayList(); - // match host part of url - for (String activeBroker : set) { - if (activeBroker.startsWith(pulsarAddress.getHostName() + ":")) { - matchBrokers.add(activeBroker); - } - } - - if (matchBrokers.isEmpty()) { - log.error("No node for broker {} under zk://loadbalance", pulsarAddress); - returnFuture.complete(Optional.empty()); - removeTopicManagerCache(topic.toString()); - return; - } + List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); + if (log.isDebugEnabled()) { + availableBrokers.forEach(loadManagerReport -> + log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, " + + "pulsarUrlTls: {}, webUrl: {}, webUrlTls: {} kafka: {}", + topic, + loadManagerReport.getPulsarServiceUrl(), + loadManagerReport.getPulsarServiceUrlTls(), + loadManagerReport.getWebServiceUrl(), + loadManagerReport.getWebServiceUrlTls(), + loadManagerReport.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME))); + } - // Get a list of ServiceLookupData for each matchBroker. - List>> list = matchBrokers.stream() - .map(matchBroker -> localBrokerDataCache.get( - String.format("%s/%s", LoadManager.LOADBALANCE_BROKERS_ROOT, matchBroker))) - .collect(Collectors.toList()); - - FutureUtil.waitForAll(list).whenComplete((ignore, th) -> { - if (th != null) { - log.error("Error in getDataAsync() for {}", pulsarAddress, th); - returnFuture.complete(Optional.empty()); - removeTopicManagerCache(topic.toString()); - return; - } - - try { - for (CompletableFuture> lookupData : list) { - ServiceLookupData data = lookupData.get().get(); - if (log.isDebugEnabled()) { - log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, " - + "pulsarUrlTls: {}, webUrl: {}, webUrlTls: {} kafka: {}", - topic, - data.getPulsarServiceUrl(), - data.getPulsarServiceUrlTls(), - data.getWebServiceUrl(), - data.getWebServiceUrlTls(), - data.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); - } - - if (lookupDataContainsAddress(data, hostAndPort)) { - KOP_ADDRESS_CACHE.put(topic.toString(), returnFuture); - returnFuture.complete(data.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); - return; - } - } - } catch (Exception e) { - log.error("Error in {} lookupFuture get: ", pulsarAddress, e); - returnFuture.complete(Optional.empty()); - removeTopicManagerCache(topic.toString()); - return; - } - - // no matching lookup data in all matchBrokers. - log.error("Not able to search {} in all child of zk://loadbalance", pulsarAddress); - returnFuture.complete(Optional.empty()); - }); - }); + String hostAndPort = pulsarAddress.getHostName() + ":" + pulsarAddress.getPort(); + Optional serviceLookupData = availableBrokers.stream() + .filter(loadManagerReport -> lookupDataContainsAddress(loadManagerReport, hostAndPort)).findAny(); + if (serviceLookupData.isPresent()) { + KOP_ADDRESS_CACHE.put(topic.toString(), returnFuture); + returnFuture.complete(serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); + } else { + log.error("No node for broker {} under loadBalance", pulsarAddress); + returnFuture.complete(Optional.empty()); + removeTopicManagerCache(topic.toString()); + } return returnFuture; } + // whether a ServiceLookupData contains wanted address. + private boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { + return (data.getPulsarServiceUrl() != null && data.getPulsarServiceUrl().contains(hostAndPort)) + || (data.getPulsarServiceUrlTls() != null && data.getPulsarServiceUrlTls().contains(hostAndPort)) + || (data.getWebServiceUrl() != null && data.getWebServiceUrl().contains(hostAndPort)) + || (data.getWebServiceUrlTls() != null && data.getWebServiceUrlTls().contains(hostAndPort)); + } + public static void removeTopicManagerCache(String topicName) { LOOKUP_CACHE.remove(topicName); KOP_ADDRESS_CACHE.remove(topicName); @@ -279,6 +235,11 @@ public void close() { return; } clear(); + try { + metadataStoreCacheLoader.close(); + } catch (IOException e) { + log.error("Close metadataStoreCacheLoader failed.", e); + } } } From 4a757127c349106e47b287690a5c44ec35bab42c Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Wed, 20 Oct 2021 17:43:43 +0800 Subject: [PATCH 2/4] Change brokerLookupTimeout to milliseconds --- docs/configuration.md | 1 + .../pulsar/handlers/kop/KafkaProtocolHandler.java | 2 +- .../pulsar/handlers/kop/KafkaServiceConfiguration.java | 4 ++-- .../pulsar/handlers/kop/KopBrokerLookupManager.java | 4 ++-- 4 files changed, 6 insertions(+), 5 deletions(-) diff --git a/docs/configuration.md b/docs/configuration.md index c4d1865d95..3a67d81037 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -80,6 +80,7 @@ This section lists configurations that may affect the performance. | requestTimeoutMs | Limit the timeout in milliseconds for request, like `request.timeout.ms` in Kafka client.
If a request was not processed in the timeout, KoP would return an error response to client. | 30000 | | connectionMaxIdleMs | The idle connection timeout in milliseconds. If the idle connection timeout (such as `connections.max.idle.ms` used in the Kafka server) is reached, the server handler will close this idle connection.
**Note**: If it is set to `-1`, it indicates that the idle connection timeout is disabled. | 600000 | | failedAuthenticationDelayMs | Connection close delay on failed authentication: this is the time (in milliseconds) by which connection close will be delayed on authentication failure, like `connection.failed.authentication.delay.ms` in Kafka server. | 300 | +| brokerLookupTimeoutMs | The timeout for broker lookups (in milliseconds). | 30000 | > **NOTE** > 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 8ab151044f..41b170f13e 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 @@ -451,7 +451,7 @@ public void start(BrokerService service) { try { kopBrokerLookupManager = new KopBrokerLookupManager( brokerService.getPulsar(), kafkaConfig.getKafkaAdvertisedListeners(), - kafkaConfig.getBrokerLookupTimeoutSeconds()); + kafkaConfig.getBrokerLookupTimeoutMs()); } catch (Exception ex) { log.error("Failed to get kopBrokerLookupManager", ex); throw new IllegalStateException(ex); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java index c055f3c2f3..f42df934e4 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java @@ -233,9 +233,9 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { @FieldContext( category = CATEGORY_KOP, - doc = "The timeout for broker lookups (in seconds)" + doc = "The timeout for broker lookups (in milliseconds)" ) - private int brokerLookupTimeoutSeconds = 30; + private int brokerLookupTimeoutMs = 30_000; // Kafka SSL configs @FieldContext( diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java index 811a3324ca..a63af88e48 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java @@ -50,11 +50,11 @@ public class KopBrokerLookupManager { KOP_ADDRESS_CACHE = new ConcurrentHashMap<>(); public KopBrokerLookupManager( - PulsarService pulsarService, String advertisedListeners, int brokerLookupTimeoutSeconds) throws Exception { + PulsarService pulsarService, String advertisedListeners, int brokerLookupTimeoutMs) throws Exception { this.advertisedListeners = advertisedListeners; this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), - brokerLookupTimeoutSeconds); + brokerLookupTimeoutMs); } public CompletableFuture> findBroker(@NonNull TopicName topic, From 2cf90476013dd5b575a07ef01fd9c7f58bbfa0bc Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 21 Oct 2021 11:18:20 +0800 Subject: [PATCH 3/4] Use endsWith to check contains address --- .../pulsar/handlers/kop/KopBrokerLookupManager.java | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java index a63af88e48..b6d2ac1659 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java @@ -24,6 +24,7 @@ import javax.annotation.Nullable; import lombok.NonNull; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; import org.apache.pulsar.common.naming.TopicName; @@ -211,10 +212,8 @@ private CompletableFuture> getProtocolDataToAdvertise( // whether a ServiceLookupData contains wanted address. private boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { - return (data.getPulsarServiceUrl() != null && data.getPulsarServiceUrl().contains(hostAndPort)) - || (data.getPulsarServiceUrlTls() != null && data.getPulsarServiceUrlTls().contains(hostAndPort)) - || (data.getWebServiceUrl() != null && data.getWebServiceUrl().contains(hostAndPort)) - || (data.getWebServiceUrlTls() != null && data.getWebServiceUrlTls().contains(hostAndPort)); + return StringUtils.endsWith(data.getPulsarServiceUrl(), hostAndPort) + || StringUtils.endsWith(data.getPulsarServiceUrlTls(), hostAndPort); } public static void removeTopicManagerCache(String topicName) { From 675a06bb0cb019009ba884e5b47bee98e1214a18 Mon Sep 17 00:00:00 2001 From: Demogorgon314 Date: Thu, 21 Oct 2021 17:46:26 +0800 Subject: [PATCH 4/4] Address reviewer's comment --- .../pulsar/handlers/kop/KopBrokerLookupManager.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java index b6d2ac1659..f32f987992 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KopBrokerLookupManager.java @@ -181,8 +181,6 @@ private CompletableFuture> getProtocolDataToAdvertise( return returnFuture; } - // advertised data is written in /loadbalance/brokers/advertisedAddress:webServicePort - // here we get the broker url, need to find related webServiceUrl. List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); if (log.isDebugEnabled()) { availableBrokers.forEach(loadManagerReport -> @@ -204,14 +202,14 @@ private CompletableFuture> getProtocolDataToAdvertise( returnFuture.complete(serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME)); } else { log.error("No node for broker {} under loadBalance", pulsarAddress); - returnFuture.complete(Optional.empty()); removeTopicManagerCache(topic.toString()); + returnFuture.complete(Optional.empty()); } return returnFuture; } // whether a ServiceLookupData contains wanted address. - private boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { + private static boolean lookupDataContainsAddress(ServiceLookupData data, String hostAndPort) { return StringUtils.endsWith(data.getPulsarServiceUrl(), hostAndPort) || StringUtils.endsWith(data.getPulsarServiceUrlTls(), hostAndPort); }