From 41d50cf44f7ed837712451a06ea8085e8bacc585 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 13 Jun 2023 16:36:49 +0800 Subject: [PATCH 01/12] [improve][log] Improve dispatcher log to trace io error when send data to client --- .../java/org/apache/pulsar/broker/service/Consumer.java | 6 ++++++ .../java/org/apache/pulsar/broker/service/ServerCnx.java | 3 ++- .../persistent/PersistentDispatcherMultipleConsumers.java | 2 +- 3 files changed, 9 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 275d685280865..17241da5f84b3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -361,6 +361,12 @@ public Future sendMessages(final List entries, EntryBatch msgOutCounter.add(totalMessages); bytesOutCounter.add(totalBytes); chunkedMessageRate.recordMultipleEvents(totalChunkedMessages, 0); + } else { + log.warn("[{}-{}] Sent messages to client failed by IO exception[{}], these messages(messages count:" + + " {}) will be redelivered after the heartbeat check fails. If the next heartbeat" + + " check is successful, these messages will not be stuck until the client reconnect" + + " or the topic is reloaded. Consumer: {}", + topicName, subscription, totalMessages, consumerId, this.toString(), status.cause()); } }); return writeAndFlushPromise; 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 194f8593e4aba..ad1b381c45500 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 @@ -1160,7 +1160,8 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { remoteAddress, getPrincipal()); } - log.info("[{}] Subscribing on topic {} / {}", remoteAddress, topicName, subscriptionName); + log.info("[{}] Subscribing on topic {} / {}. consumerId: {}", remoteAddress, topicName, + subscriptionName, consumerId); try { Metadata.validateMetadata(metadata, service.getPulsar().getConfiguration().getMaxConsumerMetadataSize()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index b3d48252efe58..5c5e688518eed 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -200,9 +200,9 @@ protected boolean isConsumersExceededOnSubscription() { public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { // decrement unack-message count for removed consumer addUnAckedMessages(-consumer.getUnackedMessages()); + log.info("Removed consumer {} with pending {} acks", consumer, consumer.getPendingAcks().size()); if (consumerSet.removeAll(consumer) == 1) { consumerList.remove(consumer); - log.info("Removed consumer {} with pending {} acks", consumer, consumer.getPendingAcks().size()); if (consumerList.isEmpty()) { cancelPendingRead(); From bde1358099580ef7d70fd1010c35b0a45c1b7095 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 13 Jun 2023 16:44:35 +0800 Subject: [PATCH 02/12] - --- .../main/java/org/apache/pulsar/broker/service/Consumer.java | 2 +- .../persistent/PersistentDispatcherMultipleConsumers.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 17241da5f84b3..c6c5ee37c7889 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -364,7 +364,7 @@ public Future sendMessages(final List entries, EntryBatch } else { log.warn("[{}-{}] Sent messages to client failed by IO exception[{}], these messages(messages count:" + " {}) will be redelivered after the heartbeat check fails. If the next heartbeat" - + " check is successful, these messages will not be stuck until the client reconnect" + + " check is successful, these messages will be stuck until the client reconnect" + " or the topic is reloaded. Consumer: {}", topicName, subscription, totalMessages, consumerId, this.toString(), status.cause()); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 5c5e688518eed..b3d48252efe58 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -200,9 +200,9 @@ protected boolean isConsumersExceededOnSubscription() { public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { // decrement unack-message count for removed consumer addUnAckedMessages(-consumer.getUnackedMessages()); - log.info("Removed consumer {} with pending {} acks", consumer, consumer.getPendingAcks().size()); if (consumerSet.removeAll(consumer) == 1) { consumerList.remove(consumer); + log.info("Removed consumer {} with pending {} acks", consumer, consumer.getPendingAcks().size()); if (consumerList.isEmpty()) { cancelPendingRead(); From 6e6487a342f888a44c4f2949354f727019ba5ef9 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 13 Jun 2023 16:45:12 +0800 Subject: [PATCH 03/12] - --- .../main/java/org/apache/pulsar/broker/service/Consumer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index c6c5ee37c7889..5429736541905 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -362,7 +362,7 @@ public Future sendMessages(final List entries, EntryBatch bytesOutCounter.add(totalBytes); chunkedMessageRate.recordMultipleEvents(totalChunkedMessages, 0); } else { - log.warn("[{}-{}] Sent messages to client failed by IO exception[{}], these messages(messages count:" + log.warn("[{}-{}] Sent messages to client fail by IO exception[{}], these messages(messages count:" + " {}) will be redelivered after the heartbeat check fails. If the next heartbeat" + " check is successful, these messages will be stuck until the client reconnect" + " or the topic is reloaded. Consumer: {}", From 0b5174a4158d1d1ae382b7ad33dbc27157b7188a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 14 Jun 2023 00:04:36 +0800 Subject: [PATCH 04/12] address comments --- .../main/java/org/apache/pulsar/broker/service/Consumer.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 5429736541905..2186b6bf88a2a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -366,7 +366,8 @@ public Future sendMessages(final List entries, EntryBatch + " {}) will be redelivered after the heartbeat check fails. If the next heartbeat" + " check is successful, these messages will be stuck until the client reconnect" + " or the topic is reloaded. Consumer: {}", - topicName, subscription, totalMessages, consumerId, this.toString(), status.cause()); + topicName, subscription, status.cause() == null ? "" : status.cause().getMessage(), + totalMessages, this.toString(), status.cause()); } }); return writeAndFlushPromise; From 8ef79c40967b727464bd2d2dff68e3fa567f937f Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 14 Jun 2023 01:48:20 +0800 Subject: [PATCH 05/12] improve log --- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 ad1b381c45500..bd8a35e23385c 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 @@ -1160,8 +1160,8 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { remoteAddress, getPrincipal()); } - log.info("[{}] Subscribing on topic {} / {}. consumerId: {}", remoteAddress, topicName, - subscriptionName, consumerId); + log.info("[{}] Subscribing on topic {} / {}. consumerId: {}", this.ctx().channel().toString(), + topicName, subscriptionName, consumerId); try { Metadata.validateMetadata(metadata, service.getPulsar().getConfiguration().getMaxConsumerMetadataSize()); From 3e8cabc919337f962c9b7d1245ab313599870a4a Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 14 Jun 2023 13:01:40 +0800 Subject: [PATCH 06/12] add warn log when discard ack command --- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 4 ++++ 1 file changed, 4 insertions(+) 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 bd8a35e23385c..03fce402b2eee 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 @@ -1784,6 +1784,10 @@ protected void handleAck(CommandAck ack) { } return null; }); + } else { + log.warn("Consumer future is not complete(not complete or error), but received command ack. so discard" + + " this command. consumerId: {}, cnx: {}, messageIdCount: {}", ack.getConsumerId(), + this.ctx().channel().toString(), ack.getMessageIdsCount()); } } From a94c201be4024460209e94378ebe2afe74cc505b Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 15 Jun 2023 12:12:49 +0800 Subject: [PATCH 07/12] change log level of discard ack to debug --- .../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 03fce402b2eee..0305973702aa8 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 @@ -1785,7 +1785,7 @@ protected void handleAck(CommandAck ack) { return null; }); } else { - log.warn("Consumer future is not complete(not complete or error), but received command ack. so discard" + log.debug("Consumer future is not complete(not complete or error), but received command ack. so discard" + " this command. consumerId: {}, cnx: {}, messageIdCount: {}", ack.getConsumerId(), this.ctx().channel().toString(), ack.getMessageIdsCount()); } From 3b3ce9450f9ee9b2e9c32df78bd45a054efbddd4 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 15 Jun 2023 12:33:55 +0800 Subject: [PATCH 08/12] do health check manually when sent messages to client fail --- .../apache/pulsar/broker/service/Consumer.java | 2 ++ .../apache/pulsar/broker/service/ServerCnx.java | 17 +++++++++++++++++ .../pulsar/broker/service/TransportCnx.java | 7 +++++++ 3 files changed, 26 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 2186b6bf88a2a..164f11b559a0b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -368,6 +368,8 @@ public Future sendMessages(final List entries, EntryBatch + " or the topic is reloaded. Consumer: {}", topicName, subscription, status.cause() == null ? "" : status.cause().getMessage(), totalMessages, this.toString(), status.cause()); + // If the health check fail, this connection will be closed. + cnx.healthCheckManually(); } }); return writeAndFlushPromise; 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 0305973702aa8..bdcfa8161abeb 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 @@ -57,6 +57,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.regex.Pattern; import java.util.stream.Collectors; import javax.naming.AuthenticationException; @@ -85,6 +86,7 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionNotFoundException; import org.apache.pulsar.broker.service.BrokerServiceException.TopicNotFoundException; +import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.service.schema.SchemaRegistryService; @@ -244,6 +246,10 @@ public class ServerCnx extends PulsarHandler implements TransportCnx { private final long connectionLivenessCheckTimeoutMillis; + protected static final AtomicLongFieldUpdater LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER = + AtomicLongFieldUpdater.newUpdater(ServerCnx.class, "lastManualHeartbeatCheckTime"); + private volatile long lastManualHeartbeatCheckTime = 0L; + // Number of bytes pending to be published from a single specific IO thread. private static final FastThreadLocal pendingBytesPerThread = new FastThreadLocal() { @Override @@ -3423,6 +3429,17 @@ public CompletableFuture checkConnectionLiveness() { } } + @Override + public void healthCheckManually() { + long lastCheckTime = LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER.get(this); + if (System.currentTimeMillis() - lastCheckTime < 5000) { + return; + } + if (LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER.compareAndSet(this, lastCheckTime, System.currentTimeMillis())) { + sendPing(); + } + } + @Override protected void messageReceived() { super.messageReceived(); 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 94f934fec681e..dffcea3bb1eba 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 @@ -90,4 +90,11 @@ public interface TransportCnx { * is null if the connection liveness check is disabled. */ CompletableFuture checkConnectionLiveness(); + + /** + * If an IO exception is found, you can call this method immediately. The connection will be closed if the + * heartbeat check fails. + * Note: to avoid calling this method frequently, this method discards non-first calls within 5 seconds. + */ + void healthCheckManually(); } From e4d0335610fc6365ff48859508e40f65c36deddd Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 15 Jun 2023 12:46:04 +0800 Subject: [PATCH 09/12] remove unnecessary imports --- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 1 - 1 file changed, 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 bdcfa8161abeb..9a449c02e9342 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 @@ -86,7 +86,6 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ServiceUnitNotReadyException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionNotFoundException; import org.apache.pulsar.broker.service.BrokerServiceException.TopicNotFoundException; -import org.apache.pulsar.broker.service.persistent.PersistentDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.persistent.PersistentSubscription; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.broker.service.schema.SchemaRegistryService; From 6e45b238cae509a754f80d4c9ace184c5c144f26 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Thu, 15 Jun 2023 23:59:35 +0800 Subject: [PATCH 10/12] add a check if (log.isDebugEnabled()) { --- .../java/org/apache/pulsar/broker/service/ServerCnx.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 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 9a449c02e9342..4844e8340f963 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 @@ -1790,9 +1790,11 @@ protected void handleAck(CommandAck ack) { return null; }); } else { - log.debug("Consumer future is not complete(not complete or error), but received command ack. so discard" - + " this command. consumerId: {}, cnx: {}, messageIdCount: {}", ack.getConsumerId(), - this.ctx().channel().toString(), ack.getMessageIdsCount()); + if (log.isDebugEnabled()) { + log.debug("Consumer future is not complete(not complete or error), but received command ack. so discard" + + " this command. consumerId: {}, cnx: {}, messageIdCount: {}", ack.getConsumerId(), + this.ctx().channel().toString(), ack.getMessageIdsCount()); + } } } From 19ebf6b58c695708cf914c232f7c9477b313aa9e Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 27 Jun 2023 19:22:02 +0800 Subject: [PATCH 11/12] address comments --- .../pulsar/broker/service/Consumer.java | 13 ++--- .../pulsar/broker/service/ServerCnx.java | 16 ------ .../pulsar/broker/service/TransportCnx.java | 7 --- .../api/SimpleProducerConsumerTest.java | 53 +++++++++++++++++++ 4 files changed, 58 insertions(+), 31 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java index 164f11b559a0b..176f033a6dcf1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Consumer.java @@ -362,14 +362,11 @@ public Future sendMessages(final List entries, EntryBatch bytesOutCounter.add(totalBytes); chunkedMessageRate.recordMultipleEvents(totalChunkedMessages, 0); } else { - log.warn("[{}-{}] Sent messages to client fail by IO exception[{}], these messages(messages count:" - + " {}) will be redelivered after the heartbeat check fails. If the next heartbeat" - + " check is successful, these messages will be stuck until the client reconnect" - + " or the topic is reloaded. Consumer: {}", - topicName, subscription, status.cause() == null ? "" : status.cause().getMessage(), - totalMessages, this.toString(), status.cause()); - // If the health check fail, this connection will be closed. - cnx.healthCheckManually(); + if (log.isDebugEnabled()) { + log.debug("[{}-{}] Sent messages to client fail by IO exception[{}], close the connection" + + " immediately. Consumer: {}", topicName, subscription, + status.cause() == null ? "" : status.cause().getMessage(), this.toString()); + } } }); return writeAndFlushPromise; 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 4844e8340f963..651aaf2cba2ea 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 @@ -57,7 +57,6 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicLongFieldUpdater; import java.util.regex.Pattern; import java.util.stream.Collectors; import javax.naming.AuthenticationException; @@ -245,10 +244,6 @@ public class ServerCnx extends PulsarHandler implements TransportCnx { private final long connectionLivenessCheckTimeoutMillis; - protected static final AtomicLongFieldUpdater LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER = - AtomicLongFieldUpdater.newUpdater(ServerCnx.class, "lastManualHeartbeatCheckTime"); - private volatile long lastManualHeartbeatCheckTime = 0L; - // Number of bytes pending to be published from a single specific IO thread. private static final FastThreadLocal pendingBytesPerThread = new FastThreadLocal() { @Override @@ -3430,17 +3425,6 @@ public CompletableFuture checkConnectionLiveness() { } } - @Override - public void healthCheckManually() { - long lastCheckTime = LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER.get(this); - if (System.currentTimeMillis() - lastCheckTime < 5000) { - return; - } - if (LAST_MANUAL_HEARTBEAT_CHECK_TIME_UPDATER.compareAndSet(this, lastCheckTime, System.currentTimeMillis())) { - sendPing(); - } - } - @Override protected void messageReceived() { super.messageReceived(); 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 dffcea3bb1eba..94f934fec681e 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 @@ -90,11 +90,4 @@ public interface TransportCnx { * is null if the connection liveness check is disabled. */ CompletableFuture checkConnectionLiveness(); - - /** - * If an IO exception is found, you can call this method immediately. The connection will be closed if the - * heartbeat check fails. - * Note: to avoid calling this method frequently, this method discards non-first calls within 5 seconds. - */ - void healthCheckManually(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index c45c2b1522f59..9a60bb2e71098 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -39,9 +39,14 @@ import com.google.common.collect.Sets; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; +import io.netty.channel.ChannelDuplexHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelOutboundHandler; +import io.netty.channel.ChannelPromise; import io.netty.util.Timeout; import java.io.ByteArrayInputStream; import java.io.IOException; +import java.net.SocketAddress; import java.nio.ByteBuffer; import java.nio.file.Files; import java.nio.file.Paths; @@ -84,7 +89,9 @@ import org.apache.commons.lang3.reflect.FieldUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.PulsarVersion; +import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.service.ServerCnx; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.schema.GenericRecord; @@ -4642,4 +4649,50 @@ public void testClientVersion() throws Exception { producer2.close(); client.close(); } + + @Test + public void testConsumeWhenDeliveryFailedByIOException() throws Exception { + final String topic = BrokerTestUtil.newUniqueName("persistent://my-property/my-ns/tp_"); + final String subscriptionName = "subscription1"; + final int messagesCount = 100; + final int receiverQueueSize = 1; + Producer producer = pulsarClient.newProducer(Schema.STRING).topic(topic).enableBatching(false).create(); + ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING).topic(topic) + .subscriptionName(subscriptionName).receiverQueueSize(receiverQueueSize).subscribe(); + for (int i = 0; i < messagesCount; i++) { + producer.send(i + ""); + } + // Wait incoming queue of the consumer is full. + Awaitility.await().untilAsserted(() -> { + assertEquals(consumer.getIncomingMessageSize(), receiverQueueSize); + }); + + // Mock an io error for sending messages out. + ServerCnx serverCnx = (ServerCnx) pulsar.getBrokerService().getTopic(topic, false).join().get() + .getSubscription(subscriptionName).getDispatcher().getConsumers().iterator().next().cnx(); + serverCnx.ctx().channel().pipeline().addFirst(new ChannelDuplexHandler() { + + @Override + public void flush(ChannelHandlerContext ctx) throws Exception { + throw new IOException("Mocked error"); + } + }); + + // Verify all messages will be consumed. + Set receivedMessages = new HashSet<>(); + while (true) { + Message msg = consumer.receive(2, TimeUnit.SECONDS); + if (msg != null) { + receivedMessages.add(msg.getValue()); + consumer.acknowledge(msg); + } else { + break; + } + } + Assert.assertEquals(receivedMessages.size(), messagesCount); + + producer.close(); + consumer.close(); + admin.topics().delete(topic, false); + } } From 1937acd53552a72f2d9b370901c1e882ddccea49 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Tue, 27 Jun 2023 19:35:35 +0800 Subject: [PATCH 12/12] remove unnecessary imports --- .../apache/pulsar/client/api/SimpleProducerConsumerTest.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 9a60bb2e71098..0c0e61fe33f88 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -41,12 +41,9 @@ import io.netty.buffer.Unpooled; import io.netty.channel.ChannelDuplexHandler; import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelOutboundHandler; -import io.netty.channel.ChannelPromise; import io.netty.util.Timeout; import java.io.ByteArrayInputStream; import java.io.IOException; -import java.net.SocketAddress; import java.nio.ByteBuffer; import java.nio.file.Files; import java.nio.file.Paths;