From c8877f40a7bf241cb02987d5d8464d30c13aab26 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 15:41:48 +0800 Subject: [PATCH 1/6] Revert "[broker] Fixed delayed delivery after read operation error (#18098)" This reverts commit 68bfd13f5c3e7171146386558f3a52dd2d51c12c. --- ...PersistentDispatcherMultipleConsumers.java | 17 ++----- .../persistent/DelayedDeliveryTest.java | 45 ------------------- 2 files changed, 4 insertions(+), 58 deletions(-) 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..0b8e7e18338ca 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 @@ -93,7 +93,8 @@ public class PersistentDispatcherMultipleConsumers extends AbstractDispatcherMul "totalAvailablePermits"); protected volatile int totalAvailablePermits = 0; protected volatile int readBatchSize; - protected final Backoff readFailureBackoff; + protected final Backoff readFailureBackoff = new Backoff(15, TimeUnit.SECONDS, + 1, TimeUnit.MINUTES, 0, TimeUnit.MILLISECONDS); private static final AtomicIntegerFieldUpdater TOTAL_UNACKED_MESSAGES_UPDATER = AtomicIntegerFieldUpdater.newUpdater(PersistentDispatcherMultipleConsumers.class, @@ -128,10 +129,6 @@ public PersistentDispatcherMultipleConsumers(PersistentTopic topic, ManagedCurso : RedeliveryTrackerDisabled.REDELIVERY_TRACKER_DISABLED; this.readBatchSize = serviceConfig.getDispatcherMaxReadBatchSize(); this.initializeDispatchRateLimiterIfNeeded(Optional.empty()); - this.readFailureBackoff = new Backoff( - topic.getBrokerService().pulsar().getConfiguration().getDispatcherReadFailureBackoffInitialTimeInMs(), - TimeUnit.MILLISECONDS, - 1, TimeUnit.MINUTES, 0, TimeUnit.MILLISECONDS); } @Override @@ -662,10 +659,7 @@ public synchronized void readEntriesFailed(ManagedLedgerException exception, Obj topic.getBrokerService().executor().schedule(() -> { synchronized (PersistentDispatcherMultipleConsumers.this) { - // If it's a replay read we need to retry even if there's already - // another scheduled read, otherwise we'd be stuck until - // more messages are published. - if (!havePendingRead || readType == ReadType.Replay) { + if (!havePendingRead) { log.info("[{}] Retrying read operation", name); readMoreEntries(); } else { @@ -875,10 +869,7 @@ protected synchronized Set getMessagesToReplayNow(int maxMessagesT return redeliveryMessages.getMessagesToReplayNow(maxMessagesToRead); } else if (delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().hasMessageAvailable()) { delayedDeliveryTracker.get().resetTickTime(topic.getDelayedDeliveryTickTimeMillis()); - Set messagesAvailableNow = - delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); - messagesAvailableNow.forEach(p -> redeliveryMessages.add(p.getLedgerId(), p.getEntryId())); - return messagesAvailableNow; + return delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); } else { return Collections.emptySet(); } 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..480da2f5b94a8 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 @@ -34,7 +34,6 @@ import lombok.Cleanup; -import org.apache.bookkeeper.client.BKException; import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.service.Dispatcher; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -62,7 +61,6 @@ public void setup() throws Exception { conf.setSystemTopicEnabled(true); conf.setTopicLevelPoliciesEnabled(true); conf.setDelayedDeliveryTickTimeMillis(1024); - conf.setDispatcherReadFailureBackoffInitialTimeInMs(1000); super.internalSetup(); super.producerBaseSetup(); } @@ -541,47 +539,4 @@ public void testDelayedDeliveryWithAllConsumersDisconnecting() throws Exception Awaitility.await().untilAsserted(() -> Assert.assertEquals(dispatcher.getNumberOfDelayedMessages(), 0)); } - - @Test - public void testDispatcherReadFailure() throws Exception { - String topic = BrokerTestUtil.newUniqueName("testDispatcherReadFailure"); - - @Cleanup - Consumer consumer = pulsarClient.newConsumer(Schema.STRING) - .topic(topic) - .subscriptionName("shared-sub") - .subscriptionType(SubscriptionType.Shared) - .subscribe(); - - @Cleanup - Producer producer = pulsarClient.newProducer(Schema.STRING) - .topic(topic) - .create(); - - for (int i = 0; i < 10; i++) { - producer.newMessage() - .value("msg-" + i) - .deliverAfter(5, TimeUnit.SECONDS) - .sendAsync(); - } - - producer.flush(); - - Message msg = consumer.receive(100, TimeUnit.MILLISECONDS); - assertNull(msg); - - // Inject failure in BK read - this.mockBookKeeper.failNow(BKException.Code.ReadException); - - Set receivedMsgs = new TreeSet<>(); - for (int i = 0; i < 10; i++) { - msg = consumer.receive(10, TimeUnit.SECONDS); - receivedMsgs.add(msg.getValue()); - } - - assertEquals(receivedMsgs.size(), 10); - for (int i = 0; i < 10; i++) { - assertTrue(receivedMsgs.contains("msg-" + i)); - } - } } From 9aa94a680a90d54831fba3f9abb9c2e8317af70b Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 16:09:06 +0800 Subject: [PATCH 2/6] [branch-2.9] Fixed key-shared delivery of messages with interleaved delays --- .../InMemoryDelayedDeliveryTracker.java | 2 +- ...PersistentDispatcherMultipleConsumers.java | 6 ++- ...tStickyKeyDispatcherMultipleConsumers.java | 51 +++++++++++-------- .../persistent/DelayedDeliveryTest.java | 49 +++++++++++++++--- 4 files changed, 78 insertions(+), 30 deletions(-) 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 0b8e7e18338ca..1c57632ab01d6 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 @@ -274,7 +274,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 480da2f5b94a8..1406c56cec576 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 @@ -24,12 +24,7 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -import java.util.ArrayList; -import java.util.HashSet; -import java.util.List; -import java.util.Set; -import java.util.TreeSet; -import java.util.UUID; +import java.util.*; import java.util.concurrent.TimeUnit; import lombok.Cleanup; @@ -539,4 +534,46 @@ 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); + } + } } From 06a4e727c53aed6ca495b39ea6c27ab685ed7a62 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 16:11:12 +0800 Subject: [PATCH 3/6] Revert "Revert "[fix][sec] File tiered storage: upgrade jettison to get rid of CVE-2022-40149 (#18022)"" This reverts commit 36a9d1324764b50dbbc14aaf1be27a732d9f548f. --- pom.xml | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index adc6c1b7415c8..ad9d464792fe9 100644 --- a/pom.xml +++ b/pom.xml @@ -223,7 +223,7 @@ flexible messaging model and an intuitive client API. 2.3.1 1.5.0 3.1 - 4.0.3 + 1.5.1 0.6.1 @@ -798,6 +798,13 @@ flexible messaging model and an intuitive client API. import + + org.codehaus.jettison + jettison + ${jettison.version} + + + org.hdrhistogram HdrHistogram From 3ab5ace6e33ca8378df814d03baff2bf3bdff160 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 16:11:59 +0800 Subject: [PATCH 4/6] Revert "Revert "Revert "[fix][sec] File tiered storage: upgrade jettison to get rid of CVE-2022-40149 (#18022)""" This reverts commit 06a4e727c53aed6ca495b39ea6c27ab685ed7a62. --- pom.xml | 9 +-------- 1 file changed, 1 insertion(+), 8 deletions(-) diff --git a/pom.xml b/pom.xml index ad9d464792fe9..adc6c1b7415c8 100644 --- a/pom.xml +++ b/pom.xml @@ -223,7 +223,7 @@ flexible messaging model and an intuitive client API. 2.3.1 1.5.0 3.1 - 1.5.1 + 4.0.3 0.6.1 @@ -798,13 +798,6 @@ flexible messaging model and an intuitive client API. import - - org.codehaus.jettison - jettison - ${jettison.version} - - - org.hdrhistogram HdrHistogram From 9659bd40c8075c89e2bd647bc2c68836f179d2da Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 16:12:58 +0800 Subject: [PATCH 5/6] Revert "Revert "[broker] Fixed delayed delivery after read operation error (#18098)"" This reverts commit c8877f40 --- ...PersistentDispatcherMultipleConsumers.java | 17 +++++-- .../persistent/DelayedDeliveryTest.java | 45 +++++++++++++++++++ 2 files changed, 58 insertions(+), 4 deletions(-) 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 1c57632ab01d6..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 @@ -93,8 +93,7 @@ public class PersistentDispatcherMultipleConsumers extends AbstractDispatcherMul "totalAvailablePermits"); protected volatile int totalAvailablePermits = 0; protected volatile int readBatchSize; - protected final Backoff readFailureBackoff = new Backoff(15, TimeUnit.SECONDS, - 1, TimeUnit.MINUTES, 0, TimeUnit.MILLISECONDS); + protected final Backoff readFailureBackoff; private static final AtomicIntegerFieldUpdater TOTAL_UNACKED_MESSAGES_UPDATER = AtomicIntegerFieldUpdater.newUpdater(PersistentDispatcherMultipleConsumers.class, @@ -129,6 +128,10 @@ public PersistentDispatcherMultipleConsumers(PersistentTopic topic, ManagedCurso : RedeliveryTrackerDisabled.REDELIVERY_TRACKER_DISABLED; this.readBatchSize = serviceConfig.getDispatcherMaxReadBatchSize(); this.initializeDispatchRateLimiterIfNeeded(Optional.empty()); + this.readFailureBackoff = new Backoff( + topic.getBrokerService().pulsar().getConfiguration().getDispatcherReadFailureBackoffInitialTimeInMs(), + TimeUnit.MILLISECONDS, + 1, TimeUnit.MINUTES, 0, TimeUnit.MILLISECONDS); } @Override @@ -663,7 +666,10 @@ public synchronized void readEntriesFailed(ManagedLedgerException exception, Obj topic.getBrokerService().executor().schedule(() -> { synchronized (PersistentDispatcherMultipleConsumers.this) { - if (!havePendingRead) { + // If it's a replay read we need to retry even if there's already + // another scheduled read, otherwise we'd be stuck until + // more messages are published. + if (!havePendingRead || readType == ReadType.Replay) { log.info("[{}] Retrying read operation", name); readMoreEntries(); } else { @@ -873,7 +879,10 @@ protected synchronized Set getMessagesToReplayNow(int maxMessagesT return redeliveryMessages.getMessagesToReplayNow(maxMessagesToRead); } else if (delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().hasMessageAvailable()) { delayedDeliveryTracker.get().resetTickTime(topic.getDelayedDeliveryTickTimeMillis()); - return delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); + Set messagesAvailableNow = + delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); + messagesAvailableNow.forEach(p -> redeliveryMessages.add(p.getLedgerId(), p.getEntryId())); + return messagesAvailableNow; } else { return Collections.emptySet(); } 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 1406c56cec576..22d7307db12fd 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 @@ -29,6 +29,7 @@ import lombok.Cleanup; +import org.apache.bookkeeper.client.BKException; import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.service.Dispatcher; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -56,6 +57,7 @@ public void setup() throws Exception { conf.setSystemTopicEnabled(true); conf.setTopicLevelPoliciesEnabled(true); conf.setDelayedDeliveryTickTimeMillis(1024); + conf.setDispatcherReadFailureBackoffInitialTimeInMs(1000); super.internalSetup(); super.producerBaseSetup(); } @@ -576,4 +578,47 @@ public void testInterleavedMessagesOnKeySharedSubscription() throws Exception { consumer.acknowledge(msg); } } + + @Test + public void testDispatcherReadFailure() throws Exception { + String topic = BrokerTestUtil.newUniqueName("testDispatcherReadFailure"); + + @Cleanup + Consumer consumer = pulsarClient.newConsumer(Schema.STRING) + .topic(topic) + .subscriptionName("shared-sub") + .subscriptionType(SubscriptionType.Shared) + .subscribe(); + + @Cleanup + Producer producer = pulsarClient.newProducer(Schema.STRING) + .topic(topic) + .create(); + + for (int i = 0; i < 10; i++) { + producer.newMessage() + .value("msg-" + i) + .deliverAfter(5, TimeUnit.SECONDS) + .sendAsync(); + } + + producer.flush(); + + Message msg = consumer.receive(100, TimeUnit.MILLISECONDS); + assertNull(msg); + + // Inject failure in BK read + this.mockBookKeeper.failNow(BKException.Code.ReadException); + + Set receivedMsgs = new TreeSet<>(); + for (int i = 0; i < 10; i++) { + msg = consumer.receive(10, TimeUnit.SECONDS); + receivedMsgs.add(msg.getValue()); + } + + assertEquals(receivedMsgs.size(), 10); + for (int i = 0; i < 10; i++) { + assertTrue(receivedMsgs.contains("msg-" + i)); + } + } } From 42566576cab289ef3446410ca101e5acd00cb27f Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 19 Oct 2022 16:13:43 +0800 Subject: [PATCH 6/6] fix checkstyle --- .../broker/service/persistent/DelayedDeliveryTest.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) 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 22d7307db12fd..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 @@ -24,7 +24,13 @@ import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; -import java.util.*; +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; import java.util.concurrent.TimeUnit; import lombok.Cleanup;