diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java index 4a8842b15c14e..7e4155c061375 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/InMemoryDelayedDeliveryTracker.java @@ -266,7 +266,7 @@ public void run(Timeout timeout) throws Exception { synchronized (dispatcher) { lastTickRun = clock.millis(); currentTimeoutTarget = -1; - timeout = null; + this.timeout = null; dispatcher.readMoreEntries(); } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 2fdb03cb1941a..a36b1eded3dc5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -277,7 +277,11 @@ public synchronized void readMoreEntries() { consumerList.size()); } havePendingRead = true; - minReplayedPosition = getMessagesToReplayNow(1).stream().findFirst().orElse(null); + Set toReplay = getMessagesToReplayNow(1); + minReplayedPosition = toReplay.stream().findFirst().orElse(null); + if (minReplayedPosition != null) { + redeliveryMessages.add(minReplayedPosition.getLedgerId(), minReplayedPosition.getEntryId()); + } cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, topic.getMaxReadPosition()); } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 73eb031d60c87..85b2f5163ce9b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -173,29 +173,36 @@ protected void sendMessagesToConsumers(ReadType readType, List entries) { // This may happen when consumer closed. See issue #12885 for details. if (!allowOutOfOrderDelivery) { Set messagesToReplayNow = this.getMessagesToReplayNow(1); - if (messagesToReplayNow != null && !messagesToReplayNow.isEmpty() && this.minReplayedPosition != null) { - PositionImpl relayPosition = messagesToReplayNow.stream().findFirst().get(); - // If relayPosition is a new entry wither smaller position is inserted for redelivery during this async - // read, it is possible that this relayPosition should dispatch to consumer first. So in order to - // preserver order delivery, we need to discard this read result, and try to trigger a replay read, - // that containing "relayPosition", by calling readMoreEntries. - if (relayPosition.compareTo(minReplayedPosition) < 0) { - if (log.isDebugEnabled()) { - log.debug("[{}] Position {} (<{}) is inserted for relay during current {} read, discard this " - + "read and retry with readMoreEntries.", - name, relayPosition, minReplayedPosition, readType); - } - if (readType == ReadType.Normal) { - entries.forEach(entry -> { - long stickyKeyHash = getStickyKeyHash(entry); - addMessageToReplay(entry.getLedgerId(), entry.getEntryId(), stickyKeyHash); - entry.release(); - }); - } else if (readType == ReadType.Replay) { - entries.forEach(Entry::release); + if (messagesToReplayNow != null && !messagesToReplayNow.isEmpty()) { + PositionImpl replayPosition = messagesToReplayNow.stream().findFirst().get(); + // We have received a message potentially from the delayed tracker and, since we're not using it + // right now, it needs to be added to the redelivery tracker or we won't attempt anymore to + // resend it (until we disconnect consumer). + redeliveryMessages.add(replayPosition.getLedgerId(), replayPosition.getEntryId()); + + if (this.minReplayedPosition != null) { + // If relayPosition is a new entry wither smaller position is inserted for redelivery during this + // async read, it is possible that this relayPosition should dispatch to consumer first. So in + // order to preserver order delivery, we need to discard this read result, and try to trigger a + // replay read, that containing "relayPosition", by calling readMoreEntries. + if (replayPosition.compareTo(minReplayedPosition) < 0) { + if (log.isDebugEnabled()) { + log.debug("[{}] Position {} (<{}) is inserted for relay during current {} read, " + + "discard this read and retry with readMoreEntries.", + name, replayPosition, minReplayedPosition, readType); + } + if (readType == ReadType.Normal) { + entries.forEach(entry -> { + long stickyKeyHash = getStickyKeyHash(entry); + addMessageToReplay(entry.getLedgerId(), entry.getEntryId(), stickyKeyHash); + entry.release(); + }); + } else if (readType == ReadType.Replay) { + entries.forEach(Entry::release); + } + readMoreEntries(); + return; } - readMoreEntries(); - return; } } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java index 041818850eda4..fc94f4d72c063 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/DelayedDeliveryTest.java @@ -27,6 +27,7 @@ import java.util.ArrayList; import java.util.HashSet; import java.util.List; +import java.util.Random; import java.util.Set; import java.util.TreeSet; import java.util.UUID; @@ -542,6 +543,48 @@ public void testDelayedDeliveryWithAllConsumersDisconnecting() throws Exception Awaitility.await().untilAsserted(() -> Assert.assertEquals(dispatcher.getNumberOfDelayedMessages(), 0)); } + @Test + public void testInterleavedMessagesOnKeySharedSubscription() throws Exception { + String topic = BrokerTestUtil.newUniqueName("testInterleavedMessagesOnKeySharedSubscription"); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("key-shared-sub") + .subscriptionType(SubscriptionType.Key_Shared) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .create(); + + Random random = new Random(0); + for (int i = 0; i < 10; i++) { + // Publish 1 message without delay and 1 with delay + producer.newMessage() + .value("immediate-msg-" + i) + .sendAsync(); + + int delayMillis = 1000 + random.nextInt(1000); + producer.newMessage() + .value("delayed-msg-" + i) + .deliverAfter(delayMillis, TimeUnit.MILLISECONDS) + .sendAsync(); + Thread.sleep(1000); + } + + producer.flush(); + + Set receivedMessages = new HashSet<>(); + + while (receivedMessages.size() < 20) { + Message msg = consumer.receive(3, TimeUnit.SECONDS); + receivedMessages.add(msg.getValue()); + consumer.acknowledge(msg); + } + } + @Test public void testDispatcherReadFailure() throws Exception { String topic = BrokerTestUtil.newUniqueName("testDispatcherReadFailure");