Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -361,6 +361,12 @@ public Future<Void> sendMessages(final List<? extends Entry> 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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down Expand Up @@ -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());
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<String> producer = pulsarClient.newProducer(Schema.STRING).topic(topic).enableBatching(false).create();
ConsumerImpl<String> consumer = (ConsumerImpl<String>) 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<String> receivedMessages = new HashSet<>();
while (true) {
Message<String> 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);
}
}