From ec907a272735215acf6c669b2548e4db6cd472a5 Mon Sep 17 00:00:00 2001 From: Nikolai Borisov Date: Wed, 1 Nov 2023 16:16:30 +0400 Subject: [PATCH] [improve][broker] Unblock stuck Key_Shared subscription after consumer reconnect --- conf/broker.conf | 8 + .../pulsar/broker/ServiceConfiguration.java | 10 ++ .../pulsar/broker/service/Consumer.java | 3 + .../pulsar/broker/service/Subscription.java | 20 +++ ...PersistentDispatcherMultipleConsumers.java | 3 + ...tStickyKeyDispatcherMultipleConsumers.java | 80 +++++++-- .../persistent/PersistentSubscription.java | 56 +++++++ .../client/api/KeySharedSubscriptionTest.java | 156 +++++++++++++++++- 8 files changed, 310 insertions(+), 26 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index a043f379ed478..9c7c106cf54df 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -336,6 +336,14 @@ maxUnackedMessagesPerSubscriptionOnBrokerBlocked=0.16 # Broker periodically checks if subscription is stuck and unblock if flag is enabled. (Default is disabled) unblockStuckSubscriptionEnabled=false +# If enabled subscriptions stores keys of messages which allows consumer not +# to stuck in case it goes to recently assigned. The setting allows to overcome +# situation when new KeyShared consumers will not get any messages until a consumer +# that did get messages disconnects or acks/nacks some messages +# The trade of is the need to track all the not acked messages in subscription +# which could potentially lead to performance degradation and higher memory consumption +rememberNotAckedMessagesKey=false + # Tick time to schedule task that checks topic publish rate limiting across all topics # Reducing to lower value can give more accuracy while throttling publish but # it uses more CPU to perform frequent check. (Disable publish throttling with value 0) diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 27b302ea8396b..0bb8d7cf1c1a5 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -889,6 +889,16 @@ The delayed message index time step(in seconds) in per bucket snapshot segment, + " unacked messages than this percentage limit and subscription will not receive any new messages " + " until that subscription acks back `limit/2` messages") private double maxUnackedMessagesPerSubscriptionOnBrokerBlocked = 0.16; + @FieldContext( + category = CATEGORY_POLICIES, + doc = "If enabled subscriptions stores keys of messages which allows consumer not " + + "to stuck in case it goes to recently assigned. The setting allows to overcome " + + "situation when new KeyShared consumers will not get any messages until a consumer " + + "that did get messages disconnects or acks/nacks some messages " + + "The trade of is the need to track all the not acked messages in subscription" + + "which could potentially lead to performance degradation and higher memory consumption" + ) + private boolean rememberNotAckedMessagesKey = false; @FieldContext( category = CATEGORY_POLICIES, doc = "Maximum size of Consumer metadata") 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 e72c805d73879..d7f925e8301fe 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 @@ -323,6 +323,7 @@ public Future sendMessages(final List entries, EntryBatch unackedMessages -= (batchSize - BitSet.valueOf(ackSet).cardinality()); } pendingAcks.put(entry.getLedgerId(), entry.getEntryId(), batchSize, stickyKeyHash); + subscription.addPendingMessageKey(entry, subscription.getName(), consumerId); if (log.isDebugEnabled()) { log.debug("[{}-{}] Added {}:{} ledger entry with batchSize of {} to pendingAcks in" + " broker.service.Consumer for consumerId: {}", @@ -984,6 +985,7 @@ private boolean removePendingAcks(PositionImpl position) { ? ackOwnedConsumer.getPendingAcks().get(position.getLedgerId(), position.getEntryId()) : null; if (ackedPosition != null) { + subscription.removePendingMessageKey(position.getEntryId()); if (!ackOwnedConsumer.getPendingAcks().remove(position.getLedgerId(), position.getEntryId())) { // Message was already removed by the other consumer return false; @@ -1102,6 +1104,7 @@ private int addAndGetUnAckedMsgs(Consumer consumer, int ackedMessages) { private void clearUnAckedMsgs() { int unaAckedMsgs = UNACKED_MESSAGES_UPDATER.getAndSet(this, 0); subscription.addUnAckedMessages(-unaAckedMsgs); + subscription.cleanPendingMessageKeys(); } public boolean isPreciseDispatcherFlowControl() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java index 9deeafdb272f5..d6004dc44e2db 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Subscription.java @@ -104,6 +104,26 @@ default long getNumberOfEntriesDelayed() { boolean isSubscriptionMigrated(); + default boolean isPendingAckMessageKeysRemembered() { + return false; + } + + default void addPendingMessageKey(Entry pendingEntry, String subscription, long consumerId) { + //Default is no op + } + + default void removePendingMessageKey(long pendingEntryId) { + //Default is no op + } + + default void cleanPendingMessageKeys() { + //Default is no op + } + + default boolean couldSendToConsumer(String messageKey, long consumerId) { + return true; + } + default void processReplicatedSubscriptionSnapshot(ReplicatedSubscriptionsSnapshot snapshot) { // Default is no-op } 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 b3d48252efe58..ba93f0eb208e6 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 @@ -200,6 +200,9 @@ protected boolean isConsumersExceededOnSubscription() { public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { // decrement unack-message count for removed consumer addUnAckedMessages(-consumer.getUnackedMessages()); + if (subscription.isPendingAckMessageKeysRemembered()) { + consumer.getPendingAcks().keys().forEach(value -> subscription.removePendingMessageKey(value.second)); + } if (consumerSet.removeAll(consumer) == 1) { consumerList.remove(consumer); log.info("Removed consumer {} with pending {} acks", consumer, consumer.getPendingAcks().size()); 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 8f05530f58bfa..f14ef9ec712f8 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 @@ -37,6 +37,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.ConsistentHashingStickyKeyConsumerSelector; @@ -52,6 +53,8 @@ import org.apache.pulsar.common.api.proto.CommandSubscribe.SubType; import org.apache.pulsar.common.api.proto.KeySharedMeta; import org.apache.pulsar.common.api.proto.KeySharedMode; +import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.util.FutureUtil; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -256,22 +259,28 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis availablePermits = Math.min(availablePermits, remainUnAckedMessages); } int maxMessagesForC = Math.min(entriesWithSameKeyCount, availablePermits); - int messagesForC = getRestrictedMaxEntriesForConsumer(consumer, entriesWithSameKey, maxMessagesForC, - readType, consumerStickyKeyHashesMap.get(consumer)); + Pair> messagesWithEntries = getRestrictedMaxEntriesForConsumer( + consumer, + entriesWithSameKey, + maxMessagesForC, + readType, consumerStickyKeyHashesMap.get(consumer) + ); + int messagesForC = messagesWithEntries.getKey(); + List toDispatch = messagesWithEntries.getValue(); if (log.isDebugEnabled()) { log.debug("[{}] select consumer {} with messages num {}, read type is {}", name, consumer.consumerName(), messagesForC, readType); } - if (messagesForC < entriesWithSameKeyCount) { + if (messagesForC < toDispatch.size()) { // We are not able to push all the messages with given key to its consumer, // so we discard for now and mark them for later redelivery - for (int i = messagesForC; i < entriesWithSameKeyCount; i++) { - Entry entry = entriesWithSameKey.get(i); + for (int i = messagesForC; i < toDispatch.size(); i++) { + Entry entry = toDispatch.get(i); long stickyKeyHash = getStickyKeyHash(entry); addMessageToReplay(entry.getLedgerId(), entry.getEntryId(), stickyKeyHash); entry.release(); - entriesWithSameKey.set(i, null); + toDispatch.set(i, null); } } @@ -279,7 +288,7 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis // remove positions first from replay list first : sendMessages recycles entries if (readType == ReadType.Replay) { for (int i = 0; i < messagesForC; i++) { - Entry entry = entriesWithSameKey.get(i); + Entry entry = toDispatch.get(i); redeliveryMessages.remove(entry.getLedgerId(), entry.getEntryId()); } } @@ -287,10 +296,10 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis SendMessageInfo sendMessageInfo = SendMessageInfo.getThreadLocal(); EntryBatchSizes batchSizes = EntryBatchSizes.get(messagesForC); EntryBatchIndexesAcks batchIndexesAcks = EntryBatchIndexesAcks.get(messagesForC); - totalEntries += filterEntriesForConsumer(entriesWithSameKey, batchSizes, sendMessageInfo, + totalEntries += filterEntriesForConsumer(toDispatch, batchSizes, sendMessageInfo, batchIndexesAcks, cursor, readType == ReadType.Replay, consumer); - consumer.sendMessages(entriesWithSameKey, batchSizes, batchIndexesAcks, + consumer.sendMessages(toDispatch, batchSizes, batchIndexesAcks, sendMessageInfo.getTotalMessages(), sendMessageInfo.getTotalBytes(), sendMessageInfo.getTotalChunkedMessages(), getRedeliveryTracker()).addListener(future -> { @@ -332,19 +341,24 @@ protected synchronized boolean trySendMessagesToConsumers(ReadType readType, Lis return false; } - private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List entries, int maxMessages, - ReadType readType, Set stickyKeyHashes) { + private Pair> getRestrictedMaxEntriesForConsumer( + Consumer consumer, + List entries, + int maxMessages, + ReadType readType, + Set stickyKeyHashes + ) { if (maxMessages == 0) { - return 0; + return Pair.of(0, entries); } if (readType == ReadType.Normal && stickyKeyHashes != null && redeliveryMessages.containsStickyKeyHashes(stickyKeyHashes)) { // If redeliveryMessages contains messages that correspond to the same hash as the messages // that the dispatcher is trying to send, do not send those messages for order guarantee - return 0; + return Pair.of(0, entries); } if (recentlyJoinedConsumers == null) { - return maxMessages; + return Pair.of(maxMessages, entries); } removeConsumersFromRecentJoinedConsumers(); PositionImpl maxReadPosition = recentlyJoinedConsumers.get(consumer); @@ -352,7 +366,12 @@ private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List en // is now ready to receive any message if (maxReadPosition == null) { // The consumer has not recently joined, so we can send all messages - return maxMessages; + return Pair.of(maxMessages, entries); + } + + if (subscription.isPendingAckMessageKeysRemembered()) { + //if pending ack messages tracked we do not need to block recently joined consumers by position + return getRestrictedEntriesForConsumerPendingAck(entries, consumer); } // If the read type is Replay, we should avoid send messages that hold by other consumer to the new consumers, @@ -381,11 +400,38 @@ private int getRestrictedMaxEntriesForConsumer(Consumer consumer, List en if (((PositionImpl) entries.get(i).getPosition()).compareTo(maxReadPosition) >= 0) { // We have already crossed the divider line. All messages in the list are now // newer than what we can currently dispatch to this consumer - return i; + return Pair.of(i, entries); + } + } + + return Pair.of(maxMessages, entries); + } + + private Pair> getRestrictedEntriesForConsumerPendingAck( + List entries, + Consumer consumer + ) { + List filtered = new ArrayList<>(entries.size()); + //if we have recently joined consumers we should skip sending messages with not acked keys + for (Entry entry : entries) { + MessageMetadata metadata = Commands.peekAndCopyMessageMetadata( + entry.getDataBuffer(), + subscription.getName(), + consumer.consumerId() + ); + + + if (metadata == null || !metadata.hasPartitionKey() + || subscription.couldSendToConsumer(metadata.getPartitionKey(), consumer.consumerId())) { + filtered.add(entry); + } else { + long stickyKeyHash = getStickyKeyHash(entry); + addMessageToReplay(entry.getLedgerId(), entry.getEntryId(), stickyKeyHash); + entry.release(); } } - return maxMessages; + return Pair.of(filtered.size(), filtered); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 88130c3c2010c..8d442a1f13451 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -30,6 +30,7 @@ import java.util.Optional; import java.util.TreeMap; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicLong; @@ -55,6 +56,7 @@ import org.apache.bookkeeper.mledger.impl.PositionImpl; import org.apache.commons.collections4.MapUtils; import org.apache.commons.lang3.tuple.MutablePair; +import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.intercept.BrokerInterceptor; import org.apache.pulsar.broker.service.AbstractSubscription; @@ -124,8 +126,11 @@ public class PersistentSubscription extends AbstractSubscription implements Subs Map.of(REPLICATED_SUBSCRIPTION_PROPERTY, 1L); private static final Map NON_REPLICATED_SUBSCRIPTION_CURSOR_PROPERTIES = Map.of(); + private final Map> pendingMessages = new ConcurrentHashMap<>(); + private volatile ReplicatedSubscriptionSnapshotCache replicatedSubscriptionSnapshotCache; private final PendingAckHandle pendingAckHandle; + private final boolean shouldRememberUnackedMessageKey; private volatile Map subscriptionProperties; private volatile CompletableFuture fenceFuture; @@ -159,6 +164,11 @@ public PersistentSubscription(PersistentTopic topic, String subscriptionName, Ma } else { this.pendingAckHandle = new PendingAckHandleDisabled(); } + shouldRememberUnackedMessageKey = topic + .getBrokerService() + .getPulsar() + .getConfig() + .isRememberNotAckedMessagesKey(); IS_FENCED_UPDATER.set(this, FALSE); } @@ -1269,6 +1279,52 @@ public boolean isSubscriptionMigrated() { return topic.isMigrated() && cursor.getNumberOfEntriesInBacklog(true) <= 0; } + @Override + public boolean isPendingAckMessageKeysRemembered() { + return shouldRememberUnackedMessageKey; + } + + @Override + public void addPendingMessageKey(Entry pendingMessage, String subscription, long consumerId) { + if (shouldRememberUnackedMessageKey && pendingMessage != null) { + MessageMetadata metadata = Commands.peekAndCopyMessageMetadata( + pendingMessage.getDataBuffer(), + subscription, + consumerId + ); + if (metadata != null && metadata.hasPartitionKey()) { + pendingMessages.put(pendingMessage.getEntryId(), Pair.of(consumerId, metadata.getPartitionKey())); + } + } + } + + @Override + public void removePendingMessageKey(long pendingEntryId) { + if (shouldRememberUnackedMessageKey) { + pendingMessages.remove(pendingEntryId); + } + } + + @Override + public void cleanPendingMessageKeys() { + if (shouldRememberUnackedMessageKey) { + pendingMessages.clear(); + } + } + + @Override + public boolean couldSendToConsumer(String messageKey, long consumerId) { + if (!shouldRememberUnackedMessageKey) { + return true; + } + for (Pair pendingMessageKey: pendingMessages.values()) { + if (messageKey.equals(pendingMessageKey.getValue()) && !pendingMessageKey.getKey().equals(consumerId)) { + return false; + } + } + return true; + } + @Override public Map getSubscriptionProperties() { return subscriptionProperties; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 18fb141be3178..ee18dd9c5895d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -95,11 +95,15 @@ public Object[][] partitionedProvider() { @DataProvider(name = "data") public Object[][] dataProvider() { return new Object[][] { - // Topic-Type and "Batching" - { "persistent", false }, - { "persistent", true }, - { "non-persistent", false }, - { "non-persistent", true }, + // Topic-Type, "Batching", "rememberUnackedMessages" + { "persistent", false, false }, + { "persistent", true, false }, + { "non-persistent", false, false }, + { "non-persistent", true, false }, + { "persistent", false, true }, + { "persistent", true, true }, + { "non-persistent", false, true }, + { "non-persistent", true, true }, }; } @@ -128,6 +132,7 @@ protected void cleanup() throws Exception { @AfterMethod(alwaysRun = true) public void resetDefaultNamespace() throws Exception { + this.conf.setRememberNotAckedMessagesKey(false); //reset to default List list = admin.namespaces().getTopics("public/default"); for (String topicName : list){ if (!pulsar.getBrokerService().isSystemTopic(topicName)) { @@ -140,8 +145,16 @@ public void resetDefaultNamespace() throws Exception { private static final int NUMBER_OF_KEYS = 300; @Test(dataProvider = "data") - public void testSendAndReceiveWithHashRangeAutoSplitStickyKeyConsumerSelector(String topicType, boolean enableBatch) + public void testSendAndReceiveWithHashRangeAutoSplitStickyKeyConsumerSelector( + String topicType, + boolean enableBatch, + boolean rememberUnackedKeys + ) throws PulsarClientException { + if (rememberUnackedKeys) { + this.conf.setRememberNotAckedMessagesKey(true); + } + String topic = topicType + "://public/default/key_shared-" + UUID.randomUUID(); @Cleanup @@ -167,7 +180,15 @@ public void testSendAndReceiveWithHashRangeAutoSplitStickyKeyConsumerSelector(St } @Test(dataProvider = "data") - public void testSendAndReceiveWithBatching(String topicType, boolean enableBatch) throws Exception { + public void testSendAndReceiveWithBatching( + String topicType, + boolean enableBatch, + boolean rememberUnackedKeys + ) throws Exception { + if (rememberUnackedKeys) { + this.conf.setRememberNotAckedMessagesKey(true); + } + String topic = topicType + "://public/default/key_shared-" + UUID.randomUUID(); @Cleanup @@ -264,8 +285,13 @@ public void testSendAndReceiveWithHashRangeExclusiveStickyKeyConsumerSelector(bo @Test(dataProvider = "data") public void testConsumerCrashSendAndReceiveWithHashRangeAutoSplitStickyKeyConsumerSelector( String topicType, - boolean enableBatch + boolean enableBatch, + boolean rememberUnackedKeys ) throws PulsarClientException, InterruptedException { + if (rememberUnackedKeys) { + this.conf.setRememberNotAckedMessagesKey(true); + } + String topic = topicType + "://public/default/key_shared_consumer_crash-" + UUID.randomUUID(); @Cleanup @@ -308,8 +334,13 @@ public void testConsumerCrashSendAndReceiveWithHashRangeAutoSplitStickyKeyConsum @Test(dataProvider = "data") public void testNonKeySendAndReceiveWithHashRangeAutoSplitStickyKeyConsumerSelector( String topicType, - boolean enableBatch + boolean enableBatch, + boolean rememberUnackedKeys ) throws PulsarClientException { + if (rememberUnackedKeys) { + this.conf.setRememberNotAckedMessagesKey(true); + } + String topic = topicType + "://public/default/key_shared_none_key-" + UUID.randomUUID(); @Cleanup @@ -1630,4 +1661,111 @@ public void testContinueDispatchMessagesWhenMessageDelayed() throws Exception { log.info("Got {} other messages...", sum); Assert.assertEquals(sum, delayedMessages + messages); } + + @Test + public void testSharedKeysOrderingWhenAddingConsumers() throws Exception { + this.conf.setRememberNotAckedMessagesKey(true); + + String topic = "testSharedKeysOrderingWhenAddingConsumers-" + UUID.randomUUID(); + + @Cleanup + Producer producer = createProducer(topic, false); + + @Cleanup + Consumer c1 = createConsumer(topic); + + for (int i = 0; i < 10; i++) { + producer.newMessage() + .key(String.valueOf(i)) + .value(i) + .send(); + } + + // All the already published messages will be pre-fetched by C1. + + // Adding a new consumer. + @Cleanup + Consumer c2 = createConsumer(topic); + + // Produce messages with the same keys as was pre-fetched by C1 + for (int i = 10; i < 20; i++) { + producer.newMessage() + .key(String.valueOf(i - 10)) + .value(i) + .send(); + } + + // Closing c1, would trigger all messages to go to c2 + c1.close(); + + for (int i = 0; i < 20; i++) { + Message msg = c2.receive(); + assertEquals(msg.getValue().intValue(), i); + + c2.acknowledge(msg); + } + } + + @Test + public void testConsumerIsNotBlockedForNewKeys() throws Exception { + this.conf.setRememberNotAckedMessagesKey(true); + + String topic = "testConsumerIsNotBlockedForNewKeys-" + UUID.randomUUID(); + String firstConsumerMessagesKey = "1"; + String secondConsumerMessagesKey = "2"; + + @Cleanup + Producer producer = createProducer(topic, false); + + @Cleanup + Consumer c1 = createConsumer(topic, KeySharedPolicy.stickyHashRange() + .ranges(Range.of(0, 45000))); + + for (int i = 0; i < 10; i++) { + producer.newMessage() + .key(firstConsumerMessagesKey) + .value(i) + .send(); + } + + // Receive first message to make the key unacked + Message firstMessage = c1.receive(); + assertEquals(0, firstMessage.getValue().intValue()); + + // Adding a new consumer. It will become recently joined one + @Cleanup + Consumer c2 = createConsumer(topic, KeySharedPolicy.stickyHashRange() + .ranges(Range.of(45001, 65535))); + + //Produce messages with the key which was not pre-fetched by C1 + for (int i = 10; i < 20; i++) { + producer.newMessage() + .key(secondConsumerMessagesKey) + .value(i) + .send(); + } + + //C2 is not blocked as the key "newMessagesKey" was not pre-fetched by C1, + //and we are safe here not to break the key ordering + for (int i = 10; i < 20; i++) { + Message msg = c2.receive(); + assertEquals(i, msg.getValue().intValue()); + + c2.acknowledge(msg); + } + + assertNull(c2.receive(1, TimeUnit.SECONDS)); + + c1.acknowledge(firstMessage); + + //verify C1 receives all the messages in the correct order + for (int i = 1; i < 10; i++) { + Message msg = c1.receive(); + assertEquals(i, msg.getValue().intValue()); + + c1.acknowledge(msg); + } + + assertNull(c1.receive(1, TimeUnit.SECONDS)); + } }