This is an automated email from the ASF dual-hosted git repository.
mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new c293ae3 Enable listener to receive messages even if receiver queue
size is zero (#1977)
c293ae3 is described below
commit c293ae33c28fad12c2ff3e90591e2a18ce95e97e
Author: massakam <[email protected]>
AuthorDate: Tue Jun 19 10:54:53 2018 +0900
Enable listener to receive messages even if receiver queue size is zero
(#1977)
* Enable listener to receive messages even if receiver queue size is zero
* Send flow if acked message is received
---
.gitignore | 1 +
.../pulsar/client/impl/ZeroQueueSizeTest.java | 48 ++++++++++++++++++++++
.../apache/pulsar/client/impl/ConsumerImpl.java | 39 ++++++++++++++----
3 files changed, 81 insertions(+), 7 deletions(-)
diff --git a/.gitignore b/.gitignore
index b009375..ed5a8b6 100644
--- a/.gitignore
+++ b/.gitignore
@@ -7,6 +7,7 @@ logs
/data
pulsar-broker/tmp.*
pulsar-broker/src/test/resources/log4j2.yaml
+pulsar-functions/worker/test-tenant/
*.log
*.versionsBackup
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ZeroQueueSizeTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ZeroQueueSizeTest.java
index 725fa46..2859047 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ZeroQueueSizeTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ZeroQueueSizeTest.java
@@ -20,8 +20,12 @@ package org.apache.pulsar.client.impl;
import static org.testng.Assert.assertEquals;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import com.google.common.collect.Lists;
+
import org.apache.pulsar.broker.service.BrokerTestBase;
import org.apache.pulsar.client.admin.PulsarAdminException;
import org.apache.pulsar.client.api.Consumer;
@@ -123,6 +127,50 @@ public class ZeroQueueSizeTest extends BrokerTestBase {
}
@Test()
+ public void zeroQueueSizeConsumerListener() throws Exception {
+ String key = "zeroQueueSizeConsumerListener";
+
+ // 1. Config
+ final String topicName = "persistent://prop/use/ns-abc/topic-" + key;
+ final String subscriptionName = "my-ex-subscription-" + key;
+ final String messagePredicate = "my-message-" + key + "-";
+
+ // 2. Create Producer
+ Producer<byte[]> producer = pulsarClient.newProducer().topic(topicName)
+ .enableBatching(false)
+ .messageRoutingMode(MessageRoutingMode.SinglePartition)
+ .create();
+
+ // 3. Create Consumer
+ List<Message<byte[]>> messages = Lists.newArrayList();
+ CountDownLatch latch = new CountDownLatch(totalMessages);
+ ConsumerImpl<byte[]> consumer = (ConsumerImpl<byte[]>)
pulsarClient.newConsumer().topic(topicName)
+
.subscriptionName(subscriptionName).receiverQueueSize(0).messageListener((cons,
msg) -> {
+ assertEquals(((ConsumerImpl) cons).numMessagesInQueue(),
0);
+ synchronized(messages) {
+ messages.add(msg);
+ }
+ log.info("Consumer received: " + new
String(msg.getData()));
+ latch.countDown();
+ }).subscribe();
+
+ // 3. producer publish messages
+ for (int i = 0; i < totalMessages; i++) {
+ String message = messagePredicate + i;
+ log.info("Producer produced: " + message);
+ producer.send(message.getBytes());
+ }
+
+ // 4. Receiver receives the message
+ latch.await();
+ assertEquals(consumer.numMessagesInQueue(), 0);
+ assertEquals(messages.size(), totalMessages);
+ for (int i = 0; i < messages.size(); i++) {
+ assertEquals(new String(messages.get(i).getData()),
messagePredicate + i);
+ }
+ }
+
+ @Test()
public void zeroQueueSizeSharedSubscription() throws PulsarClientException
{
String key = "zeroQueueSizeSharedSubscription";
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 056b17e..10fe6ff 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
@@ -485,7 +485,8 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
AVAILABLE_PERMITS_UPDATER.set(this, 0);
// For zerosize queue : If the connection is reset and
someone is waiting for the messages
// or queue was not empty: send a flow command
- if (waitingOnReceiveForZeroQueueSize ||
(conf.getReceiverQueueSize() == 0 && currentSize > 0)) {
+ if (waitingOnReceiveForZeroQueueSize
+ || (conf.getReceiverQueueSize() == 0 &&
(currentSize > 0 || listener != null))) {
sendFlowPermitsToBroker(cnx, 1);
}
} else {
@@ -678,6 +679,9 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
log.debug("[{}][{}] Ignoring message as it was already being
acked earlier by same consumer {}/{}",
topic, subscription, msgId);
}
+ if (conf.getReceiverQueueSize() == 0) {
+ increaseAvailablePermits(cnx);
+ }
return;
}
@@ -722,12 +726,12 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
// if the conf.getReceiverQueueSize() is 0 then discard
message if no one is waiting for it.
// if asyncReceive is waiting then notify callback without
adding to incomingMessages queue
unAckedMessageTracker.add((MessageIdImpl)
message.getMessageId());
- boolean asyncReceivedWaiting = !pendingReceives.isEmpty();
- if ((conf.getReceiverQueueSize() != 0 ||
waitingOnReceiveForZeroQueueSize) && !asyncReceivedWaiting) {
- incomingMessages.add(message);
- }
- if (asyncReceivedWaiting) {
+ if (!pendingReceives.isEmpty()) {
notifyPendingReceivedCallback(message, null);
+ } else if (conf.getReceiverQueueSize() != 0 ||
waitingOnReceiveForZeroQueueSize) {
+ incomingMessages.add(message);
+ } else if (conf.getReceiverQueueSize() == 0 && listener !=
null) {
+ triggerZeroQueueSizeListener(message);
}
} finally {
lock.readLock().unlock();
@@ -754,7 +758,7 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
msgMetadata.recycle();
}
- if (listener != null) {
+ if (listener != null && conf.getReceiverQueueSize() != 0) {
// Trigger the notification on the message listener in a separate
thread to avoid blocking the networking
// thread while the message processing happens
listenerExecutor.execute(() -> {
@@ -816,6 +820,27 @@ public class ConsumerImpl<T> extends ConsumerBase<T>
implements ConnectionHandle
}
}
+ private void triggerZeroQueueSizeListener(final Message<T> message) {
+ checkArgument(conf.getReceiverQueueSize() == 0);
+ checkNotNull(listener, "listener can't be null");
+ checkNotNull(message, "unqueued message can't be null");
+
+ listenerExecutor.execute(() -> {
+ stats.updateNumMsgsReceived(message);
+ try {
+ if (log.isDebugEnabled()) {
+ log.debug("[{}][{}] Calling message listener for unqueued
message {}", topic, subscription,
+ message.getMessageId());
+ }
+ listener.received(ConsumerImpl.this, message);
+ } catch (Throwable t) {
+ log.error("[{}][{}] Message listener error in processing
unqueued message: {}", topic, subscription,
+ message.getMessageId(), t);
+ }
+ increaseAvailablePermits(cnx());
+ });
+ }
+
void receiveIndividualMessagesFromBatch(MessageMetadata msgMetadata,
ByteBuf uncompressedPayload,
MessageIdData messageId, ClientCnx cnx) {
int batchSize = msgMetadata.getNumMessagesInBatch();