From d7ac95fe142c6d78b6e856ad421e48b6ae009ab0 Mon Sep 17 00:00:00 2001 From: lipenghui Date: Tue, 25 Jun 2019 16:48:01 +0800 Subject: [PATCH] Reduce unnecessary track message calls. --- .../client/api/ConsumerRedeliveryTest.java | 54 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 43 ++++++--------- .../client/impl/MultiTopicsConsumerImpl.java | 2 +- 3 files changed, 70 insertions(+), 29 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java index 63c974dec5ee6..a42d80cd86f58 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ConsumerRedeliveryTest.java @@ -18,7 +18,11 @@ */ package org.apache.pulsar.client.api; +import java.util.ArrayList; +import java.util.List; import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import org.apache.pulsar.client.impl.ConsumerImpl; @@ -29,6 +33,7 @@ import com.google.common.collect.Sets; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import static org.testng.Assert.assertEquals; @@ -123,4 +128,53 @@ public void testOrderedRedelivery() throws Exception { consumer2.close(); } + @Test + public void testUnAckMessageRedeliveryWithReceiveAsync() throws PulsarClientException, ExecutionException, InterruptedException { + String topic = "persistent://my-property/my-ns/async-unack-redelivery"; + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("s1") + .ackTimeout(3, TimeUnit.SECONDS) + .subscribe(); + + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .enableBatching(true) + .batchingMaxMessages(5) + .batchingMaxPublishDelay(1, TimeUnit.SECONDS) + .create(); + + final int messages = 10; + List>> futures = new ArrayList<>(10); + for (int i = 0; i < messages; i++) { + futures.add(consumer.receiveAsync()); + } + + for (int i = 0; i < messages; i++) { + producer.sendAsync("my-message-" + i); + } + + int messageReceived = 0; + for (CompletableFuture> future : futures) { + Message message = future.get(); + assertNotNull(message); + messageReceived++; + // Don't ack message, wait for ack timeout. + } + + assertEquals(10, messageReceived); + + for (int i = 0; i < messages; i++) { + Message message = consumer.receive(); + assertNotNull(message); + messageReceived++; + consumer.acknowledge(message); + } + + assertEquals(20, messageReceived); + + producer.close(); + consumer.close(); + } + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index a25a98a1d291d..8704fd307d406 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -312,10 +312,8 @@ protected Message internalReceive() throws PulsarClientException { Message message; try { message = incomingMessages.take(); - trackMessage(message); - Message interceptMsg = beforeConsume(message); - messageProcessed(interceptMsg); - return interceptMsg; + messageProcessed(message); + return beforeConsume(message); } catch (InterruptedException e) { stats.incrementNumReceiveFailed(); throw PulsarClientException.unwrap(e); @@ -341,10 +339,8 @@ protected CompletableFuture> internalReceiveAsync() { } if (message != null) { - trackMessage(message); - Message interceptMsg = beforeConsume(message); - messageProcessed(interceptMsg); - result.complete(interceptMsg); + messageProcessed(message); + result.complete(beforeConsume(message)); } return result; @@ -355,12 +351,11 @@ protected Message internalReceive(int timeout, TimeUnit unit) throws PulsarCl Message message; try { message = incomingMessages.poll(timeout, unit); - trackMessage(message); - Message interceptMsg = beforeConsume(message); - if (interceptMsg != null) { - messageProcessed(interceptMsg); + if (message == null) { + return null; } - return interceptMsg; + messageProcessed(message); + return beforeConsume(message); } catch (InterruptedException e) { State state = getState(); if (state != State.Closing && state != State.Closed) { @@ -821,7 +816,6 @@ void messageReceived(MessageIdData messageId, int redeliveryCount, ByteBuf heade possibleSendToDeadLetterTopicMessages.put((MessageIdImpl)message.getMessageId(), Collections.singletonList(message)); } if (!pendingReceives.isEmpty()) { - trackMessage(message); notifyPendingReceivedCallback(message, null); } else if (canEnqueueMessage(message)) { incomingMessages.add(message); @@ -1044,19 +1038,7 @@ protected synchronized void messageProcessed(Message msg) { increaseAvailablePermits(currentCnx); stats.updateNumMsgsReceived(msg); - if (conf.getAckTimeoutMillis() != 0) { - // reset timer for messages that are received by the client - MessageIdImpl id = (MessageIdImpl) msg.getMessageId(); - if (id instanceof BatchMessageIdImpl) { - id = new MessageIdImpl(id.getLedgerId(), id.getEntryId(), getPartitionIndex()); - } - if (partitionIndex != -1) { - // we should no longer track this message, TopicsConsumer will take care from now onwards - unAckedMessageTracker.remove(id); - } else { - unAckedMessageTracker.add(id); - } - } + trackMessage(msg); } protected void trackMessage(Message msg) { @@ -1068,7 +1050,12 @@ protected void trackMessage(Message msg) { // do not add each item in batch message into tracker id = new MessageIdImpl(id.getLedgerId(), id.getEntryId(), getPartitionIndex()); } - unAckedMessageTracker.add(id); + if (partitionIndex != -1) { + // we should no longer track this message, TopicsConsumer will take care from now onwards + unAckedMessageTracker.remove(id); + } else { + unAckedMessageTracker.add(id); + } } } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java index a5e513870650e..5d45ca07f29ca 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MultiTopicsConsumerImpl.java @@ -260,7 +260,6 @@ private void messageReceived(ConsumerImpl consumer, Message message) { try { TopicMessageImpl topicMessage = new TopicMessageImpl<>( consumer.getTopic(), consumer.getTopicNameWithoutPartition(), message); - unAckedMessageTracker.add(topicMessage.getMessageId()); if (log.isDebugEnabled()) { log.debug("[{}][{}] Received message from topics-consumer {}", @@ -270,6 +269,7 @@ private void messageReceived(ConsumerImpl consumer, Message message) { // if asyncReceive is waiting : return message to callback without adding to incomingMessages queue if (!pendingReceives.isEmpty()) { CompletableFuture> receivedFuture = pendingReceives.poll(); + unAckedMessageTracker.add(topicMessage.getMessageId()); listenerExecutor.execute(() -> receivedFuture.complete(topicMessage)); } else { // Enqueue the message so that it can be retrieved when application calls receive()