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..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 @@ -361,6 +361,12 @@ public Future sendMessages(final List entries, EntryBatch msgOutCounter.add(totalMessages); bytesOutCounter.add(totalBytes); chunkedMessageRate.recordMultipleEvents(totalChunkedMessages, 0); + } else { + 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 194f8593e4aba..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 @@ -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: {}", this.ctx().channel().toString(), + topicName, subscriptionName, consumerId); try { Metadata.validateMetadata(metadata, service.getPulsar().getConfiguration().getMaxConsumerMetadataSize()); @@ -1783,6 +1784,12 @@ protected void handleAck(CommandAck ack) { } return null; }); + } else { + 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()); + } } } 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..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 @@ -39,6 +39,8 @@ 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.util.Timeout; import java.io.ByteArrayInputStream; import java.io.IOException; @@ -84,7 +86,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 +4646,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); + } }