From 8988f68fa2e35b058a56cca936f0cdfd2e843050 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 14 Jun 2023 01:05:20 +0800 Subject: [PATCH 1/2] [fix] [broker]Consumption stuck due to orphan consumers after connection reconnect --- .../pulsar/broker/service/ServerCnx.java | 32 ++++++++++--------- 1 file changed, 17 insertions(+), 15 deletions(-) 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 98c0e5b497998..48473e81be917 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 @@ -1243,20 +1243,11 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { return topic.subscribe(option); } }) - .thenAccept(consumer -> { - if (consumerFuture.complete(consumer)) { - log.info("[{}] Created subscription on topic {} / {}", - remoteAddress, topicName, subscriptionName); - commandSender.sendSuccessResponse(requestId); - if (brokerInterceptor != null) { - try { - brokerInterceptor.consumerCreated(this, consumer, metadata); - } catch (Throwable t) { - log.error("Exception occur when intercept consumer created.", t); - } - } - } else { - // The consumer future was completed before by a close command + .thenAcceptAsync(consumer -> { + if (consumerFuture.complete(consumer) || !isActive()) { + // Two cases: + // 1. The consumer future was completed before by a close command. + // 2. The consumer future was completed after the ServerCnx closed. try { consumer.close(); log.info("[{}] Cleared consumer created after timeout on client side {}", @@ -1268,9 +1259,20 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { remoteAddress, consumer, e.getMessage()); } consumers.remove(consumerId, consumerFuture); + } else { + log.info("[{}] Created subscription on topic {} / {}", + remoteAddress, topicName, subscriptionName); + commandSender.sendSuccessResponse(requestId); + if (brokerInterceptor != null) { + try { + brokerInterceptor.consumerCreated(this, consumer, metadata); + } catch (Throwable t) { + log.error("Exception occur when intercept consumer created.", t); + } + } } - }) + }, ctx().channel().eventLoop()) .exceptionally(exception -> { if (exception.getCause() instanceof ConsumerBusyException) { if (log.isDebugEnabled()) { From af415baa7e8f620c65eb167396877acf568ca6eb Mon Sep 17 00:00:00 2001 From: fengyubiao <9947090@qq.com> Date: Wed, 14 Jun 2023 07:43:20 +0800 Subject: [PATCH 2/2] Update pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java Co-authored-by: Penghui Li --- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 48473e81be917..7d2e518f4a0ff 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 @@ -1244,7 +1244,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { } }) .thenAcceptAsync(consumer -> { - if (consumerFuture.complete(consumer) || !isActive()) { + if (!consumerFuture.complete(consumer) || !isActive()) { // Two cases: // 1. The consumer future was completed before by a close command. // 2. The consumer future was completed after the ServerCnx closed.