From 8eff8dd99b8aafc5e0d0a29890c6131cc1ab65db Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 15:23:22 +0800 Subject: [PATCH 1/8] Correct topic lookup implementation --- .../pulsar/handlers/kop/EndPoint.java | 13 ++ .../handlers/kop/KopBrokerLookupManager.java | 159 +++++++----------- .../pulsar/handlers/kop/LookupClient.java | 85 +--------- 3 files changed, 73 insertions(+), 184 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java index 3118f6d147..b962f32c0c 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java @@ -25,6 +25,7 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import lombok.Getter; +import lombok.NonNull; import org.apache.commons.lang3.StringUtils; import org.apache.kafka.common.security.auth.SecurityProtocol; @@ -115,6 +116,18 @@ public static Map parseListeners(final String listeners, final return parseListeners(listeners, parseProtocolMap(protocolMapString)); } + public static String findListener(@NonNull final String listeners, final String name) { + if (name == null) { + return null; + } + for (String listener : listeners.split(END_POINT_SEPARATOR)) { + if (listener.contains(":") && listener.substring(0, listener.indexOf(":")).equals(name)) { + return listener; + } + } + return null; + } + public static EndPoint getPlainTextEndPoint(final String listeners) { for (String listener : listeners.split(END_POINT_SEPARATOR)) { if (listener.startsWith(SecurityProtocol.PLAINTEXT.name()) 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 ccc4c894ca..038d7779ef 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 @@ -17,6 +17,7 @@ import java.io.IOException; import java.net.InetSocketAddress; import java.util.List; +import java.util.Map; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; @@ -25,6 +26,7 @@ import lombok.NonNull; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; +import org.apache.kafka.common.security.auth.SecurityProtocol; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; import org.apache.pulsar.common.naming.TopicName; @@ -38,21 +40,21 @@ @Slf4j public class KopBrokerLookupManager { - private final String advertisedListeners; private final LookupClient lookupClient; + private final Map protocolMap; private final MetadataStoreCacheLoader metadataStoreCacheLoader; private final AtomicBoolean closed = new AtomicBoolean(false); - public static final ConcurrentHashMap>> + public static final ConcurrentHashMap> LOOKUP_CACHE = new ConcurrentHashMap<>(); - public static final ConcurrentHashMap>> - KOP_ADDRESS_CACHE = new ConcurrentHashMap<>(); + public static final ConcurrentHashMap KOP_ADDRESS_CACHE = new ConcurrentHashMap<>(); public KopBrokerLookupManager(KafkaServiceConfiguration conf, PulsarService pulsarService) throws Exception { - this.advertisedListeners = conf.getKafkaAdvertisedListeners(); this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); + // TODO: change it to getKafkaAdvertisedProtocolMap + this.protocolMap = EndPoint.parseProtocolMap(conf.getKafkaProtocolMap()); this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), conf.getBrokerLookupTimeoutMs()); } @@ -62,51 +64,29 @@ public CompletableFuture> findBroker(@NonNull TopicN if (log.isDebugEnabled()) { log.debug("Handle Lookup for topic {}", topic); } - CompletableFuture> returnFuture = new CompletableFuture<>(); - - getTopicBroker(topic.toString(), - advertisedEndPoint != null && advertisedEndPoint.isValidInProtocolMap() - ? advertisedEndPoint.getListenerName() : null) - .thenApply(address -> getProtocolDataToAdvertise(address, topic, advertisedEndPoint)) - .thenAccept(kopAddressFuture -> kopAddressFuture.thenAccept(listenersOptional -> { - if (!listenersOptional.isPresent()) { - log.error("Not get advertise data for Kafka topic:{}.", topic); + + return getTopicBroker(topic.toString()) + .thenApply(internalListenerAddress -> { + final String listener = getAdvertisedListener( + internalListenerAddress, topic, advertisedEndPoint); + if (listener == null) { removeTopicManagerCache(topic.toString()); - returnFuture.complete(Optional.empty()); - return; + return Optional.empty(); } - // It's the `kafkaAdvertisedListeners` config that's written to ZK - final String listeners = listenersOptional.get(); - final EndPoint endPoint = - (advertisedEndPoint != null && advertisedEndPoint.isTlsEnabled() - ? EndPoint.getSslEndPoint(listeners) : EndPoint.getPlainTextEndPoint(listeners)); - - if (log.isDebugEnabled()) { - log.debug("Found broker localListeners: {} for topicName: {}, " - + "localListeners: {}, found Listeners: {}", - listeners, topic, advertisedListeners, listeners); + try { + final EndPoint endPoint = new EndPoint(listener, protocolMap); + return Optional.of(new InetSocketAddress(endPoint.getHostname(), endPoint.getPort())); + } catch (IllegalStateException e) { + log.error("Failed to create EndPoint from '{}': {}", listener, e.getMessage()); + return Optional.empty(); } - - // here we found topic broker: broker2, but this is in broker1, - // how to clean the lookup cache? - if (!advertisedListeners.contains(endPoint.getOriginalListener())) { - removeTopicManagerCache(topic.toString()); - } - returnFuture.complete(Optional.of(endPoint.getInetAddress())); - })).exceptionally(throwable -> { - log.error("Not get advertise data for Kafka topic:{}. throwable: [{}]", - topic, throwable.getMessage()); - removeTopicManagerCache(topic.toString()); - returnFuture.complete(Optional.empty()); - return null; }); - return returnFuture; } // call pulsarclient.lookup.getbroker to get and own a topic. // when error happens, the returned future will complete with null. - public CompletableFuture getTopicBroker(String topicName, String listenerName) { + public CompletableFuture getTopicBroker(String topicName) { if (closed.get()) { if (log.isDebugEnabled()) { log.debug("Return null for getTopicBroker({}) since channel closing", topicName); @@ -114,71 +94,36 @@ public CompletableFuture getTopicBroker(String topicName, Str return CompletableFuture.completedFuture(null); } - ConcurrentHashMap> topicLookupCache = - LOOKUP_CACHE.computeIfAbsent(topicName, t-> { - if (log.isDebugEnabled()) { - log.debug("Topic {} not in Lookup_cache, call lookupBroker", topicName); - } - ConcurrentHashMap> cache = new ConcurrentHashMap<>(); - cache.put(listenerName == null ? "" : listenerName, lookupBroker(topicName, listenerName)); - return cache; - }); - - return topicLookupCache.computeIfAbsent(listenerName == null ? "" : listenerName, t-> { - if (log.isDebugEnabled()) { - log.debug("Topic {} not in Lookup_cache, call lookupBroker", topicName); - } - return lookupBroker(topicName, listenerName); - }); + return LOOKUP_CACHE.computeIfAbsent(topicName, this::lookupBroker); } - private CompletableFuture lookupBroker(final String topic, String listenerName) { + private CompletableFuture lookupBroker(final String topic) { if (closed.get()) { if (log.isDebugEnabled()) { log.debug("Return null for getTopic({}) since channel closing", topic); } return CompletableFuture.completedFuture(null); } - return lookupClient.getBrokerAddress(TopicName.get(topic), listenerName); + return lookupClient.getBrokerAddress(TopicName.get(topic)); } - private CompletableFuture> getProtocolDataToAdvertise( - InetSocketAddress pulsarAddress, TopicName topic, @Nullable EndPoint advertisedEndPoint) { - CompletableFuture> returnFuture = new CompletableFuture<>(); - - if (pulsarAddress == null) { - log.error("[{}] failed get pulsar address, returned null.", topic.toString()); - - // getTopicBroker returns null. topic should be removed from LookupCache. - removeTopicManagerCache(topic.toString()); - - returnFuture.complete(Optional.empty()); - return returnFuture; + private String getAdvertisedListener(InetSocketAddress internalListenerAddress, + TopicName topic, + @Nullable EndPoint advertisedEndPoint) { + if (internalListenerAddress == null) { + log.error("[{}] failed get pulsar address, returned null.", topic); + return null; } if (log.isDebugEnabled()) { - log.debug("Found broker for topic {} puslarAddress: {}", - topic, pulsarAddress); + log.debug("[{}] Found broker's internal listener address: {}", + topic, internalListenerAddress); } - // get kop address from cache to prevent query zk each time. - final CompletableFuture> future = KOP_ADDRESS_CACHE.get(topic.toString()); - if (future != null) { - return future; - } - - if (advertisedEndPoint != null && advertisedEndPoint.isValidInProtocolMap()) { - // if kafkaProtocolMap is set, the lookup result is the advertised address - String kafkaAdvertisedAddress = String.format("%s://%s:%s", advertisedEndPoint.getSecurityProtocol().name, - pulsarAddress.getHostName(), pulsarAddress.getPort()); - KOP_ADDRESS_CACHE.put(topic.toString(), returnFuture); - returnFuture.complete(Optional.ofNullable(kafkaAdvertisedAddress)); - if (log.isDebugEnabled()) { - log.debug("{} get kafka Advertised Address through kafkaListenerName: {}", - topic, pulsarAddress); - } - return returnFuture; - } + final String protocolDataToAdvertise = KOP_ADDRESS_CACHE.get(topic.toString()); + if (protocolDataToAdvertise != null) { + return protocolDataToAdvertise; + } // else: either the key doesn't exist or a null value was put List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); if (log.isDebugEnabled()) { @@ -193,18 +138,32 @@ private CompletableFuture> getProtocolDataToAdvertise( loadManagerReport.getProtocol(KafkaProtocolHandler.PROTOCOL_NAME))); } - String hostAndPort = pulsarAddress.getHostName() + ":" + pulsarAddress.getPort(); - Optional serviceLookupData = availableBrokers.stream() + final String hostAndPort = internalListenerAddress.getHostName() + ":" + internalListenerAddress.getPort(); + final 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)); + if (!serviceLookupData.isPresent()) { + log.error("No node for broker {} under loadBalance", internalListenerAddress); + return null; + } + + final String kafkaAdvertisedListeners = + serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME).orElse(null); + if (kafkaAdvertisedListeners == null) { + log.error("No kafkaAdvertisedListeners found in broker {}", internalListenerAddress); + return null; + } + + if (log.isDebugEnabled()) { + log.debug("Found kafkaAdvertisedListeners: {}", kafkaAdvertisedListeners); + } + + final String listenerName = (advertisedEndPoint != null) ? advertisedEndPoint.getListenerName() : null; + if (listenerName == null) { + // Treat the whole kafkaAdvertisedListeners as the listener + return kafkaAdvertisedListeners; } else { - log.error("No node for broker {} under loadBalance", pulsarAddress); - removeTopicManagerCache(topic.toString()); - returnFuture.complete(Optional.empty()); + return EndPoint.findListener(kafkaAdvertisedListeners, listenerName); } - return returnFuture; } // whether a ServiceLookupData contains wanted address. diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/LookupClient.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/LookupClient.java index 4ce171f7cb..1e69a99c0a 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/LookupClient.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/LookupClient.java @@ -14,22 +14,10 @@ package io.streamnative.pulsar.handlers.kop; import java.net.InetSocketAddress; -import java.net.URI; -import java.net.URISyntaxException; -import java.util.Map; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ConcurrentHashMap; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.tuple.Pair; -import org.apache.kafka.common.security.auth.SecurityProtocol; import org.apache.pulsar.broker.PulsarService; -import org.apache.pulsar.broker.lookup.LookupResult; -import org.apache.pulsar.broker.namespace.LookupOptions; -import org.apache.pulsar.broker.namespace.NamespaceService; -import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.client.impl.ClientCnx; -import org.apache.pulsar.client.impl.PulsarClientImpl; -import org.apache.pulsar.common.api.proto.ServerError; import org.apache.pulsar.common.naming.TopicName; /** @@ -38,88 +26,17 @@ @Slf4j public class LookupClient extends AbstractPulsarClient { - private final NamespaceService namespaceService; - - private ConcurrentHashMap pulsarClientMap; - public LookupClient(final PulsarService pulsarService, final KafkaServiceConfiguration kafkaConfig) { super(createPulsarClient(pulsarService, kafkaConfig, conf -> {})); - namespaceService = pulsarService.getNamespaceService(); - try { - pulsarClientMap = createPulsarClientMap(pulsarService, kafkaConfig); - } catch (PulsarClientException e) { - log.error("Failed to create PulsarClient", e); - throw new IllegalStateException(e); - } } public LookupClient(final PulsarService pulsarService) { super(createPulsarClient(pulsarService)); log.warn("This constructor should not be called, it's only called " + "when the PulsarService doesn't exist in KafkaProtocolHandlers.LOOKUP_CLIENT_UP"); - namespaceService = pulsarService.getNamespaceService(); } public CompletableFuture getBrokerAddress(final TopicName topicName) { - return getBrokerAddress(topicName, null); - } - - public CompletableFuture getBrokerAddress(final TopicName topicName, String listenerName) { - // First try to use NamespaceService to find the broker directly. - final LookupOptions options = LookupOptions.builder() - .authoritative(false) - .advertisedListenerName(listenerName) - .loadTopicsInBundle(true) - .build(); - return namespaceService.getBrokerServiceUrlAsync(topicName, options).thenCompose(optLookupResult -> { - if (log.isDebugEnabled()) { - log.debug("[{}] Lookup result {}", topicName.toString(), optLookupResult); - } - if (!optLookupResult.isPresent()) { - return getFailedAddressFuture(ClientCnx.getPulsarClientException( - ServerError.ServiceNotReady, - "No broker was available to own " + topicName)); - } - final LookupResult lookupResult = optLookupResult.get(); - if (lookupResult.isRedirect()) { - // Kafka client can't process redirect field, so here we fallback to PulsarClient - return pulsarClientMap.getOrDefault(listenerName == null ? "" : listenerName, getPulsarClient()). - getLookup().getBroker(topicName).thenApply(Pair::getLeft); - } else { - return getAddressFutureFromBrokerUrl(lookupResult.getLookupData().getBrokerUrl()); - } - }); - } - - private ConcurrentHashMap createPulsarClientMap( - PulsarService pulsarService, KafkaServiceConfiguration kafkaConfig) throws PulsarClientException { - ConcurrentHashMap pulsarClientMap = new ConcurrentHashMap<>(); - final Map protocolMap = EndPoint.parseProtocolMap(kafkaConfig.getKafkaProtocolMap()); - if (protocolMap.isEmpty()) { - pulsarClientMap.put("", getPulsarClient()); - } else { - for (Map.Entry entry : protocolMap.entrySet()) { - pulsarClientMap.put(entry.getKey(), createPulsarClient( - pulsarService, kafkaConfig, conf -> conf.setListenerName(entry.getKey()))); - } - } - return pulsarClientMap; - } - - private static CompletableFuture getFailedAddressFuture(final Throwable throwable) { - final CompletableFuture future = new CompletableFuture<>(); - future.completeExceptionally(throwable); - return future; - } - - private static CompletableFuture getAddressFutureFromBrokerUrl(final String brokerUrl) { - final CompletableFuture future = new CompletableFuture<>(); - try { - final URI uri = new URI(brokerUrl); - future.complete(InetSocketAddress.createUnresolved(uri.getHost(), uri.getPort())); - } catch (URISyntaxException e) { - future.completeExceptionally(e); - } - return future; + return getPulsarClient().getLookup().getBroker(topicName).thenApply(Pair::getLeft); } } From b65ea62dbedfcb6179f42a3c4790a7ddfd5255c0 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 17:14:34 +0800 Subject: [PATCH 2/8] Remove KOP_ADDRESS_CACHE --- .../pulsar/handlers/kop/EndPoint.java | 25 ++++-- .../handlers/kop/KafkaRequestHandler.java | 2 +- .../handlers/kop/KopBrokerLookupManager.java | 79 +++++++------------ .../TransactionMarkerChannelManager.java | 3 +- .../handlers/kop/CacheInvalidatorTest.java | 3 - 5 files changed, 49 insertions(+), 63 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java index b962f32c0c..98ce8209a8 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/EndPoint.java @@ -25,7 +25,6 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import lombok.Getter; -import lombok.NonNull; import org.apache.commons.lang3.StringUtils; import org.apache.kafka.common.security.auth.SecurityProtocol; @@ -91,6 +90,18 @@ public InetSocketAddress getInetAddress() { return new InetSocketAddress(hostname, port); } + // listeners must be enable to be split into at least 1 token + private static String[] getListenerArray(final String listeners) { + if (StringUtils.isEmpty(listeners)) { + throw new IllegalStateException("listeners is empty"); + } + final String[] listenerArray = listeners.split(END_POINT_SEPARATOR); + if (listenerArray.length == 0) { + throw new IllegalStateException(listeners + " is split into 0 tokens by " + END_POINT_SEPARATOR); + } + return listenerArray; + } + @VisibleForTesting public static Map parseListeners(final String listeners) { return parseListeners(listeners, ""); @@ -99,7 +110,7 @@ public static Map parseListeners(final String listeners) { private static Map parseListeners(final String listeners, final Map protocolMap) { final Map endPointMap = new HashMap<>(); - for (String listener : listeners.split(END_POINT_SEPARATOR)) { + for (String listener : getListenerArray(listeners)) { final EndPoint endPoint = new EndPoint(listener, protocolMap); if (endPointMap.containsKey(endPoint.listenerName)) { throw new IllegalStateException( @@ -116,16 +127,20 @@ public static Map parseListeners(final String listeners, final return parseListeners(listeners, parseProtocolMap(protocolMapString)); } - public static String findListener(@NonNull final String listeners, final String name) { + public static String findListener(final String listeners, final String name) { if (name == null) { return null; } - for (String listener : listeners.split(END_POINT_SEPARATOR)) { + for (String listener : getListenerArray(listeners)) { if (listener.contains(":") && listener.substring(0, listener.indexOf(":")).equals(name)) { return listener; } } - return null; + throw new IllegalStateException("listener \"" + name + "\" doesn't exist in " + listeners); + } + + public static String findFirstListener(String listeners) { + return getListenerArray(listeners)[0]; } public static EndPoint getPlainTextEndPoint(final String listeners) { 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 2b5a7e7f2a..3bc857b73b 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 @@ -2481,7 +2481,7 @@ public CompletableFuture findBroker(TopicName topic) { if (log.isDebugEnabled()) { log.debug("[{}] Handle Lookup for {}", ctx.channel(), topic); } - return kopBrokerLookupManager.findBroker(topic, advertisedEndPoint) + return kopBrokerLookupManager.findBroker(topic.toString(), advertisedEndPoint) .thenApply(listenerInetSocketAddressOpt -> listenerInetSocketAddressOpt .map(inetSocketAddress -> newPartitionMetadata(topic, newNode(inetSocketAddress))) .orElse(null) 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 038d7779ef..d42307517d 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 @@ -23,7 +23,6 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicBoolean; import javax.annotation.Nullable; -import lombok.NonNull; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.apache.kafka.common.security.auth.SecurityProtocol; @@ -49,8 +48,6 @@ public class KopBrokerLookupManager { public static final ConcurrentHashMap> LOOKUP_CACHE = new ConcurrentHashMap<>(); - public static final ConcurrentHashMap KOP_ADDRESS_CACHE = new ConcurrentHashMap<>(); - public KopBrokerLookupManager(KafkaServiceConfiguration conf, PulsarService pulsarService) throws Exception { this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); // TODO: change it to getKafkaAdvertisedProtocolMap @@ -59,33 +56,35 @@ public KopBrokerLookupManager(KafkaServiceConfiguration conf, PulsarService puls conf.getBrokerLookupTimeoutMs()); } - public CompletableFuture> findBroker(@NonNull TopicName topic, + public CompletableFuture> findBroker(String topic, @Nullable EndPoint advertisedEndPoint) { - if (log.isDebugEnabled()) { - log.debug("Handle Lookup for topic {}", topic); - } - - return getTopicBroker(topic.toString()) + return getTopicBroker(topic) .thenApply(internalListenerAddress -> { - final String listener = getAdvertisedListener( - internalListenerAddress, topic, advertisedEndPoint); - if (listener == null) { - removeTopicManagerCache(topic.toString()); + if (internalListenerAddress == null) { + log.error("[{}] failed get pulsar address, returned null.", topic); + removeTopicManagerCache(topic); return Optional.empty(); + } else if (log.isDebugEnabled()) { + log.debug("[{}] Found broker's internal listener address: {}", + topic, internalListenerAddress); } try { + final String listener = getAdvertisedListener( + internalListenerAddress, topic, advertisedEndPoint); + if (log.isDebugEnabled()) { + log.debug("Found listener {} for topic {}", listener, topic); + } final EndPoint endPoint = new EndPoint(listener, protocolMap); return Optional.of(new InetSocketAddress(endPoint.getHostname(), endPoint.getPort())); } catch (IllegalStateException e) { - log.error("Failed to create EndPoint from '{}': {}", listener, e.getMessage()); + log.error("Failed to find the advertised listener: {}", e.getMessage()); + removeTopicManagerCache(topic); return Optional.empty(); } }); } - // call pulsarclient.lookup.getbroker to get and own a topic. - // when error happens, the returned future will complete with null. public CompletableFuture getTopicBroker(String topicName) { if (closed.get()) { if (log.isDebugEnabled()) { @@ -94,6 +93,9 @@ public CompletableFuture getTopicBroker(String topicName) { return CompletableFuture.completedFuture(null); } + if (log.isDebugEnabled()) { + log.debug("Handle Lookup for topic {}", topicName); + } return LOOKUP_CACHE.computeIfAbsent(topicName, this::lookupBroker); } @@ -108,24 +110,10 @@ private CompletableFuture lookupBroker(final String topic) { } private String getAdvertisedListener(InetSocketAddress internalListenerAddress, - TopicName topic, + String topic, @Nullable EndPoint advertisedEndPoint) { - if (internalListenerAddress == null) { - log.error("[{}] failed get pulsar address, returned null.", topic); - return null; - } - if (log.isDebugEnabled()) { - log.debug("[{}] Found broker's internal listener address: {}", - topic, internalListenerAddress); - } - - final String protocolDataToAdvertise = KOP_ADDRESS_CACHE.get(topic.toString()); - if (protocolDataToAdvertise != null) { - return protocolDataToAdvertise; - } // else: either the key doesn't exist or a null value was put - - List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); + final List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); if (log.isDebugEnabled()) { availableBrokers.forEach(loadManagerReport -> log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, " @@ -146,24 +134,13 @@ private String getAdvertisedListener(InetSocketAddress internalListenerAddress, return null; } - final String kafkaAdvertisedListeners = - serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME).orElse(null); - if (kafkaAdvertisedListeners == null) { - log.error("No kafkaAdvertisedListeners found in broker {}", internalListenerAddress); - return null; - } - - if (log.isDebugEnabled()) { - log.debug("Found kafkaAdvertisedListeners: {}", kafkaAdvertisedListeners); - } - - final String listenerName = (advertisedEndPoint != null) ? advertisedEndPoint.getListenerName() : null; - if (listenerName == null) { - // Treat the whole kafkaAdvertisedListeners as the listener - return kafkaAdvertisedListeners; - } else { - return EndPoint.findListener(kafkaAdvertisedListeners, listenerName); - } + return serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME).map(kafkaAdvertisedListeners -> + Optional.ofNullable(advertisedEndPoint) + .map(endPoint -> EndPoint.findListener(kafkaAdvertisedListeners, endPoint.getListenerName())) + .orElse(EndPoint.findFirstListener(kafkaAdvertisedListeners)) + ).orElseGet(() -> { + throw new IllegalStateException("No kafkaAdvertisedListeners found in broker " + internalListenerAddress); + }); } // whether a ServiceLookupData contains wanted address. @@ -174,12 +151,10 @@ private static boolean lookupDataContainsAddress(ServiceLookupData data, String public static void removeTopicManagerCache(String topicName) { LOOKUP_CACHE.remove(topicName); - KOP_ADDRESS_CACHE.remove(topicName); } public static void clear() { LOOKUP_CACHE.clear(); - KOP_ADDRESS_CACHE.clear(); } public void close() { diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java index ab0079d18c..9c11e48f62 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java @@ -44,7 +44,6 @@ import org.apache.kafka.common.requests.TransactionResult; import org.apache.kafka.common.requests.WriteTxnMarkersRequest; import org.apache.kafka.common.requests.WriteTxnMarkersRequest.TxnMarkerEntry; -import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.netty.ChannelFutures; import org.eclipse.jetty.util.BlockingArrayQueue; @@ -225,7 +224,7 @@ public void addTxnMarkersToBrokerQueue(String transactionalId, for (TopicPartition topicPartition : topicPartitions) { String pulsarTopic = new KopTopic(topicPartition.topic()).getPartitionName(topicPartition.partition()); CompletableFuture> addressFuture = - kopBrokerLookupManager.findBroker(TopicName.get(pulsarTopic), sslEndPoint); + kopBrokerLookupManager.findBroker(pulsarTopic, sslEndPoint); CompletableFuture addFuture = new CompletableFuture<>(); addressFutureList.add(addFuture); addressFuture.whenComplete((address, throwable) -> { diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/CacheInvalidatorTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/CacheInvalidatorTest.java index 245663c4e2..810339b340 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/CacheInvalidatorTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/CacheInvalidatorTest.java @@ -74,7 +74,6 @@ public void testCacheInvalidatorIsTriggered() throws Exception { assertEquals("value", record.value()); } - assertFalse(KopBrokerLookupManager.KOP_ADDRESS_CACHE.isEmpty()); assertFalse(KopBrokerLookupManager.LOOKUP_CACHE.isEmpty()); BundlesData bundles = pulsar.getAdminClient().namespaces().getBundles( @@ -88,8 +87,6 @@ public void testCacheInvalidatorIsTriggered() throws Exception { Awaitility.await().untilAsserted(() -> { log.info("LOOKUP_CACHE {}", KopBrokerLookupManager.LOOKUP_CACHE); - log.info("KOP_ADDRESS_CACHE {}", KopBrokerLookupManager.KOP_ADDRESS_CACHE); - assertTrue(KopBrokerLookupManager.KOP_ADDRESS_CACHE.isEmpty()); assertTrue(KopBrokerLookupManager.LOOKUP_CACHE.isEmpty()); }); From a0b21c4b1205f1ab166a3b5fb26993e0bf21b8df Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 19:45:47 +0800 Subject: [PATCH 3/8] Store channel initializers for tests --- .../pulsar/handlers/kop/KafkaChannelInitializer.java | 12 ++++++++---- .../pulsar/handlers/kop/KafkaProtocolHandler.java | 6 +++++- .../pulsar/handlers/kop/KafkaRequestHandler.java | 2 -- .../handlers/kop/KafkaServiceConfiguration.java | 5 +++-- 4 files changed, 16 insertions(+), 9 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 3f006cdce2..e3fa64c890 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 @@ -15,6 +15,7 @@ import static io.streamnative.pulsar.handlers.kop.KafkaProtocolHandler.TLS_HANDLER; +import com.google.common.annotations.VisibleForTesting; import io.netty.channel.ChannelInitializer; import io.netty.channel.socket.SocketChannel; import io.netty.handler.codec.LengthFieldBasedFrameDecoder; @@ -92,10 +93,13 @@ protected void initChannel(SocketChannel ch) throws Exception { ch.pipeline().addLast(new LengthFieldPrepender(4)); ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder(MAX_FRAME_LENGTH, 0, 4, 0, 4)); - ch.pipeline().addLast("handler", - new KafkaRequestHandler(pulsarService, kafkaConfig, - tenantContextManager, kopBrokerLookupManager, adminManager, - enableTls, advertisedEndPoint, statsLogger)); + ch.pipeline().addLast("handler", newCnx()); } + @VisibleForTesting + public KafkaRequestHandler newCnx() throws Exception { + return new KafkaRequestHandler(pulsarService, kafkaConfig, + tenantContextManager, kopBrokerLookupManager, adminManager, + 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 3f639778a1..d1c4dd0107 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 @@ -85,6 +85,9 @@ public class KafkaProtocolHandler implements ProtocolHandler, TenantContextManag private KopBrokerLookupManager kopBrokerLookupManager; private AdminManager adminManager = null; private SystemTopicClient txnTopicClient; + @VisibleForTesting + @Getter + private Map> channelInitializerMap; @Getter @VisibleForTesting @@ -577,7 +580,8 @@ public Map> newChannelIniti forEach((listener, endPoint) -> builder.put(endPoint.getInetAddress(), newKafkaChannelInitializer(endPoint)) ); - return builder.build(); + channelInitializerMap = builder.build(); + return channelInitializerMap; } catch (Exception e){ log.error("KafkaProtocolHandler newChannelInitializers failed with ", e); return null; 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 3bc857b73b..c377032f54 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 @@ -203,7 +203,6 @@ public class KafkaRequestHandler extends KafkaCommandDecoder { private final Boolean tlsEnabled; private final EndPoint advertisedEndPoint; - private final String advertisedListeners; private final int defaultNumPartitions; public final int maxReadEntriesNum; private final int failedAuthenticationDelayMs; @@ -308,7 +307,6 @@ public KafkaRequestHandler(PulsarService pulsarService, this.adminManager = adminManager; this.tlsEnabled = tlsEnabled; this.advertisedEndPoint = advertisedEndPoint; - this.advertisedListeners = kafkaConfig.getKafkaAdvertisedListeners(); this.topicManager = new KafkaTopicManager(this); this.defaultNumPartitions = kafkaConfig.getDefaultNumPartitions(); this.maxReadEntriesNum = kafkaConfig.getMaxReadEntriesNum(); 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 7726b241d5..108328c56a 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 @@ -195,10 +195,11 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { ) private String kafkaProtocolMap; - @Deprecated @FieldContext( category = CATEGORY_KOP, - doc = "Use kafkaProtocolMap, kafkaListeners and advertisedAddress instead." + doc = "Specify the internal listener name for the broker.\n" + + "The listener name must be contained in the advertisedListeners.\n" + + "This config is used as the listener name in topic lookup." ) private String kafkaAdvertisedListeners; From 2e4828f54fc43c5185485a06822d0a684a14a17c Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 20:28:48 +0800 Subject: [PATCH 4/8] Add tests for findBroker --- .../handlers/kop/KafkaListenerNameTest.java | 60 +++++++++++++++++++ 1 file changed, 60 insertions(+) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java index 5d3b75fa66..510d045493 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java @@ -13,6 +13,13 @@ */ package io.streamnative.pulsar.handlers.kop; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; + +import io.netty.channel.ChannelHandlerContext; +import java.net.InetSocketAddress; +import java.util.HashMap; +import java.util.Map; import java.util.Properties; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; @@ -22,8 +29,11 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.requests.MetadataResponse; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.broker.ServiceConfigurationUtils; +import org.apache.pulsar.common.naming.TopicName; import org.testng.Assert; import org.testng.annotations.Test; @@ -44,6 +54,56 @@ protected void cleanup() throws Exception { // Clean up in the test method } + @Test(timeOut = 30000) + public void testFindBrokerForMultipleListeners() throws Exception { + final Map bindPortToAdvertisedAddress = new HashMap<>(); + final int anotherKafkaPort = PortManager.nextFreePort(); + bindPortToAdvertisedAddress.put(kafkaBrokerPort, + InetSocketAddress.createUnresolved("192.168.0.1", PortManager.nextFreePort())); + bindPortToAdvertisedAddress.put(anotherKafkaPort, + InetSocketAddress.createUnresolved("192.168.0.2", PortManager.nextFreePort())); + + super.resetConfig(); + conf.setKafkaListeners("PLAINTEXT://0.0.0.0:" + kafkaBrokerPort + ",GW://0.0.0.0:" + anotherKafkaPort); + conf.setKafkaProtocolMap("PLAINTEXT:PLAINTEXT,GW:PLAINTEXT"); + + conf.setKafkaAdvertisedListeners(String.format("PLAINTEXT://%s,GW://%s", + bindPortToAdvertisedAddress.get(kafkaBrokerPort), + bindPortToAdvertisedAddress.get(anotherKafkaPort))); + super.internalSetup(); + + final String topic = "persistent://public/default/test-find-broker-for-multiple-listeners"; + admin.topics().createPartitionedTopic(topic, 1); + final String partitionName = topic + "-partition-0"; + + final KafkaProtocolHandler protocolHandler = (KafkaProtocolHandler) + pulsar.getProtocolHandlers().protocol(KafkaProtocolHandler.PROTOCOL_NAME); + protocolHandler.getChannelInitializerMap().forEach((inetSocketAddress, channelInitializer) -> { + try { + final KafkaRequestHandler requestHandler = ((KafkaChannelInitializer) channelInitializer).newCnx(); + ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); + io.netty.channel.Channel mockChannel = mock(io.netty.channel.Channel.class); + doReturn(mockChannel).when(mockCtx).channel(); + requestHandler.ctx = mockCtx; + + final MetadataResponse.PartitionMetadata partitionMetadata = + requestHandler.findBroker(TopicName.get(partitionName)).get(); + Assert.assertEquals(partitionMetadata.error(), Errors.NONE); + + final InetSocketAddress expectedAddress = bindPortToAdvertisedAddress.get(inetSocketAddress.getPort()); + log.info("[{}] Expected advertised listener: {}, partition metadata: {}", + inetSocketAddress, expectedAddress, partitionMetadata.leader()); + + Assert.assertEquals(partitionMetadata.leader().host(), expectedAddress.getHostName()); + Assert.assertEquals(partitionMetadata.leader().port(), expectedAddress.getPort()); + } catch (Exception e) { + Assert.fail(e.getMessage()); + } + }); + + super.internalCleanup(); + } + @Test(timeOut = 30000) public void testListenerName() throws Exception { super.resetConfig(); From 0c8764f5f8f864453eb3b1bfb7d21c394b771886 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 20:53:21 +0800 Subject: [PATCH 5/8] Test metadata request --- .../pulsar/handlers/kop/KafkaApisTest.java | 7 ++- .../handlers/kop/KafkaListenerNameTest.java | 55 ++++++++++++++----- 2 files changed, 46 insertions(+), 16 deletions(-) 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 1b53fdc5f3..e770ae3bd9 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 @@ -164,7 +164,12 @@ protected void cleanup() throws Exception { super.internalCleanup(); } - KafkaHeaderAndRequest buildRequest(AbstractRequest.Builder builder) { + private KafkaHeaderAndRequest buildRequest(AbstractRequest.Builder builder) { + return buildRequest(builder, serviceAddress); + } + + static KafkaHeaderAndRequest buildRequest(AbstractRequest.Builder builder, + SocketAddress serviceAddress) { AbstractRequest request = builder.build(); builder.apiKey(); diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java index 510d045493..c1ff9631a1 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/KafkaListenerNameTest.java @@ -13,14 +13,21 @@ */ package io.streamnative.pulsar.handlers.kop; +import static io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder.KafkaHeaderAndRequest; +import static org.apache.kafka.common.requests.MetadataResponse.PartitionMetadata; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; +import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; import java.net.InetSocketAddress; +import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Properties; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; @@ -29,11 +36,13 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.Node; import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.requests.AbstractResponse; +import org.apache.kafka.common.requests.MetadataRequest; import org.apache.kafka.common.requests.MetadataResponse; import org.apache.kafka.common.serialization.StringSerializer; import org.apache.pulsar.broker.ServiceConfigurationUtils; -import org.apache.pulsar.common.naming.TopicName; import org.testng.Assert; import org.testng.annotations.Test; @@ -55,7 +64,7 @@ protected void cleanup() throws Exception { } @Test(timeOut = 30000) - public void testFindBrokerForMultipleListeners() throws Exception { + public void testMetadataRequestForMultiListeners() throws Exception { final Map bindPortToAdvertisedAddress = new HashMap<>(); final int anotherKafkaPort = PortManager.nextFreePort(); bindPortToAdvertisedAddress.put(kafkaBrokerPort, @@ -72,9 +81,9 @@ public void testFindBrokerForMultipleListeners() throws Exception { bindPortToAdvertisedAddress.get(anotherKafkaPort))); super.internalSetup(); - final String topic = "persistent://public/default/test-find-broker-for-multiple-listeners"; - admin.topics().createPartitionedTopic(topic, 1); - final String partitionName = topic + "-partition-0"; + final String topic = "persistent://public/default/test-metadata-request-for-multi-listeners"; + final int numPartitions = 3; + admin.topics().createPartitionedTopic(topic, numPartitions); final KafkaProtocolHandler protocolHandler = (KafkaProtocolHandler) pulsar.getProtocolHandlers().protocol(KafkaProtocolHandler.PROTOCOL_NAME); @@ -82,20 +91,36 @@ public void testFindBrokerForMultipleListeners() throws Exception { try { final KafkaRequestHandler requestHandler = ((KafkaChannelInitializer) channelInitializer).newCnx(); ChannelHandlerContext mockCtx = mock(ChannelHandlerContext.class); - io.netty.channel.Channel mockChannel = mock(io.netty.channel.Channel.class); - doReturn(mockChannel).when(mockCtx).channel(); + doReturn(mock(Channel.class)).when(mockCtx).channel(); requestHandler.ctx = mockCtx; - final MetadataResponse.PartitionMetadata partitionMetadata = - requestHandler.findBroker(TopicName.get(partitionName)).get(); - Assert.assertEquals(partitionMetadata.error(), Errors.NONE); - final InetSocketAddress expectedAddress = bindPortToAdvertisedAddress.get(inetSocketAddress.getPort()); - log.info("[{}] Expected advertised listener: {}, partition metadata: {}", - inetSocketAddress, expectedAddress, partitionMetadata.leader()); - Assert.assertEquals(partitionMetadata.leader().host(), expectedAddress.getHostName()); - Assert.assertEquals(partitionMetadata.leader().port(), expectedAddress.getPort()); + final KafkaHeaderAndRequest metadataRequest = KafkaApisTest.buildRequest( + new MetadataRequest.Builder(Collections.singletonList(topic), true), + inetSocketAddress); + final CompletableFuture future = new CompletableFuture<>(); + requestHandler.handleTopicMetadataRequest(metadataRequest, future); + final MetadataResponse metadataResponse = (MetadataResponse) future.get(); + Assert.assertEquals(metadataResponse.brokers().size(), 1); + final List brokers = new ArrayList<>(metadataResponse.brokers()); + Assert.assertEquals(brokers.size(), 1); + Assert.assertEquals(brokers.get(0).host(), expectedAddress.getHostName()); + Assert.assertEquals(brokers.get(0).port(), expectedAddress.getPort()); + + final List topicMetadataList = + new ArrayList<>(metadataResponse.topicMetadata()); + Assert.assertEquals(topicMetadataList.size(), 1); + Assert.assertEquals(topicMetadataList.get(0).topic(), topic); + + final List partitionMetadataList = topicMetadataList.get(0).partitionMetadata(); + Assert.assertEquals(partitionMetadataList.size(), numPartitions); + for (int i = 0; i < numPartitions; i++) { + final PartitionMetadata partitionMetadata = partitionMetadataList.get(i); + Assert.assertEquals(partitionMetadata.error(), Errors.NONE); + Assert.assertEquals(partitionMetadata.leader().host(), expectedAddress.getHostName()); + Assert.assertEquals(partitionMetadata.leader().port(), expectedAddress.getPort()); + } } catch (Exception e) { Assert.fail(e.getMessage()); } From d0b09fefd1a6057cf6f0535209fe662a7db9e589 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 28 Oct 2021 21:33:21 +0800 Subject: [PATCH 6/8] Remove unused protocol map --- .../pulsar/handlers/kop/KopBrokerLookupManager.java | 13 +++++-------- 1 file changed, 5 insertions(+), 8 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 d42307517d..d113e8d13e 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 @@ -17,15 +17,14 @@ import java.io.IOException; import java.net.InetSocketAddress; 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.atomic.AtomicBoolean; +import java.util.regex.Matcher; import javax.annotation.Nullable; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.apache.kafka.common.security.auth.SecurityProtocol; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.resources.MetadataStoreCacheLoader; import org.apache.pulsar.common.naming.TopicName; @@ -40,7 +39,6 @@ public class KopBrokerLookupManager { private final LookupClient lookupClient; - private final Map protocolMap; private final MetadataStoreCacheLoader metadataStoreCacheLoader; private final AtomicBoolean closed = new AtomicBoolean(false); @@ -50,8 +48,6 @@ public class KopBrokerLookupManager { public KopBrokerLookupManager(KafkaServiceConfiguration conf, PulsarService pulsarService) throws Exception { this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); - // TODO: change it to getKafkaAdvertisedProtocolMap - this.protocolMap = EndPoint.parseProtocolMap(conf.getKafkaProtocolMap()); this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), conf.getBrokerLookupTimeoutMs()); } @@ -75,9 +71,10 @@ public CompletableFuture> findBroker(String topic, if (log.isDebugEnabled()) { log.debug("Found listener {} for topic {}", listener, topic); } - final EndPoint endPoint = new EndPoint(listener, protocolMap); - return Optional.of(new InetSocketAddress(endPoint.getHostname(), endPoint.getPort())); - } catch (IllegalStateException e) { + final Matcher matcher = EndPoint.matcherListener(listener, + listener + " cannot be split into 3 parts"); + return Optional.of(new InetSocketAddress(matcher.group(2), Integer.parseInt(matcher.group(3)))); + } catch (IllegalStateException | NumberFormatException e) { log.error("Failed to find the advertised listener: {}", e.getMessage()); removeTopicManagerCache(topic); return Optional.empty(); From 11950e11808b956ebd18a1b9d119042394c822c6 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 29 Oct 2021 10:46:07 +0800 Subject: [PATCH 7/8] Use orElseThrow instead of orElseGet --- .../pulsar/handlers/kop/KopBrokerLookupManager.java | 5 ++--- 1 file changed, 2 insertions(+), 3 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 d113e8d13e..6c70206009 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 @@ -135,9 +135,8 @@ private String getAdvertisedListener(InetSocketAddress internalListenerAddress, Optional.ofNullable(advertisedEndPoint) .map(endPoint -> EndPoint.findListener(kafkaAdvertisedListeners, endPoint.getListenerName())) .orElse(EndPoint.findFirstListener(kafkaAdvertisedListeners)) - ).orElseGet(() -> { - throw new IllegalStateException("No kafkaAdvertisedListeners found in broker " + internalListenerAddress); - }); + ).orElseThrow(() -> new IllegalStateException( + "No kafkaAdvertisedListeners found in broker " + internalListenerAddress)); } // whether a ServiceLookupData contains wanted address. From 75c05bd2cc3d2182d7d0ed42dee41afe8bdca821 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Fri, 29 Oct 2021 11:29:00 +0800 Subject: [PATCH 8/8] Fix config description --- .../pulsar/handlers/kop/KafkaServiceConfiguration.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) 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 108328c56a..cec0d353f5 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 @@ -184,6 +184,8 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { category = CATEGORY_KOP, doc = "Comma-separated list of URIs we will listen on and the listener names.\n" + "e.g. PLAINTEXT://localhost:9092,SSL://localhost:9093.\n" + + "Each URI's scheme represents a listener name if `kafkaProtocolMap` is configured.\n" + + "Otherwise, the scheme must be a valid protocol in [PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL].\n" + "If hostname is not set, bind to the default interface." ) private String kafkaListeners; @@ -197,9 +199,8 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { @FieldContext( category = CATEGORY_KOP, - doc = "Specify the internal listener name for the broker.\n" - + "The listener name must be contained in the advertisedListeners.\n" - + "This config is used as the listener name in topic lookup." + doc = "Listeners to publish to ZooKeeper for clients to use.\n" + + "The format is the same as `kafkaListeners`.\n" ) private String kafkaAdvertisedListeners;