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..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 @@ -90,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, ""); @@ -98,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( @@ -115,6 +127,22 @@ public static Map parseListeners(final String listeners, final return parseListeners(listeners, parseProtocolMap(protocolMapString)); } + public static String findListener(final String listeners, final String name) { + if (name == null) { + return null; + } + for (String listener : getListenerArray(listeners)) { + if (listener.contains(":") && listener.substring(0, listener.indexOf(":")).equals(name)) { + return listener; + } + } + 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) { 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/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 2b5a7e7f2a..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(); @@ -2481,7 +2479,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/KafkaServiceConfiguration.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaServiceConfiguration.java index 7726b241d5..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; @@ -195,10 +197,10 @@ public class KafkaServiceConfiguration extends ServiceConfiguration { ) private String kafkaProtocolMap; - @Deprecated @FieldContext( category = CATEGORY_KOP, - doc = "Use kafkaProtocolMap, kafkaListeners and advertisedAddress instead." + doc = "Listeners to publish to ZooKeeper for clients to use.\n" + + "The format is the same as `kafkaListeners`.\n" ) private String kafkaAdvertisedListeners; 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..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 @@ -21,8 +21,8 @@ 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.NonNull; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarService; @@ -38,75 +38,51 @@ @Slf4j public class KopBrokerLookupManager { - private final String advertisedListeners; private final LookupClient lookupClient; 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 KopBrokerLookupManager(KafkaServiceConfiguration conf, PulsarService pulsarService) throws Exception { - this.advertisedListeners = conf.getKafkaAdvertisedListeners(); this.lookupClient = KafkaProtocolHandler.getLookupClient(pulsarService); this.metadataStoreCacheLoader = new MetadataStoreCacheLoader(pulsarService.getPulsarResources(), 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); - } - 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); - removeTopicManagerCache(topic.toString()); - returnFuture.complete(Optional.empty()); - return; + return getTopicBroker(topic) + .thenApply(internalListenerAddress -> { + 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); } - // 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 String listener = getAdvertisedListener( + internalListenerAddress, topic, advertisedEndPoint); + if (log.isDebugEnabled()) { + log.debug("Found listener {} for topic {}", listener, topic); + } + 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(); } - - // 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,73 +90,27 @@ 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); - }); + if (log.isDebugEnabled()) { + log.debug("Handle Lookup for topic {}", topicName); + } + 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, + String topic, + @Nullable EndPoint advertisedEndPoint) { - if (log.isDebugEnabled()) { - log.debug("Found broker for topic {} puslarAddress: {}", - topic, pulsarAddress); - } - - // 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; - } - - List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); + final List availableBrokers = metadataStoreCacheLoader.getAvailableBrokers(); if (log.isDebugEnabled()) { availableBrokers.forEach(loadManagerReport -> log.debug("Handle getProtocolDataToAdvertise for {}, pulsarUrl: {}, " @@ -193,18 +123,20 @@ 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)); - } else { - log.error("No node for broker {} under loadBalance", pulsarAddress); - removeTopicManagerCache(topic.toString()); - returnFuture.complete(Optional.empty()); + if (!serviceLookupData.isPresent()) { + log.error("No node for broker {} under loadBalance", internalListenerAddress); + return null; } - return returnFuture; + + return serviceLookupData.get().getProtocol(KafkaProtocolHandler.PROTOCOL_NAME).map(kafkaAdvertisedListeners -> + Optional.ofNullable(advertisedEndPoint) + .map(endPoint -> EndPoint.findListener(kafkaAdvertisedListeners, endPoint.getListenerName())) + .orElse(EndPoint.findFirstListener(kafkaAdvertisedListeners)) + ).orElseThrow(() -> new IllegalStateException( + "No kafkaAdvertisedListeners found in broker " + internalListenerAddress)); } // whether a ServiceLookupData contains wanted address. @@ -215,12 +147,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/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); } } 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()); }); 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 5d3b75fa66..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,7 +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; @@ -22,6 +36,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.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.testng.Assert; @@ -44,6 +63,72 @@ protected void cleanup() throws Exception { // Clean up in the test method } + @Test(timeOut = 30000) + public void testMetadataRequestForMultiListeners() 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-metadata-request-for-multi-listeners"; + final int numPartitions = 3; + admin.topics().createPartitionedTopic(topic, numPartitions); + + 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); + doReturn(mock(Channel.class)).when(mockCtx).channel(); + requestHandler.ctx = mockCtx; + + final InetSocketAddress expectedAddress = bindPortToAdvertisedAddress.get(inetSocketAddress.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()); + } + }); + + super.internalCleanup(); + } + @Test(timeOut = 30000) public void testListenerName() throws Exception { super.resetConfig();