From 902a41076a580a9583a9aa85361ba92bec9702dd Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Thu, 19 Oct 2023 14:45:16 -0700 Subject: [PATCH 1/4] [improve][broker] PIP-307 Added assignedBrokerUrl to CloseProducerCmd to skip lookups upon producer reconnections during unloading --- .../extensions/ExtensibleLoadManagerImpl.java | 45 +++++++ .../channel/ServiceUnitStateChannel.java | 10 ++ .../channel/ServiceUnitStateChannelImpl.java | 66 +++++++++-- .../pulsar/broker/namespace/OwnedBundle.java | 2 +- .../pulsar/broker/service/AbstractTopic.java | 5 + .../pulsar/broker/service/BrokerService.java | 9 +- .../pulsar/broker/service/Producer.java | 9 +- .../pulsar/broker/service/ServerCnx.java | 44 ++++++- .../apache/pulsar/broker/service/Topic.java | 5 + .../pulsar/broker/service/TransportCnx.java | 3 + .../nonpersistent/NonPersistentTopic.java | 26 +++- .../service/persistent/PersistentTopic.java | 33 +++++- .../ExtensibleLoadManagerImplTest.java | 111 +++++++++++++++++- .../channel/ServiceUnitStateChannelTest.java | 1 + .../broker/namespace/OwnershipCacheTest.java | 2 +- .../apache/pulsar/client/impl/ClientCnx.java | 23 +++- .../pulsar/client/impl/ConnectionHandler.java | 21 +++- .../pulsar/common/protocol/Commands.java | 26 +++- pulsar-common/src/main/proto/PulsarApi.proto | 2 + 19 files changed, 404 insertions(+), 39 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java index d3119365ddfea..8e166f00bc95f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java @@ -38,6 +38,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; import java.util.stream.Collectors; @@ -178,6 +179,10 @@ public class ExtensibleLoadManagerImpl implements ExtensibleLoadManager { private final UnloadCounter unloadCounter = new UnloadCounter(); private final SplitCounter splitCounter = new SplitCounter(); + // Record the ignored send msg count during unloading + @Getter + private final AtomicLong ignoredSendMsgCounter = new AtomicLong(); + // record unload metrics private final AtomicReference> unloadMetrics = new AtomicReference<>(); // record split metrics @@ -269,6 +274,15 @@ public static ExtensibleLoadManagerImpl get(LoadManager loadManager) { return loadManagerWrapper.get(); } + /** + * A static util func to get the ExtensibleLoadManagerImpl instance. + * @param pulsar PulsarService + * @return the ExtensibleLoadManagerImpl instance + */ + public static ExtensibleLoadManagerImpl get(PulsarService pulsar) { + return get(pulsar.getLoadManager().get()); + } + public static boolean debug(ServiceConfiguration config, Logger log) { return config.isLoadBalancerDebugModeEnabled() || log.isDebugEnabled(); } @@ -286,6 +300,37 @@ public static void createSystemTopic(PulsarService pulsar, String topic) throws } } + /** + * Gets the assigned broker for the given topic. + * @param pulsar PulsarService instance + * @param topic Topic Name + * @return the assigned broker's BrokerLookupData instance. Empty, if not assigned by Extensible LoadManager. + */ + public static CompletableFuture> getAssignedBrokerLookupData(PulsarService pulsar, + String topic) { + if (ExtensibleLoadManagerImpl.isLoadManagerExtensionEnabled(pulsar.getConfig())) { + var topicName = TopicName.get(topic); + try { + return pulsar.getNamespaceService().getBundleAsync(topicName) + .thenCompose(bundle -> { + var loadManager = ExtensibleLoadManagerImpl.get(pulsar); + var assigned = loadManager.getServiceUnitStateChannel() + .getAssigned(bundle.toString()); + if (assigned.isPresent()) { + return loadManager.getBrokerRegistry().lookupAsync(assigned.get()); + } else { + return CompletableFuture.completedFuture(Optional.empty()); + } + } + ); + } catch (Throwable e) { + log.error("Failed to DestinationBrokerLookupData for topic:{}", topic, e); + return CompletableFuture.completedFuture(Optional.empty()); + } + } + return CompletableFuture.completedFuture(Optional.empty()); + } + @Override public void start() throws PulsarServerException { if (this.started) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannel.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannel.java index 6e75fe91a914f..9be76e1b0f44d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannel.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannel.java @@ -131,6 +131,16 @@ public interface ServiceUnitStateChannel extends Closeable { */ CompletableFuture> getOwnerAsync(String serviceUnit); + /** + * Gets the assigned broker of the service unit. + * + * + * @param serviceUnit (e.g. bundle)) + * @return the future object of the assigned broker + */ + Optional getAssigned(String serviceUnit); + + /** * Checks if the target broker is the owner of the service unit. * diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java index 68501d201f0d4..d0a74c0eaef8a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java @@ -538,6 +538,44 @@ public CompletableFuture> getOwnerAsync(String serviceUnit) { } } + @Override + public Optional getAssigned(String serviceUnit) { + if (!validateChannelState(Started, true)) { + return Optional.empty(); + } + + ServiceUnitStateData data = tableview.get(serviceUnit); + if (data == null) { + return Optional.empty(); + } + ServiceUnitState state = state(data); + switch (state) { + case Owned, Assigning -> { + return Optional.of(data.dstBroker()); + } + case Releasing -> { + if (data.dstBroker() != null) { + return Optional.of(data.dstBroker()); + } + return Optional.empty(); + } + case Splitting -> { + return Optional.of(data.sourceBroker()); + } + case Init, Free -> { + return Optional.empty(); + } + case Deleted -> { + log.warn("Trying to get the assigned broker from the deleted serviceUnit:{}", serviceUnit); + return Optional.empty(); + } + default -> { + log.warn("Trying to get the assigned broker from unknown state:{} serviceUnit:{}", state, serviceUnit); + return Optional.empty(); + } + } + } + private long getNextVersionId(String serviceUnit) { var data = tableview.get(serviceUnit); return getNextVersionId(data); @@ -732,15 +770,20 @@ private void handleOwnEvent(String serviceUnit, ServiceUnitStateData data) { if (getOwnerRequest != null) { getOwnerRequest.complete(data.dstBroker()); } - stateChangeListeners.notify(serviceUnit, data, null); + CompletableFuture ownFuture = null; if (isTargetBroker(data.dstBroker())) { - log(null, serviceUnit, data, null); pulsar.getNamespaceService() .onNamespaceBundleOwned(LoadManagerShared.getNamespaceBundle(pulsar, serviceUnit)); lastOwnEventHandledAt = System.currentTimeMillis(); - } else if (data.force() && isTargetBroker(data.sourceBroker())) { - closeServiceUnit(serviceUnit); + ownFuture = CompletableFuture.completedFuture(null); + } else if ((data.force() || isTransferCommand(data)) && isTargetBroker(data.sourceBroker())) { + ownFuture = closeServiceUnit(serviceUnit, false); + } else { + ownFuture = CompletableFuture.completedFuture(null); } + + stateChangeListeners.notifyOnCompletion(ownFuture, serviceUnit, data) + .whenComplete((__, e) -> log(e, serviceUnit, data, null)); } private void handleAssignEvent(String serviceUnit, ServiceUnitStateData data) { @@ -755,15 +798,17 @@ private void handleAssignEvent(String serviceUnit, ServiceUnitStateData data) { private void handleReleaseEvent(String serviceUnit, ServiceUnitStateData data) { if (isTargetBroker(data.sourceBroker())) { ServiceUnitStateData next; + CompletableFuture unloadFuture; if (isTransferCommand(data)) { next = new ServiceUnitStateData( Assigning, data.dstBroker(), data.sourceBroker(), getNextVersionId(data)); - // TODO: when close, pass message to clients to connect to the new broker + unloadFuture = closeServiceUnit(serviceUnit, true); } else { next = new ServiceUnitStateData( Free, null, data.sourceBroker(), getNextVersionId(data)); + unloadFuture = closeServiceUnit(serviceUnit, false); } - stateChangeListeners.notifyOnCompletion(closeServiceUnit(serviceUnit) + stateChangeListeners.notifyOnCompletion(unloadFuture .thenCompose(__ -> pubAsync(serviceUnit, next)), serviceUnit, data) .whenComplete((__, e) -> log(e, serviceUnit, data, next)); } @@ -861,12 +906,13 @@ private CompletableFuture deferGetOwnerRequest(String serviceUnit) { } } - private CompletableFuture closeServiceUnit(String serviceUnit) { + private CompletableFuture closeServiceUnit(String serviceUnit, boolean closeWithoutDisconnectingClients) { long startTime = System.nanoTime(); MutableInt unloadedTopics = new MutableInt(); NamespaceBundle bundle = LoadManagerShared.getNamespaceBundle(pulsar, serviceUnit); return pulsar.getBrokerService().unloadServiceUnit( bundle, + closeWithoutDisconnectingClients, true, pulsar.getConfig().getNamespaceBundleUnloadingTimeoutMs(), TimeUnit.MILLISECONDS) @@ -875,8 +921,10 @@ private CompletableFuture closeServiceUnit(String serviceUnit) { return numUnloadedTopics; }) .whenComplete((__, ex) -> { - // clean up topics that failed to unload from the broker ownership cache - pulsar.getBrokerService().cleanUnloadedTopicFromCache(bundle); + if (!closeWithoutDisconnectingClients) { + // clean up topics that failed to unload from the broker ownership cache + pulsar.getBrokerService().cleanUnloadedTopicFromCache(bundle); + } pulsar.getNamespaceService().onNamespaceBundleUnload(bundle); double unloadBundleTime = TimeUnit.NANOSECONDS .toMillis((System.nanoTime() - startTime)); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java index e7cf23a042750..a87d45395db01 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java @@ -136,7 +136,7 @@ public CompletableFuture handleUnloadRequest(PulsarService pulsar, long ti return pulsar.getNamespaceService().getOwnershipCache() .updateBundleState(this.bundle, false) .thenCompose(v -> pulsar.getBrokerService().unloadServiceUnit( - bundle, closeWithoutWaitingClientDisconnect, timeout, timeoutUnit)) + bundle, false, closeWithoutWaitingClientDisconnect, timeout, timeoutUnit)) .handle((numUnloadedTopics, ex) -> { if (ex != null) { // ignore topic-close failure to unload bundle diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java index 9a5771e18cec4..741013a75dd50 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractTopic.java @@ -956,6 +956,11 @@ protected void checkTopicFenced() throws BrokerServiceException { } } + @Override + public boolean isFenced() { + return isFenced; + } + protected CompletableFuture internalAddProducer(Producer producer) { if (isProducersExceeded(producer)) { log.warn("[{}] Attempting to add producer to topic which reached max producers limit", topic); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 045f4d6467a64..ed9a4d20c87ef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2216,8 +2216,10 @@ public CompletableFuture checkTopicNsOwnership(final String topic) { } public CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, + boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect, long timeout, TimeUnit unit) { - CompletableFuture future = unloadServiceUnit(serviceUnit, closeWithoutWaitingClientDisconnect); + CompletableFuture future = unloadServiceUnit( + serviceUnit, closeWithoutDisconnectingClients, closeWithoutWaitingClientDisconnect); ScheduledFuture taskTimeout = executor().schedule(() -> { if (!future.isDone()) { log.warn("Unloading of {} has timed out", serviceUnit); @@ -2234,11 +2236,13 @@ public CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, * Unload all the topic served by the broker service under the given service unit. * * @param serviceUnit + * @param closeWithoutDisconnectingClients don't disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for clients to disconnect * and forcefully close managed-ledger * @return */ private CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, + boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect) { List> closeFutures = new ArrayList<>(); topics.forEach((name, topicFuture) -> { @@ -2262,7 +2266,8 @@ private CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit } } closeFutures.add(topicFuture - .thenCompose(t -> t.isPresent() ? t.get().close(closeWithoutWaitingClientDisconnect) + .thenCompose(t -> t.isPresent() ? t.get().close( + closeWithoutDisconnectingClients, closeWithoutWaitingClientDisconnect) : CompletableFuture.completedFuture(null))); } }); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index 53b79f06e8e24..16aa13d50860f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -39,6 +39,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; +import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData; import org.apache.pulsar.broker.service.BrokerServiceException.TopicClosedException; import org.apache.pulsar.broker.service.BrokerServiceException.TopicTerminatedException; import org.apache.pulsar.broker.service.Topic.PublishContext; @@ -703,17 +704,21 @@ public void closeNow(boolean removeFromTopic) { isDisconnecting.set(false); } + public CompletableFuture disconnect() { + return disconnect(Optional.empty()); + } + /** * It closes the producer from server-side and sends command to client to disconnect producer from existing * connection without closing that connection. * * @return Completable future indicating completion of producer close */ - public CompletableFuture disconnect() { + public CompletableFuture disconnect(Optional assignedBrokerLookupData) { if (!closeFuture.isDone() && isDisconnecting.compareAndSet(false, true)) { log.info("Disconnecting producer: {}", this); cnx.execute(() -> { - cnx.closeProducer(this); + cnx.closeProducer(this, assignedBrokerLookupData); closeNow(true); }); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index de96a317205bb..c48b31518d1b8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -81,6 +81,8 @@ import org.apache.pulsar.broker.authentication.AuthenticationState; import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.limiter.ConnectionController; +import org.apache.pulsar.broker.loadbalance.extensions.ExtensibleLoadManagerImpl; +import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; import org.apache.pulsar.broker.service.BrokerServiceException.ServerMetadataException; import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; @@ -1717,6 +1719,20 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { printSendCommandDebug(send, headersAndPayload); } + + ServiceConfiguration conf = getBrokerService().pulsar().getConfiguration(); + if (producer.getTopic().isFenced() + && ExtensibleLoadManagerImpl.isLoadManagerExtensionEnabled(conf)) { + long ignoredMsgCount = ExtensibleLoadManagerImpl.get(getBrokerService().pulsar()) + .getIgnoredSendMsgCounter().incrementAndGet(); + if (log.isDebugEnabled()) { + log.debug("Ignored send msg from:{}:{} to fenced topic:{} during unloading." + + " Ignored message count:{}.", + remoteAddress, send.getProducerId(), producer.getTopic().getName(), ignoredMsgCount); + } + return; + } + if (producer.isNonPersistentTopic()) { // avoid processing non-persist message if reached max concurrent-message limit if (nonPersistentPendingMessages > maxNonPersistentPendingMessages) { @@ -2963,15 +2979,31 @@ public void closeProducer(Producer producer) { assert ctx.executor().inEventLoop(); // removes producer-connection from map and send close command to producer safelyRemoveProducer(producer); + closeProducer(producer.getProducerId(), producer.getEpoch(), Optional.empty()); + } + + @Override + public void closeProducer(Producer producer, Optional assignedBrokerLookupData) { + // removes producer-connection from map and send close command to producer + safelyRemoveProducer(producer); + closeProducer(producer.getProducerId(), producer.getEpoch(), assignedBrokerLookupData); + } + + private void closeProducer(long producerId, long epoch, Optional assignedBrokerLookupData) { if (getRemoteEndpointProtocolVersion() >= v5.getValue()) { - writeAndFlush(Commands.newCloseProducer(producer.getProducerId(), -1L)); + if (assignedBrokerLookupData.isPresent()) { + writeAndFlush(Commands.newCloseProducer(producerId, -1L, + assignedBrokerLookupData.get().pulsarServiceUrl(), + assignedBrokerLookupData.get().pulsarServiceUrlTls())); + } else { + writeAndFlush(Commands.newCloseProducer(producerId, -1L)); + } + // The client does not necessarily know that the producer is closed, but the connection is still // active, and there could be messages in flight already. We want to ignore these messages for a time // because they are expected. Once the interval has passed, the client should have received the // CloseProducer command and should not send any additional messages until it sends a create Producer // command. - final long epoch = producer.getEpoch(); - final long producerId = producer.getProducerId(); recentlyClosedProducers.put(producerId, epoch); ctx.executor().schedule(() -> { recentlyClosedProducers.remove(producerId, epoch); @@ -2986,8 +3018,12 @@ public void closeProducer(Producer producer) { public void closeConsumer(Consumer consumer) { // removes consumer-connection from map and send close command to consumer safelyRemoveConsumer(consumer); + closeConsumer(consumer.consumerId()); + } + + private void closeConsumer(long consumerId) { if (getRemoteEndpointProtocolVersion() >= v5.getValue()) { - writeAndFlush(Commands.newCloseConsumer(consumer.consumerId(), -1L)); + writeAndFlush(Commands.newCloseConsumer(consumerId, -1L)); } else { close(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index 7657d77e1299f..ead6dd6982d1a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -196,6 +196,9 @@ CompletableFuture createSubscription(String subscriptionName, Init CompletableFuture close(boolean closeWithoutWaitingClientDisconnect); + CompletableFuture close( + boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect); + void checkGC(); CompletableFuture checkClusterMigration(); @@ -329,6 +332,8 @@ default boolean isSystemTopic() { boolean isPersistent(); + boolean isFenced(); + /* ------ Transaction related ------ */ /** diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TransportCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TransportCnx.java index d267160652ae4..d0b5c67767b1f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TransportCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/TransportCnx.java @@ -21,8 +21,10 @@ import io.netty.handler.codec.haproxy.HAProxyMessage; import io.netty.util.concurrent.Promise; import java.net.SocketAddress; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; +import org.apache.pulsar.broker.loadbalance.extensions.data.BrokerLookupData; public interface TransportCnx { @@ -55,6 +57,7 @@ public interface TransportCnx { void removedProducer(Producer producer); void closeProducer(Producer producer); + void closeProducer(Producer producer, Optional assignedBrokerLookupData); void cancelPublishRateLimiting(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index cd09f18736814..d232a1451787d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -40,6 +40,7 @@ import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.Position; import org.apache.pulsar.broker.PulsarServerException; +import org.apache.pulsar.broker.loadbalance.extensions.ExtensibleLoadManagerImpl; import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.resources.NamespaceResources; import org.apache.pulsar.broker.service.AbstractReplicator; @@ -485,14 +486,22 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, boolean c return deleteFuture; } + + @Override + public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { + return close(false, closeWithoutWaitingClientDisconnect); + } + /** * Close this topic - close all producers and subscriptions associated with this topic. * + * @param closeWithoutDisconnectingClients don't disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for client disconnect and forcefully close managed-ledger * @return Completable future indicating completion of close operation */ @Override - public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { + public CompletableFuture close( + boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect) { CompletableFuture closeFuture = new CompletableFuture<>(); lock.writeLock().lock(); @@ -511,7 +520,12 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect List> futures = new ArrayList<>(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + if (!closeWithoutDisconnectingClients) { + futures.add(ExtensibleLoadManagerImpl.getAssignedBrokerLookupData( + brokerService.getPulsar(), topic).thenAccept(lookupData -> + producers.values().forEach(producer -> futures.add(producer.disconnect(lookupData))) + )); + } if (topicPublishRateLimiter != null) { topicPublishRateLimiter.close(); } @@ -539,9 +553,13 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect // unload topic iterates over topics map and removing from the map with the same thread creates deadlock. // so, execute it in different thread brokerService.executor().execute(() -> { - brokerService.removeTopicFromCache(NonPersistentTopic.this); - unregisterTopicPolicyListener(); + + if (!closeWithoutDisconnectingClients) { + brokerService.removeTopicFromCache(NonPersistentTopic.this); + unregisterTopicPolicyListener(); + } closeFuture.complete(null); + }); }).exceptionally(exception -> { log.error("[{}] Error closing topic", topic, exception); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 7ca57ac8657e8..13f50b6a38a8e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -83,6 +83,7 @@ import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.delayed.BucketDelayedDeliveryTrackerFactory; import org.apache.pulsar.broker.delayed.DelayedDeliveryTrackerFactory; +import org.apache.pulsar.broker.loadbalance.extensions.ExtensibleLoadManagerImpl; import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateChannelImpl; import org.apache.pulsar.broker.loadbalance.extensions.channel.ServiceUnitStateCompactionStrategy; import org.apache.pulsar.broker.namespace.NamespaceService; @@ -1429,17 +1430,24 @@ public void deleteLedgerComplete(Object ctx) { } public CompletableFuture close() { - return close(false); + return close(false, false); + } + + @Override + public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { + return close(false, closeWithoutWaitingClientDisconnect); } /** * Close this topic - close all producers and subscriptions associated with this topic. * + * @param closeWithoutDisconnectingClients don't disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for client disconnect and forcefully close managed-ledger * @return Completable future indicating completion of close operation */ @Override - public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { + public CompletableFuture close( + boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect) { CompletableFuture closeFuture = new CompletableFuture<>(); lock.writeLock().lock(); @@ -1462,7 +1470,12 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect futures.add(transactionBuffer.closeAsync()); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); shadowReplicators.forEach((__, replicator) -> futures.add(replicator.disconnect())); - producers.values().forEach(producer -> futures.add(producer.disconnect())); + if (!closeWithoutDisconnectingClients) { + futures.add(ExtensibleLoadManagerImpl.getAssignedBrokerLookupData( + brokerService.getPulsar(), topic).thenAccept(lookupData -> + producers.values().forEach(producer -> futures.add(producer.disconnect(lookupData))) + )); + } if (topicPublishRateLimiter != null) { topicPublishRateLimiter.close(); } @@ -1491,14 +1504,22 @@ public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect ledger.asyncClose(new CloseCallback() { @Override public void closeComplete(Object ctx) { - // Everything is now closed, remove the topic from map - disposeTopic(closeFuture); + if (closeWithoutDisconnectingClients) { + closeFuture.complete(null); + } else { + // Everything is now closed, remove the topic from map + disposeTopic(closeFuture); + } } @Override public void closeFailed(ManagedLedgerException exception, Object ctx) { log.error("[{}] Failed to close managed ledger, proceeding anyway.", topic, exception); - disposeTopic(closeFuture); + if (closeWithoutDisconnectingClients) { + closeFuture.complete(null); + } else { + disposeTopic(closeFuture); + } } }, null); }).exceptionally(exception -> { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java index f499998fd3d6c..dce0d38df20fa 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java @@ -69,6 +69,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import org.apache.commons.lang3.mutable.MutableInt; import org.apache.commons.lang3.reflect.FieldUtils; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; @@ -97,6 +98,8 @@ import org.apache.pulsar.broker.namespace.NamespaceService; import org.apache.pulsar.broker.testcontext.PulsarTestContext; import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.impl.LookupService; import org.apache.pulsar.client.impl.TableViewImpl; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; @@ -118,6 +121,7 @@ import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; import org.testng.annotations.BeforeMethod; +import org.testng.annotations.DataProvider; import org.testng.annotations.Test; /** @@ -125,6 +129,7 @@ */ @Slf4j @Test(groups = "broker") +@SuppressWarnings("unchecked") public class ExtensibleLoadManagerImplTest extends MockedPulsarServiceBaseTest { private PulsarService pulsar1; @@ -383,7 +388,7 @@ public boolean test(NamespaceBundle namespaceBundle) { admin.namespaces().unloadNamespaceBundle(topicName.getNamespace(), bundle.getBundleRange(), dstBrokerUrl); Awaitility.await().untilAsserted(() -> { assertEquals(onloadCount.get(), 3); - assertEquals(unloadCount.get(), 2); + assertEquals(unloadCount.get(), 3); //one from releasing and one from owned }); assertEquals(admin.lookups().lookupTopic(topicName.toString()), dstBrokerServiceUrl); @@ -397,6 +402,110 @@ public boolean test(NamespaceBundle namespaceBundle) { assertTrue(ex.getMessage().contains("cannot be transfer to same broker")); } } + @DataProvider(name = "isPersistentTopicTest") + public Object[][] isPersistentTopicTest() { + return new Object[][] { { true }, { false }}; + } + @Test(timeOut = 30 * 1000, dataProvider = "isPersistentTopicTest") + public void testTransferClientReconnectionWithoutLookup(boolean isPersistentTopicTest) throws Exception { + String topicType = isPersistentTopicTest? "persistent" : "non-persistent"; + String topic = topicType + "://" + defaultTestNamespace + "/test-transfer-client-reconnect"; + TopicName topicName = TopicName.get(topic); + + AtomicInteger lookupCount = new AtomicInteger(); + var lookup = spyLookupService(lookupCount, topicName); + var producer = pulsarClient.newProducer().topic(topic).create(); + int lookupCountBeforeUnload = lookupCount.get(); + + NamespaceBundle bundle = getBundleAsync(pulsar1, TopicName.get(topic)).get(); + String broker = admin.lookups().lookupTopic(topic); + String dstBrokerUrl = pulsar1.getLookupServiceAddress(); + String dstBrokerServiceUrl; + if (broker.equals(pulsar1.getBrokerServiceUrl())) { + dstBrokerUrl = pulsar2.getLookupServiceAddress(); + dstBrokerServiceUrl = pulsar2.getBrokerServiceUrl(); + } else { + dstBrokerServiceUrl = pulsar1.getBrokerServiceUrl(); + } + checkOwnershipState(broker, bundle); + + final String finalDstBrokerUrl = dstBrokerUrl; + CompletableFuture.runAsync(() -> { + try { + admin.namespaces().unloadNamespaceBundle( + defaultTestNamespace, bundle.getBundleRange(), finalDstBrokerUrl); + } catch (PulsarAdminException e) { + throw new RuntimeException(e); + } + } + ); + + Awaitility.await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> { + try { + producer.send("hi".getBytes()); + String newOwner = admin.lookups().lookupTopic(topic); + assertEquals(dstBrokerServiceUrl, newOwner); + } catch (PulsarClientException e) { + throw new RuntimeException(e); + } catch (PulsarAdminException e) { + throw new RuntimeException(e); + } + }); + assertTrue(producer.isConnected()); + verify(lookup, times(lookupCountBeforeUnload)).getBroker(topicName); + producer.close(); + } + + + + @Test(timeOut = 30 * 1000, dataProvider = "isPersistentTopicTest") + public void testUnloadClientReconnectionWithLookup(boolean isPersistentTopicTest) throws Exception { + String topicType = isPersistentTopicTest? "persistent" : "non-persistent"; + String topic = topicType + "://" + defaultTestNamespace + "/test-unload-client-reconnect-" + + isPersistentTopicTest; + TopicName topicName = TopicName.get(topic); + + AtomicInteger lookupCount = new AtomicInteger(); + var lookup = spyLookupService(lookupCount, topicName); + + var producer = pulsarClient.newProducer().topic(topic).create(); + int lookupCountBeforeUnload = lookupCount.get(); + + NamespaceBundle bundle = getBundleAsync(pulsar1, TopicName.get(topic)).get(); + CompletableFuture.runAsync(() -> { + try { + admin.namespaces().unloadNamespaceBundle( + defaultTestNamespace, bundle.getBundleRange()); + } catch (PulsarAdminException e) { + throw new RuntimeException(e); + } + } + ); + MutableInt sendCount = new MutableInt(); + Awaitility.await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + try { + producer.send("hi".getBytes()); + assertEquals(sendCount.incrementAndGet(), 10); + } catch (PulsarClientException e) { + throw new RuntimeException(e); + } + }); + assertTrue(producer.isConnected()); + verify(lookup, times(lookupCountBeforeUnload + 1)).getBroker(topicName); + producer.close(); + } + + private LookupService spyLookupService(AtomicInteger lookupCount, TopicName topicName) + throws IllegalAccessException { + var lookup = spy((LookupService) + FieldUtils.readDeclaredField(pulsarClient, "lookup", true)); + FieldUtils.writeDeclaredField(pulsarClient, "lookup", lookup, true); + doAnswer(invocationOnMock -> { + lookupCount.incrementAndGet(); + return invocationOnMock.callRealMethod(); + }).when(lookup).getBroker(topicName); + return lookup; + } private void checkOwnershipState(String broker, NamespaceBundle bundle) throws ExecutionException, InterruptedException { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java index e8ccd7b01ca55..5636015188f1c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelTest.java @@ -100,6 +100,7 @@ import org.testng.annotations.Test; @Test(groups = "broker") +@SuppressWarnings("unchecked") public class ServiceUnitStateChannelTest extends MockedPulsarServiceBaseTest { private PulsarService pulsar1; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/OwnershipCacheTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/OwnershipCacheTest.java index 9e3d9e3a41340..c92127457aaf2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/OwnershipCacheTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/namespace/OwnershipCacheTest.java @@ -102,7 +102,7 @@ public void setup() throws Exception { nsService = mock(NamespaceService.class); brokerService = mock(BrokerService.class); doReturn(CompletableFuture.completedFuture(1)).when(brokerService) - .unloadServiceUnit(any(), anyBoolean(), anyLong(), any()); + .unloadServiceUnit(any(), anyBoolean(), anyBoolean(), anyLong(), any()); doReturn(config).when(pulsar).getConfiguration(); doReturn(nsService).when(pulsar).getNamespaceService(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index f3e8b2354b344..abdbe2f872cfe 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -34,6 +34,7 @@ import io.netty.util.concurrent.Promise; import java.net.InetSocketAddress; import java.net.SocketAddress; +import java.net.URI; import java.net.URISyntaxException; import java.nio.channels.ClosedChannelException; import java.util.Arrays; @@ -801,11 +802,29 @@ protected void handleError(CommandError error) { @Override protected void handleCloseProducer(CommandCloseProducer closeProducer) { - log.info("[{}] Broker notification of Closed producer: {}", remoteAddress, closeProducer.getProducerId()); final long producerId = closeProducer.getProducerId(); ProducerImpl producer = producers.remove(producerId); if (producer != null) { - producer.connectionClosed(this); + if (closeProducer.hasAssignedBrokerServiceUrl() || closeProducer.hasAssignedBrokerServiceUrlTls()) { + try { + final URI uri = new URI(producer.client.conf.isUseTls() + ? closeProducer.getAssignedBrokerServiceUrlTls() + : closeProducer.getAssignedBrokerServiceUrl()); + log.info("[{}] Broker notification of Closed producer: {}. Redirecting to {}.", + remoteAddress, closeProducer.getProducerId(), uri); + producer.getConnectionHandler().connectionClosed(this, 0L, Optional.of(uri)); + } catch (URISyntaxException e) { + log.error("[{}] Invalid redirect url {}/{} for {}", remoteAddress, + closeProducer.getAssignedBrokerServiceUrl(), + closeProducer.getAssignedBrokerServiceUrlTls(), + closeProducer.getRequestId()); + producer.connectionClosed(this); + } + } else { + log.info("[{}] Broker notification of Closed producer: {}.", + remoteAddress, closeProducer.getProducerId()); + producer.connectionClosed(this); + } } else { log.warn("Producer with id {} not found while closing producer ", producerId); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java index fc7c89c3ce693..550526be6d718 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java @@ -19,6 +19,8 @@ package org.apache.pulsar.client.impl; import java.net.InetSocketAddress; +import java.net.URI; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -66,6 +68,10 @@ protected ConnectionHandler(HandlerState state, Backoff backoff, Connection conn } protected void grabCnx() { + grabCnx(Optional.empty()); + } + + protected void grabCnx(Optional hostURI) { if (!duringConnect.compareAndSet(false, true)) { log.info("[{}] [{}] Skip grabbing the connection since there is a pending connection", state.topic, state.getHandlerName()); @@ -87,7 +93,12 @@ protected void grabCnx() { try { CompletableFuture cnxFuture; - if (state.redirectedClusterURI != null) { + if (hostURI.isPresent()) { + InetSocketAddress address = InetSocketAddress.createUnresolved( + hostURI.get().getHost(), + hostURI.get().getPort()); + cnxFuture = state.client.getConnection(address, address, randomKeyForSelectConnection); + } else if (state.redirectedClusterURI != null) { InetSocketAddress address = InetSocketAddress.createUnresolved(state.redirectedClusterURI.getHost(), state.redirectedClusterURI.getPort()); cnxFuture = state.client.getConnection(address, address, randomKeyForSelectConnection); @@ -149,6 +160,10 @@ void reconnectLater(Throwable exception) { } public void connectionClosed(ClientCnx cnx) { + connectionClosed(cnx, null, Optional.empty()); + } + + public void connectionClosed(ClientCnx cnx, Long initialConnectionDelayMs, Optional hostUrl) { lastConnectionClosedTimestamp = System.currentTimeMillis(); duringConnect.set(false); state.client.getCnxPool().releaseConnection(cnx); @@ -158,14 +173,14 @@ public void connectionClosed(ClientCnx cnx) { state.topic, state.getHandlerName(), state.getState()); return; } - long delayMs = backoff.next(); + long delayMs = initialConnectionDelayMs != null ? initialConnectionDelayMs.longValue() : backoff.next(); state.setState(State.Connecting); log.info("[{}] [{}] Closed connection {} -- Will try again in {} s", state.topic, state.getHandlerName(), cnx.channel(), delayMs / 1000.0); state.client.timer().newTimeout(timeout -> { log.info("[{}] [{}] Reconnecting after timeout", state.topic, state.getHandlerName()); - grabCnx(); + grabCnx(hostUrl); }, delayMs, TimeUnit.MILLISECONDS); } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index cf0cd820a6d10..ff116d2406b40 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -60,6 +60,7 @@ import org.apache.pulsar.common.api.proto.CommandAddSubscriptionToTxn; import org.apache.pulsar.common.api.proto.CommandAddSubscriptionToTxnResponse; import org.apache.pulsar.common.api.proto.CommandAuthChallenge; +import org.apache.pulsar.common.api.proto.CommandCloseProducer; import org.apache.pulsar.common.api.proto.CommandConnect; import org.apache.pulsar.common.api.proto.CommandConnected; import org.apache.pulsar.common.api.proto.CommandEndTxnOnPartitionResponse; @@ -761,11 +762,28 @@ public static ByteBuf newTopicMigrated(ResourceType type, long resourceId, Strin return serializeWithSize(cmd); } - public static ByteBuf newCloseProducer(long producerId, long requestId) { + public static ByteBuf newCloseProducer( + long producerId, long requestId) { + return newCloseProducer(producerId, requestId, null, null); + } + + public static ByteBuf newCloseProducer( + long producerId, long requestId, String assignedBrokerUrl, String assignedBrokerUrlTls) { BaseCommand cmd = localCmd(Type.CLOSE_PRODUCER); - cmd.setCloseProducer() - .setProducerId(producerId) - .setRequestId(requestId); + CommandCloseProducer commandCloseProducer = cmd.setCloseProducer() + .setProducerId(producerId) + .setRequestId(requestId); + + if (assignedBrokerUrl != null) { + commandCloseProducer + .setAssignedBrokerServiceUrl(assignedBrokerUrl); + } + + if (assignedBrokerUrlTls != null){ + commandCloseProducer + .setAssignedBrokerServiceUrlTls(assignedBrokerUrlTls); + } + return serializeWithSize(cmd); } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index afe193eeb7e9d..2c350aaf8a10e 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -641,6 +641,8 @@ message CommandTopicMigrated { message CommandCloseProducer { required uint64 producer_id = 1; required uint64 request_id = 2; + optional string assignedBrokerServiceUrl = 3; + optional string assignedBrokerServiceUrlTls = 4; } message CommandCloseConsumer { From 5605425d286db62446fef98320b6d9555000f46f Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Fri, 27 Oct 2023 14:51:43 -0700 Subject: [PATCH 2/4] Resolved comments --- .../extensions/ExtensibleLoadManagerImpl.java | 2 +- .../channel/ServiceUnitStateChannelImpl.java | 19 ++++++++----------- .../ExtensibleLoadManagerImplTest.java | 9 ++++++--- .../apache/pulsar/client/impl/ClientCnx.java | 9 ++++++--- .../pulsar/client/impl/ConnectionHandler.java | 4 ++-- 5 files changed, 23 insertions(+), 20 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java index 8e166f00bc95f..14c81a6a49215 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java @@ -324,7 +324,7 @@ public static CompletableFuture> getAssignedBrokerLoo } ); } catch (Throwable e) { - log.error("Failed to DestinationBrokerLookupData for topic:{}", topic, e); + log.error("Failed to lookup destination broker for topic:{}", topic, e); return CompletableFuture.completedFuture(Optional.empty()); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java index d0a74c0eaef8a..8f3ba216747dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java @@ -554,10 +554,7 @@ public Optional getAssigned(String serviceUnit) { return Optional.of(data.dstBroker()); } case Releasing -> { - if (data.dstBroker() != null) { - return Optional.of(data.dstBroker()); - } - return Optional.empty(); + return Optional.ofNullable(data.dstBroker()); } case Splitting -> { return Optional.of(data.sourceBroker()); @@ -770,20 +767,20 @@ private void handleOwnEvent(String serviceUnit, ServiceUnitStateData data) { if (getOwnerRequest != null) { getOwnerRequest.complete(data.dstBroker()); } - CompletableFuture ownFuture = null; + if (isTargetBroker(data.dstBroker())) { pulsar.getNamespaceService() .onNamespaceBundleOwned(LoadManagerShared.getNamespaceBundle(pulsar, serviceUnit)); lastOwnEventHandledAt = System.currentTimeMillis(); - ownFuture = CompletableFuture.completedFuture(null); + stateChangeListeners.notify(serviceUnit, data, null); + log(null, serviceUnit, data, null); } else if ((data.force() || isTransferCommand(data)) && isTargetBroker(data.sourceBroker())) { - ownFuture = closeServiceUnit(serviceUnit, false); + stateChangeListeners.notifyOnCompletion( + closeServiceUnit(serviceUnit, false), serviceUnit, data) + .whenComplete((__, e) -> log(e, serviceUnit, data, null)); } else { - ownFuture = CompletableFuture.completedFuture(null); + stateChangeListeners.notify(serviceUnit, data, null); } - - stateChangeListeners.notifyOnCompletion(ownFuture, serviceUnit, data) - .whenComplete((__, e) -> log(e, serviceUnit, data, null)); } private void handleAssignEvent(String serviceUnit, ServiceUnitStateData data) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java index dce0d38df20fa..97ee13c451d86 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java @@ -146,6 +146,8 @@ public class ExtensibleLoadManagerImplTest extends MockedPulsarServiceBaseTest { private final String defaultTestNamespace = "public/test"; + private LookupService lookupService; + @BeforeClass @Override public void setup() throws Exception { @@ -194,6 +196,7 @@ public void setup() throws Exception { admin.namespaces().createNamespace(defaultTestNamespace); admin.namespaces().setNamespaceReplicationClusters(defaultTestNamespace, Sets.newHashSet(this.conf.getClusterName())); + lookupService = (LookupService) FieldUtils.readDeclaredField(pulsarClient, "lookup", true); } } @@ -207,9 +210,10 @@ protected void cleanup() throws Exception { } @BeforeMethod(alwaysRun = true) - protected void initializeState() throws PulsarAdminException { + protected void initializeState() throws PulsarAdminException, IllegalAccessException { admin.namespaces().unload(defaultTestNamespace); reset(primaryLoadManager, secondaryLoadManager); + FieldUtils.writeDeclaredField(pulsarClient, "lookup", lookupService, true); } @Test @@ -497,8 +501,7 @@ public void testUnloadClientReconnectionWithLookup(boolean isPersistentTopicTest private LookupService spyLookupService(AtomicInteger lookupCount, TopicName topicName) throws IllegalAccessException { - var lookup = spy((LookupService) - FieldUtils.readDeclaredField(pulsarClient, "lookup", true)); + var lookup = spy(lookupService); FieldUtils.writeDeclaredField(pulsarClient, "lookup", lookup, true); doAnswer(invocationOnMock -> { lookupCount.incrementAndGet(); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index abdbe2f872cfe..794ef318bba95 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -812,11 +812,14 @@ protected void handleCloseProducer(CommandCloseProducer closeProducer) { : closeProducer.getAssignedBrokerServiceUrl()); log.info("[{}] Broker notification of Closed producer: {}. Redirecting to {}.", remoteAddress, closeProducer.getProducerId(), uri); - producer.getConnectionHandler().connectionClosed(this, 0L, Optional.of(uri)); + producer.getConnectionHandler().connectionClosed( + this, Optional.of(0L), Optional.of(uri)); } catch (URISyntaxException e) { log.error("[{}] Invalid redirect url {}/{} for {}", remoteAddress, - closeProducer.getAssignedBrokerServiceUrl(), - closeProducer.getAssignedBrokerServiceUrlTls(), + closeProducer.hasAssignedBrokerServiceUrl() + ? closeProducer.getAssignedBrokerServiceUrl() : "", + closeProducer.hasAssignedBrokerServiceUrlTls() + ? closeProducer.getAssignedBrokerServiceUrlTls() : "", closeProducer.getRequestId()); producer.connectionClosed(this); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java index 550526be6d718..267329015895f 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java @@ -163,7 +163,7 @@ public void connectionClosed(ClientCnx cnx) { connectionClosed(cnx, null, Optional.empty()); } - public void connectionClosed(ClientCnx cnx, Long initialConnectionDelayMs, Optional hostUrl) { + public void connectionClosed(ClientCnx cnx, Optional initialConnectionDelayMs, Optional hostUrl) { lastConnectionClosedTimestamp = System.currentTimeMillis(); duringConnect.set(false); state.client.getCnxPool().releaseConnection(cnx); @@ -173,7 +173,7 @@ public void connectionClosed(ClientCnx cnx, Long initialConnectionDelayMs, Optio state.topic, state.getHandlerName(), state.getState()); return; } - long delayMs = initialConnectionDelayMs != null ? initialConnectionDelayMs.longValue() : backoff.next(); + long delayMs = initialConnectionDelayMs.orElse(backoff.next()); state.setState(State.Connecting); log.info("[{}] [{}] Closed connection {} -- Will try again in {} s", state.topic, state.getHandlerName(), cnx.channel(), From f7310b3ec31d3f0ebca6a065fe7668353aea5a6f Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Fri, 3 Nov 2023 15:08:07 -0700 Subject: [PATCH 3/4] resolved comments --- .../extensions/ExtensibleLoadManagerImpl.java | 4 ++++ .../channel/ServiceUnitStateChannelImpl.java | 12 +++++----- .../pulsar/broker/namespace/OwnedBundle.java | 2 +- .../pulsar/broker/service/BrokerService.java | 10 ++++----- .../pulsar/broker/service/ServerCnx.java | 7 +++--- .../apache/pulsar/broker/service/Topic.java | 2 +- .../nonpersistent/NonPersistentTopic.java | 10 ++++----- .../service/persistent/PersistentTopic.java | 22 +++++++++---------- .../ExtensibleLoadManagerImplTest.java | 20 ++++++++++++----- .../pulsar/client/impl/ConnectionHandler.java | 2 +- 10 files changed, 52 insertions(+), 39 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java index 14c81a6a49215..67bab9b12ffb1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImpl.java @@ -267,6 +267,10 @@ public static boolean isLoadManagerExtensionEnabled(ServiceConfiguration conf) { return ExtensibleLoadManagerImpl.class.getName().equals(conf.getLoadManagerClassName()); } + public static boolean isLoadManagerExtensionEnabled(PulsarService pulsar) { + return pulsar.getLoadManager().get() instanceof ExtensibleLoadManagerImpl; + } + public static ExtensibleLoadManagerImpl get(LoadManager loadManager) { if (!(loadManager instanceof ExtensibleLoadManagerWrapper loadManagerWrapper)) { throw new IllegalArgumentException("The load manager should be 'ExtensibleLoadManagerWrapper'."); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java index 8f3ba216747dc..e7567bccc131e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/extensions/channel/ServiceUnitStateChannelImpl.java @@ -776,7 +776,7 @@ private void handleOwnEvent(String serviceUnit, ServiceUnitStateData data) { log(null, serviceUnit, data, null); } else if ((data.force() || isTransferCommand(data)) && isTargetBroker(data.sourceBroker())) { stateChangeListeners.notifyOnCompletion( - closeServiceUnit(serviceUnit, false), serviceUnit, data) + closeServiceUnit(serviceUnit, true), serviceUnit, data) .whenComplete((__, e) -> log(e, serviceUnit, data, null)); } else { stateChangeListeners.notify(serviceUnit, data, null); @@ -799,11 +799,11 @@ private void handleReleaseEvent(String serviceUnit, ServiceUnitStateData data) { if (isTransferCommand(data)) { next = new ServiceUnitStateData( Assigning, data.dstBroker(), data.sourceBroker(), getNextVersionId(data)); - unloadFuture = closeServiceUnit(serviceUnit, true); + unloadFuture = closeServiceUnit(serviceUnit, false); } else { next = new ServiceUnitStateData( Free, null, data.sourceBroker(), getNextVersionId(data)); - unloadFuture = closeServiceUnit(serviceUnit, false); + unloadFuture = closeServiceUnit(serviceUnit, true); } stateChangeListeners.notifyOnCompletion(unloadFuture .thenCompose(__ -> pubAsync(serviceUnit, next)), serviceUnit, data) @@ -903,13 +903,13 @@ private CompletableFuture deferGetOwnerRequest(String serviceUnit) { } } - private CompletableFuture closeServiceUnit(String serviceUnit, boolean closeWithoutDisconnectingClients) { + private CompletableFuture closeServiceUnit(String serviceUnit, boolean disconnectClients) { long startTime = System.nanoTime(); MutableInt unloadedTopics = new MutableInt(); NamespaceBundle bundle = LoadManagerShared.getNamespaceBundle(pulsar, serviceUnit); return pulsar.getBrokerService().unloadServiceUnit( bundle, - closeWithoutDisconnectingClients, + disconnectClients, true, pulsar.getConfig().getNamespaceBundleUnloadingTimeoutMs(), TimeUnit.MILLISECONDS) @@ -918,7 +918,7 @@ private CompletableFuture closeServiceUnit(String serviceUnit, boolean return numUnloadedTopics; }) .whenComplete((__, ex) -> { - if (!closeWithoutDisconnectingClients) { + if (disconnectClients) { // clean up topics that failed to unload from the broker ownership cache pulsar.getBrokerService().cleanUnloadedTopicFromCache(bundle); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java index a87d45395db01..cdedac1136e4d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/namespace/OwnedBundle.java @@ -136,7 +136,7 @@ public CompletableFuture handleUnloadRequest(PulsarService pulsar, long ti return pulsar.getNamespaceService().getOwnershipCache() .updateBundleState(this.bundle, false) .thenCompose(v -> pulsar.getBrokerService().unloadServiceUnit( - bundle, false, closeWithoutWaitingClientDisconnect, timeout, timeoutUnit)) + bundle, true, closeWithoutWaitingClientDisconnect, timeout, timeoutUnit)) .handle((numUnloadedTopics, ex) -> { if (ex != null) { // ignore topic-close failure to unload bundle diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index ed9a4d20c87ef..17a4ade60e6ab 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -2216,10 +2216,10 @@ public CompletableFuture checkTopicNsOwnership(final String topic) { } public CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, - boolean closeWithoutDisconnectingClients, + boolean disconnectClients, boolean closeWithoutWaitingClientDisconnect, long timeout, TimeUnit unit) { CompletableFuture future = unloadServiceUnit( - serviceUnit, closeWithoutDisconnectingClients, closeWithoutWaitingClientDisconnect); + serviceUnit, disconnectClients, closeWithoutWaitingClientDisconnect); ScheduledFuture taskTimeout = executor().schedule(() -> { if (!future.isDone()) { log.warn("Unloading of {} has timed out", serviceUnit); @@ -2236,13 +2236,13 @@ public CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, * Unload all the topic served by the broker service under the given service unit. * * @param serviceUnit - * @param closeWithoutDisconnectingClients don't disconnect clients + * @param disconnectClients disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for clients to disconnect * and forcefully close managed-ledger * @return */ private CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit, - boolean closeWithoutDisconnectingClients, + boolean disconnectClients, boolean closeWithoutWaitingClientDisconnect) { List> closeFutures = new ArrayList<>(); topics.forEach((name, topicFuture) -> { @@ -2267,7 +2267,7 @@ private CompletableFuture unloadServiceUnit(NamespaceBundle serviceUnit } closeFutures.add(topicFuture .thenCompose(t -> t.isPresent() ? t.get().close( - closeWithoutDisconnectingClients, closeWithoutWaitingClientDisconnect) + disconnectClients, closeWithoutWaitingClientDisconnect) : CompletableFuture.completedFuture(null))); } }); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index c48b31518d1b8..e0a063cc65e05 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1719,11 +1719,10 @@ protected void handleSend(CommandSend send, ByteBuf headersAndPayload) { printSendCommandDebug(send, headersAndPayload); } - - ServiceConfiguration conf = getBrokerService().pulsar().getConfiguration(); + PulsarService pulsar = getBrokerService().pulsar(); if (producer.getTopic().isFenced() - && ExtensibleLoadManagerImpl.isLoadManagerExtensionEnabled(conf)) { - long ignoredMsgCount = ExtensibleLoadManagerImpl.get(getBrokerService().pulsar()) + && ExtensibleLoadManagerImpl.isLoadManagerExtensionEnabled(pulsar)) { + long ignoredMsgCount = ExtensibleLoadManagerImpl.get(pulsar) .getIgnoredSendMsgCounter().incrementAndGet(); if (log.isDebugEnabled()) { log.debug("Ignored send msg from:{}:{} to fenced topic:{} during unloading." diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index ead6dd6982d1a..3e5aff03710cf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -197,7 +197,7 @@ CompletableFuture createSubscription(String subscriptionName, Init CompletableFuture close(boolean closeWithoutWaitingClientDisconnect); CompletableFuture close( - boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect); + boolean disconnectClients, boolean closeWithoutWaitingClientDisconnect); void checkGC(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index d232a1451787d..c1cd322055340 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -489,19 +489,19 @@ private CompletableFuture delete(boolean failIfHasSubscriptions, boolean c @Override public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { - return close(false, closeWithoutWaitingClientDisconnect); + return close(true, closeWithoutWaitingClientDisconnect); } /** * Close this topic - close all producers and subscriptions associated with this topic. * - * @param closeWithoutDisconnectingClients don't disconnect clients + * @param disconnectClients disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for client disconnect and forcefully close managed-ledger * @return Completable future indicating completion of close operation */ @Override public CompletableFuture close( - boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect) { + boolean disconnectClients, boolean closeWithoutWaitingClientDisconnect) { CompletableFuture closeFuture = new CompletableFuture<>(); lock.writeLock().lock(); @@ -520,7 +520,7 @@ public CompletableFuture close( List> futures = new ArrayList<>(); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); - if (!closeWithoutDisconnectingClients) { + if (disconnectClients) { futures.add(ExtensibleLoadManagerImpl.getAssignedBrokerLookupData( brokerService.getPulsar(), topic).thenAccept(lookupData -> producers.values().forEach(producer -> futures.add(producer.disconnect(lookupData))) @@ -554,7 +554,7 @@ public CompletableFuture close( // so, execute it in different thread brokerService.executor().execute(() -> { - if (!closeWithoutDisconnectingClients) { + if (disconnectClients) { brokerService.removeTopicFromCache(NonPersistentTopic.this); unregisterTopicPolicyListener(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 13f50b6a38a8e..8aa5df8b51096 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -1430,24 +1430,24 @@ public void deleteLedgerComplete(Object ctx) { } public CompletableFuture close() { - return close(false, false); + return close(true, false); } @Override public CompletableFuture close(boolean closeWithoutWaitingClientDisconnect) { - return close(false, closeWithoutWaitingClientDisconnect); + return close(true, closeWithoutWaitingClientDisconnect); } /** * Close this topic - close all producers and subscriptions associated with this topic. * - * @param closeWithoutDisconnectingClients don't disconnect clients + * @param disconnectClients disconnect clients * @param closeWithoutWaitingClientDisconnect don't wait for client disconnect and forcefully close managed-ledger * @return Completable future indicating completion of close operation */ @Override public CompletableFuture close( - boolean closeWithoutDisconnectingClients, boolean closeWithoutWaitingClientDisconnect) { + boolean disconnectClients, boolean closeWithoutWaitingClientDisconnect) { CompletableFuture closeFuture = new CompletableFuture<>(); lock.writeLock().lock(); @@ -1470,7 +1470,7 @@ public CompletableFuture close( futures.add(transactionBuffer.closeAsync()); replicators.forEach((cluster, replicator) -> futures.add(replicator.disconnect())); shadowReplicators.forEach((__, replicator) -> futures.add(replicator.disconnect())); - if (!closeWithoutDisconnectingClients) { + if (disconnectClients) { futures.add(ExtensibleLoadManagerImpl.getAssignedBrokerLookupData( brokerService.getPulsar(), topic).thenAccept(lookupData -> producers.values().forEach(producer -> futures.add(producer.disconnect(lookupData))) @@ -1504,21 +1504,21 @@ public CompletableFuture close( ledger.asyncClose(new CloseCallback() { @Override public void closeComplete(Object ctx) { - if (closeWithoutDisconnectingClients) { - closeFuture.complete(null); - } else { + if (disconnectClients) { // Everything is now closed, remove the topic from map disposeTopic(closeFuture); + } else { + closeFuture.complete(null); } } @Override public void closeFailed(ManagedLedgerException exception, Object ctx) { log.error("[{}] Failed to close managed ledger, proceeding anyway.", topic, exception); - if (closeWithoutDisconnectingClients) { - closeFuture.complete(null); - } else { + if (disconnectClients) { disposeTopic(closeFuture); + } else { + closeFuture.complete(null); } } }, null); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java index 97ee13c451d86..2acbf2367064f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/loadbalance/extensions/ExtensibleLoadManagerImplTest.java @@ -635,9 +635,11 @@ public void testCheckOwnershipPresentWithSystemNamespace() throws Exception { @Test public void testMoreThenOneFilter() throws Exception { - TopicName topicName = TopicName.get(defaultTestNamespace + "/test-filter-has-exception"); + // Use a different namespace to avoid flaky test failures + // from unloading the default namespace and the following topic policy lookups at the init state step + String namespace = "public/my-namespace"; + TopicName topicName = TopicName.get(namespace + "/test-filter-has-exception"); NamespaceBundle bundle = getBundleAsync(pulsar1, topicName).get(); - String lookupServiceAddress1 = pulsar1.getLookupServiceAddress(); doReturn(List.of(new MockBrokerFilter() { @Override @@ -655,10 +657,18 @@ public CompletableFuture> filterAsync(Map brokerLookupData = primaryLoadManager.assign(Optional.empty(), bundle).get(); - assertTrue(brokerLookupData.isPresent()); - assertEquals(brokerLookupData.get().getWebServiceUrl(), pulsar2.getWebServiceAddress()); + Awaitility.waitAtMost(5, TimeUnit.SECONDS).untilAsserted(() -> { + assertTrue(brokerLookupData.isPresent()); + assertEquals(brokerLookupData.get().getWebServiceUrl(), pulsar2.getWebServiceAddress()); + assertEquals(brokerLookupData.get().getPulsarServiceUrl(), + pulsar1.getAdminClient().lookups().lookupTopic(topicName.toString())); + assertEquals(brokerLookupData.get().getPulsarServiceUrl(), + pulsar2.getAdminClient().lookups().lookupTopic(topicName.toString())); + }); + + admin.namespaces().deleteNamespace(namespace, true); } @Test diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java index 267329015895f..d319b2ba0c63b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConnectionHandler.java @@ -160,7 +160,7 @@ void reconnectLater(Throwable exception) { } public void connectionClosed(ClientCnx cnx) { - connectionClosed(cnx, null, Optional.empty()); + connectionClosed(cnx, Optional.empty(), Optional.empty()); } public void connectionClosed(ClientCnx cnx, Optional initialConnectionDelayMs, Optional hostUrl) { From 17a2a3579bbc83959282c00bcfd45638ff90cdef Mon Sep 17 00:00:00 2001 From: Heesung Sohn Date: Sat, 4 Nov 2023 14:58:43 -0700 Subject: [PATCH 4/4] resolved comment --- .../main/java/org/apache/pulsar/client/impl/ClientCnx.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java index 794ef318bba95..ba79f1b824765 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ClientCnx.java @@ -814,13 +814,13 @@ protected void handleCloseProducer(CommandCloseProducer closeProducer) { remoteAddress, closeProducer.getProducerId(), uri); producer.getConnectionHandler().connectionClosed( this, Optional.of(0L), Optional.of(uri)); - } catch (URISyntaxException e) { + } catch (Throwable e) { log.error("[{}] Invalid redirect url {}/{} for {}", remoteAddress, closeProducer.hasAssignedBrokerServiceUrl() ? closeProducer.getAssignedBrokerServiceUrl() : "", closeProducer.hasAssignedBrokerServiceUrlTls() ? closeProducer.getAssignedBrokerServiceUrlTls() : "", - closeProducer.getRequestId()); + closeProducer.getRequestId(), e); producer.connectionClosed(this); } } else {