From fe35208292a55f36004d78df628e1f82f219f09c Mon Sep 17 00:00:00 2001 From: Ezequiel Lovelle Date: Mon, 26 Nov 2018 21:13:16 -0300 Subject: [PATCH 1/2] Deferred messages for consumers This feature offers the capability to configure consumer subscription with an arbitrary receiver delay. [pulsar-broker] - Add field for DelayQueue to store next positions pending of delivery. - Add field for ScheduledFuture to schedule the head of DelayQueue in order to start the process of delivering expired entries. - Fix on sendMessages method usage of received list of entries when such entries are being filtered by updatePermitsAndPendingAcks method. Received list of entries might be filtered by updatePermitsAndPendingAcks method and afterwards entries are processed by execute() with lambda which is probably that it's executing thread will be other than its calling thread. Therefore entries are now wrapped with CopyOnWriteArrayList in order to prevent ConcurrentModificationException by lambda on execute(), another approach could be to copy the entire list for lambda or to use a synchronized list, but this would result in a performance penalty even when entries are not being filtered, CopyOnWriteArrayList prevent this from happening. Another step further trying to fix this might be using Streams and applying transformations to inner list. This path was not taken because would require major changes. - Add processDelayEntries() method to process all elements added in DelayQueue which are ready to be delivered, at any given time just one task should be schedule using this method. - Add readPublishFrom() method to get the parameter of publish time from metadata of a message without changing its reference offset, this method will only be used if the consumer has enabled the receiver delay parameter on the subscription. - Add inner private class DelayPositionInfo to represent each position to be schedule in DelayQueue. - Add method to clean-up previous mentioned fields related to deferred messages. DelayQueue from java.util.concurrent is used in order to store each messages position next to be expired, the advantage of using this queue is that at any given time one and only one task is scheduled per consumer avoiding to schedule an unbounded number of tasks. [pulsar-client] - Add receiverDelay() method at subscription level to configure this parameter. [pulsar-common] - Set receiver delay whether it was configured by user on subscription. - Add optional receiver delay parameter to protobuf pulsar schema, code generated by generate_protobuf_docker.sh script. --- .../pulsar/broker/service/Consumer.java | 165 ++++++++++++++++-- .../pulsar/broker/service/ServerCnx.java | 5 +- .../apache/pulsar/broker/service/Topic.java | 2 +- .../nonpersistent/NonPersistentTopic.java | 4 +- .../service/persistent/PersistentTopic.java | 4 +- ...sistentDispatcherFailoverConsumerTest.java | 12 +- .../PersistentTopicConcurrentTest.java | 8 +- .../broker/service/PersistentTopicTest.java | 52 +++--- .../pulsar/client/api/ConsumerBuilder.java | 13 ++ .../client/impl/ConsumerBuilderImpl.java | 7 + .../pulsar/client/impl/ConsumerImpl.java | 6 +- .../impl/conf/ConsumerConfigurationData.java | 2 + .../apache/pulsar/common/api/Commands.java | 9 +- .../pulsar/common/api/proto/PulsarApi.java | 57 ++++++ pulsar-common/src/main/proto/PulsarApi.proto | 2 + 15 files changed, 291 insertions(+), 57 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 e8d6057c8a1e7..50fc1fc94ab32 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 @@ -29,18 +29,13 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelPromise; -import java.util.ArrayList; -import java.util.Collections; -import java.util.Iterator; -import java.util.List; -import java.util.Map; -import java.util.Objects; +import java.util.*; +import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.stream.Collectors; -import lombok.Data; - +import io.netty.util.concurrent.ScheduledFuture; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.PositionImpl; @@ -67,6 +62,8 @@ * A Consumer is a consumer currently connected and associated with a Subscription */ public class Consumer { + private static final int MAX_REDELIVERY_AT_ONCE = 100; + private static final int SCHED_CHECK_DEFERRED_MS = 1000; private final Subscription subscription; private final SubType subType; private final ServerCnx cnx; @@ -76,6 +73,7 @@ public class Consumer { private final int partitionIdx; private final InitialPosition subscriptionInitialPosition; + private final long receiverDelay; private final long consumerId; private final int priorityLevel; private final boolean readCompacted; @@ -108,6 +106,9 @@ public class Consumer { private final Map metadata; + private final DelayQueue delayedPositions; + private volatile ScheduledFuture schedProcessDelay; + public interface SendListener { void sendComplete(ChannelFuture future, SendMessageInfo sendMessageInfo); } @@ -115,7 +116,8 @@ public interface SendListener { public Consumer(Subscription subscription, SubType subType, String topicName, long consumerId, int priorityLevel, String consumerName, int maxUnackedMessages, ServerCnx cnx, String appId, - Map metadata, boolean readCompacted, InitialPosition subscriptionInitialPosition) throws BrokerServiceException { + Map metadata, boolean readCompacted, InitialPosition subscriptionInitialPosition, + long receiverDelay) throws BrokerServiceException { this.subscription = subscription; this.subType = subType; @@ -132,6 +134,7 @@ public Consumer(Subscription subscription, SubType subType, String topicName, lo this.msgRedeliver = new Rate(); this.appId = appId; this.authenticationData = cnx.authenticationData; + this.receiverDelay = receiverDelay; PERMITS_RECEIVED_WHILE_CONSUMER_BLOCKED_UPDATER.set(this, 0); MESSAGE_PERMITS_UPDATER.set(this, 0); UNACKED_MESSAGES_UPDATER.set(this, 0); @@ -151,6 +154,66 @@ public Consumer(Subscription subscription, SubType subType, String topicName, lo // We don't need to keep track of pending acks if the subscription is not shared this.pendingAcks = null; } + + if (isReceiverDelayEnabled()) { + this.delayedPositions = new DelayQueue<>(); + this.schedProcessDelay = ctx().channel().eventLoop().schedule( + this::processDelayEntries, SCHED_CHECK_DEFERRED_MS, TimeUnit.MILLISECONDS); + } else { + this.delayedPositions = null; + this.schedProcessDelay = null; + } + } + + private void processDelayEntries() { + if (delayedPositions.isEmpty()) { + log.debug("[{}-{}] no pending messages to deliver for deferred consumer", topicName, subscription); + this.schedProcessDelay = ctx().channel().eventLoop().schedule( + this::processDelayEntries, SCHED_CHECK_DEFERRED_MS, TimeUnit.MILLISECONDS); + return; + } + + List redeliverPositions = new ArrayList<>(); + DelayPositionInfo entry; + int totalRedeliveryMessages = 0; + boolean hasBeenScheduled = false; + for (int i = 0; (entry = delayedPositions.peek()) != null && i <= MAX_REDELIVERY_AT_ONCE; i++) { + final long delay = entry.getDelay(TimeUnit.MILLISECONDS); + + if (delay > 0) { + this.schedProcessDelay = ctx().channel().eventLoop().schedule(this::processDelayEntries, + Math.min(delay, SCHED_CHECK_DEFERRED_MS), TimeUnit.MILLISECONDS); + hasBeenScheduled = true; + break; + } + + totalRedeliveryMessages += entry.getBatchSize(); + redeliverPositions.add(entry.getPosition()); + delayedPositions.remove(); + } + + if (!redeliverPositions.isEmpty()) { + log.debug("[{}-{}] redelivering {} messages on delayed consumer and available permits are: {}", + topicName, subscription, totalRedeliveryMessages, getAvailablePermits()); + + addAndGetUnAckedMsgs(this, -totalRedeliveryMessages); + blockedConsumerOnUnackedMsgs = false; + subscription.redeliverUnacknowledgedMessages(this, redeliverPositions); + msgRedeliver.recordMultipleEvents(totalRedeliveryMessages, totalRedeliveryMessages); + + PERMITS_RECEIVED_WHILE_CONSUMER_BLOCKED_UPDATER.getAndAdd(this, -totalRedeliveryMessages); + MESSAGE_PERMITS_UPDATER.getAndAdd(this, totalRedeliveryMessages); + subscription.consumerFlow(this, totalRedeliveryMessages); + } + + if (!hasBeenScheduled) { + if (delayedPositions.isEmpty()) { + this.schedProcessDelay = ctx().channel().eventLoop().schedule( + this::processDelayEntries, SCHED_CHECK_DEFERRED_MS, TimeUnit.MILLISECONDS); + } else { + ctx().channel().eventLoop().execute(this::processDelayEntries); + } + } } public SubType subType() { @@ -201,10 +264,11 @@ public SendMessageInfo sendMessages(final List entries) { * * @return a SendMessageInfo object that contains the detail of what was sent to consumer */ - public SendMessageInfo sendMessages(final List entries, SendListener listener) { + public SendMessageInfo sendMessages(final List e, SendListener listener) { final ChannelHandlerContext ctx = cnx.ctx(); final SendMessageInfo sentMessages = new SendMessageInfo(); final ChannelPromise writePromise = listener != null ? ctx.newPromise() : ctx.voidPromise(); + final CopyOnWriteArrayList entries = new CopyOnWriteArrayList<>(e); if (listener != null) { writePromise.addListener(future -> listener.sendComplete(writePromise, sentMessages)); @@ -303,6 +367,7 @@ public static int getBatchSizeforEntry(ByteBuf metadataAndPayload, Subscription } void updatePermitsAndPendingAcks(final List entries, SendMessageInfo sentMessages) throws PulsarServerException { + final List entriesToDiscard = new ArrayList<>(); int permitsToReduce = 0; Iterator iter = entries.iterator(); boolean unsupportedVersion = false; @@ -314,12 +379,26 @@ void updatePermitsAndPendingAcks(final List entries, SendMessageInfo sent int batchSize = getBatchSizeforEntry(metadataAndPayload, subscription, consumerId); if (batchSize == -1) { // this would suggest that the message might have been corrupted - iter.remove(); + entriesToDiscard.add(entry); PositionImpl pos = (PositionImpl) entry.getPosition(); entry.release(); subscription.acknowledgeMessage(Collections.singletonList(pos), AckType.Individual, Collections.emptyMap()); continue; } + + if (isReceiverDelayEnabled()) { + final long timeToDeliver = readPublishFrom(metadataAndPayload) + receiverDelay; + + if (isExpired(timeToDeliver)) { + entriesToDiscard.add(entry); + PositionImpl pos = (PositionImpl) entry.getPosition(); + permitsToReduce += batchSize; + entry.release(); + delayedPositions.add(new DelayPositionInfo(pos, timeToDeliver, batchSize)); + continue; + } + } + if (pendingAcks != null) { pendingAcks.put(entry.getLedgerId(), entry.getEntryId(), batchSize, 0); } @@ -343,11 +422,32 @@ void updatePermitsAndPendingAcks(final List entries, SendMessageInfo sent } } + if (!entriesToDiscard.isEmpty()) { + entries.removeAll(entriesToDiscard); + } + msgOut.recordMultipleEvents(permitsToReduce, totalReadableBytes); sentMessages.totalSentMessages = permitsToReduce; sentMessages.totalSentMessageBytes = totalReadableBytes; } + private static boolean isExpired(long timeToDeliver) { + return System.currentTimeMillis() < timeToDeliver; + } + + private static long readPublishFrom(final ByteBuf metadataAndPayload) { + metadataAndPayload.markReaderIndex(); + final PulsarApi.MessageMetadata metadata = Commands.parseMessageMetadata(metadataAndPayload); + metadataAndPayload.resetReaderIndex(); + long publishTime = metadata.getPublishTime(); + metadata.recycle(); + return publishTime; + } + + private boolean isReceiverDelayEnabled() { + return receiverDelay > 0; + } + public boolean isWritable() { return cnx.isWritable(); } @@ -361,12 +461,23 @@ public void sendError(ByteBuf error) { * pending message acks */ public void close() throws BrokerServiceException { + tearDownDeferredData(); subscription.removeConsumer(this); cnx.removedConsumer(this); } + private void tearDownDeferredData() { + if (isReceiverDelayEnabled()) { + if (!schedProcessDelay.isDone()) { + schedProcessDelay.cancel(false); + } + delayedPositions.clear(); + } + } + public void disconnect() { log.info("Disconnecting consumer: {}", this); + tearDownDeferredData(); cnx.closeConsumer(this); try { close(); @@ -377,6 +488,7 @@ public void disconnect() { void doUnsubscribe(final long requestId) { final ChannelHandlerContext ctx = cnx.ctx(); + tearDownDeferredData(); subscription.doUnsubscribe(this).thenAccept(v -> { log.info("Unsubscribed successfully from {}", subscription); @@ -698,5 +810,36 @@ public void setTotalSentMessageBytes(long totalSentMessageBytes) { } } + private static final class DelayPositionInfo implements Delayed { + + private final PositionImpl position; + private final long timeToDeliver; + private final int batchSize; + + DelayPositionInfo(PositionImpl position, long timeToDeliver, int batchSize) { + this.position = position; + this.batchSize = batchSize; + this.timeToDeliver = timeToDeliver; + } + + PositionImpl getPosition() { + return position; + } + + int getBatchSize() { + return batchSize; + } + + @Override + public long getDelay(TimeUnit unit) { + return unit.toMillis(timeToDeliver - System.currentTimeMillis()); + } + + @Override + public int compareTo(Delayed d) { + return Long.compare(timeToDeliver, ((DelayPositionInfo) d).timeToDeliver); + } + } + private static final Logger log = LoggerFactory.getLogger(Consumer.class); } 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 e6819bf57913e..2395036ab8053 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 @@ -527,6 +527,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { final Map metadata = CommandUtils.metadataFromCommand(subscribe); final InitialPosition initialPosition = subscribe.getInitialPosition(); final SchemaData schema = subscribe.hasSchema() ? getSchema(subscribe.getSchema()) : null; + final long receiverDelay = subscribe.hasReceiverDelay() ? subscribe.getReceiverDelay() : 0; CompletableFuture isProxyAuthorizedFuture; if (service.isAuthorizationEnabled() && originalPrincipal != null) { @@ -596,7 +597,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { return topic.subscribe(ServerCnx.this, subscriptionName, consumerId, subType, priorityLevel, consumerName, isDurable, startMessageId, metadata, - readCompacted, initialPosition); + readCompacted, initialPosition, receiverDelay); } else { return FutureUtil.failedFuture( new IncompatibleSchemaException( @@ -607,7 +608,7 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { } else { return topic.subscribe(ServerCnx.this, subscriptionName, consumerId, subType, priorityLevel, consumerName, isDurable, - startMessageId, metadata, readCompacted, initialPosition); + startMessageId, metadata, readCompacted, initialPosition, receiverDelay); } }) .thenAccept(consumer -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java index 4209af79183b7..a4fbf24bc6bb2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Topic.java @@ -85,7 +85,7 @@ default long getOriginalSequenceId() { CompletableFuture subscribe(ServerCnx cnx, String subscriptionName, long consumerId, SubType subType, int priorityLevel, String consumerName, boolean isDurable, MessageId startMessageId, - Map metadata, boolean readCompacted, InitialPosition initialPosition); + Map metadata, boolean readCompacted, InitialPosition initialPosition, long receiverDelay); CompletableFuture createSubscription(String subscriptionName, InitialPosition initialPosition); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java index bd60d954a34c0..fd9c5dda51ac8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopic.java @@ -309,7 +309,7 @@ public void removeProducer(Producer producer) { @Override public CompletableFuture subscribe(final ServerCnx cnx, String subscriptionName, long consumerId, SubType subType, int priorityLevel, String consumerName, boolean isDurable, MessageId startMessageId, - Map metadata, boolean readCompacted, InitialPosition initialPosition) { + Map metadata, boolean readCompacted, InitialPosition initialPosition, long receiverDelay) { final CompletableFuture future = new CompletableFuture<>(); @@ -360,7 +360,7 @@ public CompletableFuture subscribe(final ServerCnx cnx, String subscri try { Consumer consumer = new Consumer(subscription, subType, topic, consumerId, priorityLevel, consumerName, 0, cnx, - cnx.getRole(), metadata, readCompacted, initialPosition); + cnx.getRole(), metadata, readCompacted, initialPosition, receiverDelay); subscription.addConsumer(consumer); if (!cnx.isActive()) { consumer.close(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 814bf8508e29f..0adeac08b4e6f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -482,7 +482,7 @@ public void removeProducer(Producer producer) { @Override public CompletableFuture subscribe(final ServerCnx cnx, String subscriptionName, long consumerId, SubType subType, int priorityLevel, String consumerName, boolean isDurable, MessageId startMessageId, - Map metadata, boolean readCompacted, InitialPosition initialPosition) { + Map metadata, boolean readCompacted, InitialPosition initialPosition, long receiverDelay) { final CompletableFuture future = new CompletableFuture<>(); @@ -557,7 +557,7 @@ public CompletableFuture subscribe(final ServerCnx cnx, String subscri subscriptionFuture.thenAccept(subscription -> { try { Consumer consumer = new Consumer(subscription, subType, topic, consumerId, priorityLevel, consumerName, - maxUnackedMessages, cnx, cnx.getRole(), metadata, readCompacted, initialPosition); + maxUnackedMessages, cnx, cnx.getRole(), metadata, readCompacted, initialPosition, receiverDelay); subscription.addConsumer(consumer); if (!cnx.isActive()) { consumer.close(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentDispatcherFailoverConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentDispatcherFailoverConsumerTest.java index 9db657a059774..feae026e4d4ae 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentDispatcherFailoverConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentDispatcherFailoverConsumerTest.java @@ -269,7 +269,7 @@ public void testConsumerGroupChangesWithOldNewConsumers() throws Exception { // 2. Add old consumer Consumer consumer1 = new Consumer(sub, SubType.Exclusive, topic.getName(), 1 /* consumer id */, 0, - "Cons1"/* consumer name */, 50000, serverCnxWithOldVersion, "myrole-1", Collections.emptyMap(), false, InitialPosition.Latest); + "Cons1"/* consumer name */, 50000, serverCnxWithOldVersion, "myrole-1", Collections.emptyMap(), false, InitialPosition.Latest, 0); pdfc.addConsumer(consumer1); List consumers = pdfc.getConsumers(); assertTrue(consumers.get(0).consumerName() == consumer1.consumerName()); @@ -280,7 +280,7 @@ public void testConsumerGroupChangesWithOldNewConsumers() throws Exception { // 3. Add new consumer Consumer consumer2 = new Consumer(sub, SubType.Exclusive, topic.getName(), 2 /* consumer id */, 0, - "Cons2"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), false, InitialPosition.Latest); + "Cons2"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), false, InitialPosition.Latest, 0); pdfc.addConsumer(consumer2); consumers = pdfc.getConsumers(); assertTrue(consumers.get(0).consumerName() == consumer1.consumerName()); @@ -309,7 +309,7 @@ public void testAddRemoveConsumer() throws Exception { // 2. Add consumer Consumer consumer1 = spy(new Consumer(sub, SubType.Exclusive, topic.getName(), 1 /* consumer id */, 0, "Cons1"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest)); + false /* read compacted */, InitialPosition.Latest, 0)); pdfc.addConsumer(consumer1); List consumers = pdfc.getConsumers(); assertTrue(consumers.get(0).consumerName() == consumer1.consumerName()); @@ -333,7 +333,7 @@ public void testAddRemoveConsumer() throws Exception { // 5. Add another consumer which does not change active consumer Consumer consumer2 = spy(new Consumer(sub, SubType.Exclusive, topic.getName(), 2 /* consumer id */, 0, "Cons2"/* consumer name */, - 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest)); + 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0)); pdfc.addConsumer(consumer2); consumers = pdfc.getConsumers(); assertTrue(pdfc.getActiveConsumer().consumerName() == consumer1.consumerName()); @@ -347,7 +347,7 @@ public void testAddRemoveConsumer() throws Exception { // 6. Add a consumer which changes active consumer Consumer consumer0 = spy(new Consumer(sub, SubType.Exclusive, topic.getName(), 0 /* consumer id */, 0, "Cons0"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest)); + false /* read compacted */, InitialPosition.Latest, 0)); pdfc.addConsumer(consumer0); consumers = pdfc.getConsumers(); assertTrue(pdfc.getActiveConsumer().consumerName() == consumer0.consumerName()); @@ -581,7 +581,7 @@ private Consumer getNextConsumer(PersistentDispatcherMultipleConsumers dispatche private Consumer createConsumer(int priority, int permit, boolean blocked, int id) throws Exception { Consumer consumer = new Consumer(null, SubType.Shared, "test-topic", id, priority, ""+id, 5000, - serverCnx, "appId", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest); + serverCnx, "appId", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0); try { consumer.flowPermits(permit); } catch (Exception e) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicConcurrentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicConcurrentTest.java index 7c5ca2077da5e..212833714a57d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicConcurrentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicConcurrentTest.java @@ -123,7 +123,7 @@ public void testConcurrentTopicAndSubscriptionDelete() throws Exception { .setSubType(PulsarApi.CommandSubscribe.SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); @@ -181,7 +181,7 @@ public void testConcurrentTopicGCAndSubscriptionDelete() throws Exception { .setSubType(PulsarApi.CommandSubscribe.SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); @@ -243,7 +243,7 @@ public void testConcurrentTopicDeleteAndUnsubscribe() throws Exception { .setSubType(PulsarApi.CommandSubscribe.SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); @@ -301,7 +301,7 @@ public void testConcurrentTopicDeleteAndSubsUnsubscribe() throws Exception { .setSubType(PulsarApi.CommandSubscribe.SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java index 9895e7bc52b79..32e8041bc835d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PersistentTopicTest.java @@ -428,7 +428,7 @@ public void testSubscribeFail() throws Exception { .setSubscription("").setRequestId(1).setSubType(SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); try { f1.get(); fail("should fail with exception"); @@ -447,12 +447,12 @@ public void testSubscribeUnsubscribe() throws Exception { // 1. simple subscribe Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); // 2. duplicate subscribe Future f2 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); try { f2.get(); @@ -476,7 +476,7 @@ public void testAddRemoveConsumer() throws Exception { // 1. simple add consumer Consumer consumer = new Consumer(sub, SubType.Exclusive, topic.getName(), 1 /* consumer id */, 0, "Cons1"/* consumer name */, - 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest); + 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer); assertTrue(sub.getDispatcher().isConsumerConnected()); @@ -517,14 +517,14 @@ public void testMaxConsumersShared() throws Exception { // 1. add consumer1 Consumer consumer = new Consumer(sub, SubType.Shared, topic.getName(), 1 /* consumer id */, 0, "Cons1"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer); assertEquals(sub.getConsumers().size(), 1); // 2. add consumer2 Consumer consumer2 = new Consumer(sub, SubType.Shared, topic.getName(), 2 /* consumer id */, 0, "Cons2"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer2); assertEquals(sub.getConsumers().size(), 2); @@ -532,7 +532,7 @@ public void testMaxConsumersShared() throws Exception { try { Consumer consumer3 = new Consumer(sub, SubType.Shared, topic.getName(), 3 /* consumer id */, 0, "Cons3"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer3); fail("should have failed"); } catch (BrokerServiceException e) { @@ -545,7 +545,7 @@ public void testMaxConsumersShared() throws Exception { // 4. add consumer4 to sub2 Consumer consumer4 = new Consumer(sub2, SubType.Shared, topic.getName(), 4 /* consumer id */, 0, "Cons4"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub2.addConsumer(consumer4); assertEquals(sub2.getConsumers().size(), 1); @@ -556,7 +556,7 @@ public void testMaxConsumersShared() throws Exception { try { Consumer consumer5 = new Consumer(sub2, SubType.Shared, topic.getName(), 5 /* consumer id */, 0, "Cons5"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub2.addConsumer(consumer5); fail("should have failed"); } catch (BrokerServiceException e) { @@ -608,14 +608,14 @@ public void testMaxConsumersFailover() throws Exception { // 1. add consumer1 Consumer consumer = new Consumer(sub, SubType.Failover, topic.getName(), 1 /* consumer id */, 0, "Cons1"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer); assertEquals(sub.getConsumers().size(), 1); // 2. add consumer2 Consumer consumer2 = new Consumer(sub, SubType.Failover, topic.getName(), 2 /* consumer id */, 0, "Cons2"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer2); assertEquals(sub.getConsumers().size(), 2); @@ -623,7 +623,7 @@ public void testMaxConsumersFailover() throws Exception { try { Consumer consumer3 = new Consumer(sub, SubType.Failover, topic.getName(), 3 /* consumer id */, 0, "Cons3"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer3); fail("should have failed"); } catch (BrokerServiceException e) { @@ -636,7 +636,7 @@ public void testMaxConsumersFailover() throws Exception { // 4. add consumer4 to sub2 Consumer consumer4 = new Consumer(sub2, SubType.Failover, topic.getName(), 4 /* consumer id */, 0, "Cons4"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub2.addConsumer(consumer4); assertEquals(sub2.getConsumers().size(), 1); @@ -647,7 +647,7 @@ public void testMaxConsumersFailover() throws Exception { try { Consumer consumer5 = new Consumer(sub2, SubType.Failover, topic.getName(), 5 /* consumer id */, 0, "Cons5"/* consumer name */, 50000, serverCnx, "myrole-1", Collections.emptyMap(), - false /* read compacted */, InitialPosition.Latest); + false /* read compacted */, InitialPosition.Latest, 0); sub2.addConsumer(consumer5); fail("should have failed"); } catch (BrokerServiceException e) { @@ -687,7 +687,7 @@ public void testUbsubscribeRaceConditions() throws Exception { PersistentTopic topic = new PersistentTopic(successTopicName, ledgerMock, brokerService); PersistentSubscription sub = new PersistentSubscription(topic, "sub-1", cursorMock); Consumer consumer1 = new Consumer(sub, SubType.Exclusive, topic.getName(), 1 /* consumer id */, 0, "Cons1"/* consumer name */, - 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest); + 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0); sub.addConsumer(consumer1); doAnswer(new Answer() { @@ -709,7 +709,7 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { try { Thread.sleep(10); /* delay to ensure that the ubsubscribe gets executed first */ new Consumer(sub, SubType.Exclusive, topic.getName(), 2 /* consumer id */, 0, "Cons2"/* consumer name */, - 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest); + 50000, serverCnx, "myrole-1", Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0); } catch (BrokerServiceException e) { assertTrue(e instanceof BrokerServiceException.SubscriptionFencedException); } @@ -739,7 +739,7 @@ public void testDeleteTopic() throws Exception { .setSubscription(successSubName).setRequestId(1).setSubType(SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), false /* read compacted */, InitialPosition.Latest, 0); f1.get(); assertTrue(topic.delete().isCompletedExceptionally()); @@ -754,7 +754,7 @@ public void testDeleteAndUnsubscribeTopic() throws Exception { .setSubscription(successSubName).setRequestId(1).setSubType(SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); @@ -808,7 +808,7 @@ public void testConcurrentTopicAndSubscriptionDelete() throws Exception { .setSubscription(successSubName).setRequestId(1).setSubType(SubType.Exclusive).build(); Future f1 = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); final CyclicBarrier barrier = new CyclicBarrier(2); @@ -895,7 +895,7 @@ public Object answer(InvocationOnMock invocationOnMock) throws Throwable { .setSubscription(successSubName).setRequestId(1).setSubType(SubType.Exclusive).build(); Future f = topic.subscribe(serverCnx, cmd.getSubscription(), cmd.getConsumerId(), cmd.getSubType(), - 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest); + 0, cmd.getConsumerName(), cmd.getDurable(), null, Collections.emptyMap(), cmd.getReadCompacted(), InitialPosition.Latest, 0); try { f.get(); @@ -1013,7 +1013,7 @@ public void testFailoverSubscription() throws Exception { // 1. Subscribe with non partition topic Future f1 = topic1.subscribe(serverCnx, cmd1.getSubscription(), cmd1.getConsumerId(), cmd1.getSubType(), 0, cmd1.getConsumerName(), cmd1.getDurable(), null, Collections.emptyMap(), - cmd1.getReadCompacted(), InitialPosition.Latest); + cmd1.getReadCompacted(), InitialPosition.Latest, 0); f1.get(); // 2. Subscribe with partition topic @@ -1025,7 +1025,7 @@ public void testFailoverSubscription() throws Exception { Future f2 = topic2.subscribe(serverCnx, cmd2.getSubscription(), cmd2.getConsumerId(), cmd2.getSubType(), 0, cmd2.getConsumerName(), cmd2.getDurable(), null, Collections.emptyMap(), - cmd2.getReadCompacted(), InitialPosition.Latest); + cmd2.getReadCompacted(), InitialPosition.Latest, 0); f2.get(); // 3. Subscribe and create second consumer @@ -1035,7 +1035,7 @@ public void testFailoverSubscription() throws Exception { Future f3 = topic2.subscribe(serverCnx, cmd3.getSubscription(), cmd3.getConsumerId(), cmd3.getSubType(), 0, cmd3.getConsumerName(), cmd3.getDurable(), null, Collections.emptyMap(), - cmd3.getReadCompacted(), InitialPosition.Latest); + cmd3.getReadCompacted(), InitialPosition.Latest, 0); f3.get(); assertEquals( @@ -1056,7 +1056,7 @@ public void testFailoverSubscription() throws Exception { Future f4 = topic2.subscribe(serverCnx, cmd4.getSubscription(), cmd4.getConsumerId(), cmd4.getSubType(), 0, cmd4.getConsumerName(), cmd4.getDurable(), null, Collections.emptyMap(), - cmd4.getReadCompacted(), InitialPosition.Latest); + cmd4.getReadCompacted(), InitialPosition.Latest, 0); f4.get(); assertEquals( @@ -1082,7 +1082,7 @@ public void testFailoverSubscription() throws Exception { Future f5 = topic2.subscribe(serverCnx, cmd5.getSubscription(), cmd5.getConsumerId(), cmd5.getSubType(), 0, cmd5.getConsumerName(), cmd5.getDurable(), null, Collections.emptyMap(), - cmd5.getReadCompacted(), InitialPosition.Latest); + cmd5.getReadCompacted(), InitialPosition.Latest, 0); try { f5.get(); @@ -1099,7 +1099,7 @@ public void testFailoverSubscription() throws Exception { Future f6 = topic2.subscribe(serverCnx, cmd6.getSubscription(), cmd6.getConsumerId(), cmd6.getSubType(), 0, cmd6.getConsumerName(), cmd6.getDurable(), null, Collections.emptyMap(), - cmd6.getReadCompacted(), InitialPosition.Latest); + cmd6.getReadCompacted(), InitialPosition.Latest, 0); f6.get(); // 7. unsubscribe exclusive sub diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java index f94064002ed0a..0cc4821701216 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerBuilder.java @@ -377,4 +377,17 @@ public interface ConsumerBuilder extends Cloneable { * whether to auto update partition increasement */ ConsumerBuilder autoUpdatePartitions(boolean autoUpdate); + + /** + * Set given receiver delay for consumer + * + * By default messages are received by consumer as soon as they are published by producer. + * With this parameter enabled the consumer will only receive messages that are older than receiver delay from the + * moment they were published. + * Messages which are not older are scheduled for delivery at the moment of time of being older than receiver delay. + * + * @param receiverDelay amount of time for messages to be received. + * @param timeUnit unit of time in which receiver delay is configured. + */ + ConsumerBuilder receiverDelay(long receiverDelay, TimeUnit timeUnit); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java index 7e809ee5ee787..2401d8d9a1c75 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerBuilderImpl.java @@ -295,6 +295,13 @@ public ConsumerBuilder deadLetterPolicy(DeadLetterPolicy deadLetterPolicy) { @Override public ConsumerBuilder autoUpdatePartitions(boolean autoUpdate) { conf.setAutoUpdatePartitions(autoUpdate); + } + + @Override + public ConsumerBuilder receiverDelay(long delay, TimeUnit unit) { + final long millis = unit.toMillis(delay); + checkArgument(millis > 0, "Receiver delay cannot be smaller than 1 millisecond."); + conf.setReceiverDelay(unit.toMillis(delay)); return this; } 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 5d5b1fa79467a..4f4d65dfcdb8a 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 @@ -140,6 +140,8 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle private final DeadLetterPolicy deadLetterPolicy; + private final long receiverDelay; + private Producer deadLetterProducer; protected volatile boolean paused; @@ -174,6 +176,7 @@ enum SubscriptionMode { this.priorityLevel = conf.getPriorityLevel(); this.readCompacted = conf.isReadCompacted(); this.subscriptionInitialPosition = conf.getSubscriptionInitialPosition(); + this.receiverDelay = conf.getReceiverDelay(); if (client.getConfiguration().getStatsIntervalSeconds() > 0) { stats = new ConsumerStatsRecorderImpl(client, conf, this); @@ -559,7 +562,8 @@ public void connectionOpened(final ClientCnx cnx) { si = null; } ByteBuf request = Commands.newSubscribe(topic, subscription, consumerId, requestId, getSubType(), priorityLevel, - consumerName, isDurable, startMessageIdData, metadata, readCompacted, InitialPosition.valueOf(subscriptionInitialPosition.getValue()), si); + consumerName, isDurable, startMessageIdData, metadata, readCompacted, InitialPosition.valueOf(subscriptionInitialPosition.getValue()), + si, receiverDelay); if (startMessageIdData != null) { startMessageIdData.recycle(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java index 0f61cbbf9190e..f6217cf4cad7e 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/conf/ConsumerConfigurationData.java @@ -92,6 +92,8 @@ public class ConsumerConfigurationData implements Serializable, Cloneable { private boolean autoUpdatePartitions = true; + private long receiverDelay = 0; + @JsonIgnore public String getSingleTopic() { checkArgument(topicNames.size() == 1); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java index 2adb274697289..ca49f602edf35 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/Commands.java @@ -322,12 +322,13 @@ public static ByteBufPair newSend(long producerId, long sequenceId, int numMessa public static ByteBuf newSubscribe(String topic, String subscription, long consumerId, long requestId, SubType subType, int priorityLevel, String consumerName) { return newSubscribe(topic, subscription, consumerId, requestId, subType, priorityLevel, consumerName, - true /* isDurable */, null /* startMessageId */, Collections.emptyMap(), false, InitialPosition.Earliest, null); + true /* isDurable */, null /* startMessageId */, Collections.emptyMap(), false, InitialPosition.Earliest, null, 0); } public static ByteBuf newSubscribe(String topic, String subscription, long consumerId, long requestId, SubType subType, int priorityLevel, String consumerName, boolean isDurable, MessageIdData startMessageId, - Map metadata, boolean readCompacted, InitialPosition subscriptionInitialPosition, SchemaInfo schemaInfo) { + Map metadata, boolean readCompacted, InitialPosition subscriptionInitialPosition, + SchemaInfo schemaInfo, long receiverDelay) { CommandSubscribe.Builder subscribeBuilder = CommandSubscribe.newBuilder(); subscribeBuilder.setTopic(topic); subscribeBuilder.setSubscription(subscription); @@ -350,6 +351,10 @@ public static ByteBuf newSubscribe(String topic, String subscription, long consu subscribeBuilder.setSchema(schema); } + if (receiverDelay > 0) { + subscribeBuilder.setReceiverDelay(receiverDelay); + } + CommandSubscribe subscribe = subscribeBuilder.build(); ByteBuf res = serializeWithSize(BaseCommand.newBuilder().setType(Type.SUBSCRIBE).setSubscribe(subscribe)); subscribeBuilder.recycle(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java index 66180fbb69344..b226dbc0b4afa 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/api/proto/PulsarApi.java @@ -6842,6 +6842,10 @@ public interface CommandSubscribeOrBuilder // optional .pulsar.proto.CommandSubscribe.InitialPosition initialPosition = 13 [default = Latest]; boolean hasInitialPosition(); org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition getInitialPosition(); + + // optional uint64 receiver_delay = 14; + boolean hasReceiverDelay(); + long getReceiverDelay(); } public static final class CommandSubscribe extends org.apache.pulsar.shaded.com.google.protobuf.v241.GeneratedMessageLite @@ -7170,6 +7174,16 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosi return initialPosition_; } + // optional uint64 receiver_delay = 14; + public static final int RECEIVER_DELAY_FIELD_NUMBER = 14; + private long receiverDelay_; + public boolean hasReceiverDelay() { + return ((bitField0_ & 0x00001000) == 0x00001000); + } + public long getReceiverDelay() { + return receiverDelay_; + } + private void initFields() { topic_ = ""; subscription_ = ""; @@ -7184,6 +7198,7 @@ private void initFields() { readCompacted_ = false; schema_ = org.apache.pulsar.common.api.proto.PulsarApi.Schema.getDefaultInstance(); initialPosition_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition.Latest; + receiverDelay_ = 0L; } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { @@ -7279,6 +7294,9 @@ public void writeTo(org.apache.pulsar.common.util.protobuf.ByteBufCodedOutputStr if (((bitField0_ & 0x00000800) == 0x00000800)) { output.writeEnum(13, initialPosition_.getNumber()); } + if (((bitField0_ & 0x00001000) == 0x00001000)) { + output.writeUInt64(14, receiverDelay_); + } } private int memoizedSerializedSize = -1; @@ -7339,6 +7357,10 @@ public int getSerializedSize() { size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream .computeEnumSize(13, initialPosition_.getNumber()); } + if (((bitField0_ & 0x00001000) == 0x00001000)) { + size += org.apache.pulsar.shaded.com.google.protobuf.v241.CodedOutputStream + .computeUInt64Size(14, receiverDelay_); + } memoizedSerializedSize = size; return size; } @@ -7478,6 +7500,8 @@ public Builder clear() { bitField0_ = (bitField0_ & ~0x00000800); initialPosition_ = org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition.Latest; bitField0_ = (bitField0_ & ~0x00001000); + receiverDelay_ = 0L; + bitField0_ = (bitField0_ & ~0x00002000); return this; } @@ -7564,6 +7588,10 @@ public org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe buildPartia to_bitField0_ |= 0x00000800; } result.initialPosition_ = initialPosition_; + if (((from_bitField0_ & 0x00002000) == 0x00002000)) { + to_bitField0_ |= 0x00001000; + } + result.receiverDelay_ = receiverDelay_; result.bitField0_ = to_bitField0_; return result; } @@ -7616,6 +7644,9 @@ public Builder mergeFrom(org.apache.pulsar.common.api.proto.PulsarApi.CommandSub if (other.hasInitialPosition()) { setInitialPosition(other.getInitialPosition()); } + if (other.hasReceiverDelay()) { + setReceiverDelay(other.getReceiverDelay()); + } return this; } @@ -7767,6 +7798,11 @@ public Builder mergeFrom( } break; } + case 112: { + bitField0_ |= 0x00002000; + receiverDelay_ = input.readUInt64(); + break; + } } } } @@ -8209,6 +8245,27 @@ public Builder clearInitialPosition() { return this; } + // optional uint64 receiver_delay = 14; + private long receiverDelay_ ; + public boolean hasReceiverDelay() { + return ((bitField0_ & 0x00002000) == 0x00002000); + } + public long getReceiverDelay() { + return receiverDelay_; + } + public Builder setReceiverDelay(long value) { + bitField0_ |= 0x00002000; + receiverDelay_ = value; + + return this; + } + public Builder clearReceiverDelay() { + bitField0_ = (bitField0_ & ~0x00002000); + receiverDelay_ = 0L; + + return this; + } + // @@protoc_insertion_point(builder_scope:pulsar.proto.CommandSubscribe) } diff --git a/pulsar-common/src/main/proto/PulsarApi.proto b/pulsar-common/src/main/proto/PulsarApi.proto index 761d8a1f526de..20f19aafc0a04 100644 --- a/pulsar-common/src/main/proto/PulsarApi.proto +++ b/pulsar-common/src/main/proto/PulsarApi.proto @@ -236,6 +236,8 @@ message CommandSubscribe { // Signal wthether the subscription will initialize on latest // or not -- earliest optional InitialPosition initialPosition = 13 [default = Latest]; + + optional uint64 receiver_delay = 14; } message CommandPartitionedTopicMetadata { From 4ef0ebeae60d27b80768f5ba1c799de29361a57c Mon Sep 17 00:00:00 2001 From: Ezequiel Lovelle Date: Thu, 3 Jan 2019 19:02:28 -0300 Subject: [PATCH 2/2] Tests for Delay messages on consumers - Add functional tests to verify correct behaviour from consumers with messages delay enabled. --- .../client/api/ProducerConsumerDelayTest.java | 98 +++++++++++++++++++ 1 file changed, 98 insertions(+) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerDelayTest.java diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerDelayTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerDelayTest.java new file mode 100644 index 0000000000000..84ea9b6efd5d1 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/ProducerConsumerDelayTest.java @@ -0,0 +1,98 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.client.api; + +import org.testng.Assert; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +import java.util.concurrent.TimeUnit; + +public class ProducerConsumerDelayTest extends ProducerConsumerBase { + + @BeforeMethod + @Override + protected void setup() throws Exception { + super.internalSetup(); + super.producerBaseSetup(); + } + + @AfterMethod + @Override + protected void cleanup() throws Exception { + super.internalCleanup(); + } + + @Test(timeOut = 30000) + public void testDelayConsumerOneSec() throws PulsarClientException { + final int numOfMsgs = 10; + final int delayInMs = 1000; + final String topicName = "persistent://my-property/my-ns/my-topic"; + + Producer producer = pulsarClient.newProducer().topic(topicName).create(); + Consumer consumer = pulsarClient.newConsumer().topic(topicName) + .receiverDelay(delayInMs, TimeUnit.MILLISECONDS) + .subscriptionName("subscriber-exclusive-with-delay").subscribe(); + + testConsumeWithDelay(numOfMsgs, delayInMs, producer, consumer); + } + + @Test(timeOut = 30000) + public void testDelayConsumerFiveSecs() throws PulsarClientException { + final int numOfMsg = 100; + final int delayInMs = 5000; + final String topicName = "persistent://my-property/my-ns/my-topic"; + + Producer producer = pulsarClient.newProducer().topic(topicName).create(); + Consumer consumerA = pulsarClient.newConsumer().topic(topicName) + .receiverDelay(delayInMs, TimeUnit.MILLISECONDS) + .subscriptionName("subscriber-shared-with-delay").subscribe(); + + testConsumeWithDelay(numOfMsg, delayInMs, producer, consumerA); + } + + private static void testConsumeWithDelay(int numOfMsg, int delayInMs, Producer producer, + Consumer consumer) + throws PulsarClientException { + // Produce messages normally + for (int i = 0; i < numOfMsg; i++) { + producer.sendAsync(String.format("Message num %d", i).getBytes()); + } + + int recvMsgs = consumeMsgWithDelay(numOfMsg, delayInMs, consumer); + consumer.unsubscribe(); + consumer.close(); + producer.close(); + Assert.assertEquals(numOfMsg, recvMsgs); + } + + private static int consumeMsgWithDelay(int numOfMsg, int delayInMs, Consumer consumer) + throws PulsarClientException { + int i; + for (i = 0; i < numOfMsg; i++) { + Message msg = consumer.receive(); + Assert.assertNotNull(msg); + long delay = System.currentTimeMillis() - msg.getPublishTime(); + Assert.assertTrue(delay >= delayInMs); + consumer.acknowledge(msg); + } + return i; + } +}