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-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; + } +} 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 {