From 2d17cb8d51df6bba50971c42472a3681974c673a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Arnaud=20Cogolu=C3=A8gnes?= <514737+acogoluegnes@users.noreply.github.com> Date: Fri, 23 Jan 2026 16:30:02 +0100 Subject: [PATCH] Use ConcurrentSkipListMap for unconfirmed messages tracking This eliminates TreeMap snapshot copies in the confirm timeout task. The monotonically increasing publishing IDs mean insertions always occur at the tail while removals (confirms) happen near the head, which suits the skip list's concurrent access patterns well. --- .../rabbitmq/stream/impl/StreamProducer.java | 22 ++++++++----------- 1 file changed, 9 insertions(+), 13 deletions(-) diff --git a/src/main/java/com/rabbitmq/stream/impl/StreamProducer.java b/src/main/java/com/rabbitmq/stream/impl/StreamProducer.java index 57916bfdf3..33f6f215a2 100644 --- a/src/main/java/com/rabbitmq/stream/impl/StreamProducer.java +++ b/src/main/java/com/rabbitmq/stream/impl/StreamProducer.java @@ -1,4 +1,4 @@ -// Copyright (c) 2020-2025 Broadcom. All Rights Reserved. +// Copyright (c) 2020-2026 Broadcom. All Rights Reserved. // The term "Broadcom" refers to Broadcom Inc. and/or its subsidiaries. // // This software, the RabbitMQ Stream Java client library, is dual-licensed under the @@ -48,10 +48,8 @@ import java.util.Map; import java.util.Map.Entry; import java.util.Objects; -import java.util.SortedMap; -import java.util.TreeMap; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.ConcurrentNavigableMap; +import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; @@ -73,8 +71,7 @@ final class StreamProducer extends ResourceBase implements Producer { private static final ConfirmationHandler NO_OP_CONFIRMATION_HANDLER = confirmationStatus -> {}; private final long id; private final MessageAccumulator accumulator; - // FIXME investigate a more optimized data structure to handle pending messages - private final ConcurrentMap unconfirmedMessages; + private final ConcurrentNavigableMap unconfirmedMessages; private final int batchSize; private final String name; private final String stream; @@ -148,7 +145,7 @@ final class StreamProducer extends ResourceBase implements Producer { this.maxUnconfirmedMessages = maxUnconfirmedMessages; this.unconfirmedMessagesSemaphore = new Semaphore(maxUnconfirmedMessages, true); - this.unconfirmedMessages = new ConcurrentHashMap<>(this.maxUnconfirmedMessages, 0.75f, 2); + this.unconfirmedMessages = new ConcurrentSkipListMap<>(); if (filterValueExtractor == null) { this.publishVersion = VERSION_1; @@ -290,10 +287,8 @@ public int fragmentLength(Object entity) { private Runnable confirmTimeoutTask(Duration confirmTimeout) { return () -> { long limit = this.environment.clock().time() - confirmTimeout.toNanos(); - SortedMap unconfirmedSnapshot = - new TreeMap<>(this.unconfirmedMessages); int count = 0; - for (Entry unconfirmedEntry : unconfirmedSnapshot.entrySet()) { + for (Entry unconfirmedEntry : this.unconfirmedMessages.entrySet()) { if (unconfirmedEntry.getValue().time() < limit) { if (Thread.currentThread().isInterrupted()) { return; @@ -522,7 +517,7 @@ void running() { "Re-publishing {} unconfirmed message(s)", this.unconfirmedMessages.size()); if (!this.unconfirmedMessages.isEmpty()) { Map messagesToResend = - new TreeMap<>(this.unconfirmedMessages); + new ConcurrentSkipListMap<>(this.unconfirmedMessages); this.unconfirmedMessages.clear(); Iterator> resendIterator = messagesToResend.entrySet().iterator(); @@ -550,7 +545,8 @@ void running() { LOGGER.debug( "Skipping republishing of {} unconfirmed messages", this.unconfirmedMessages.size()); - Map messagesToFail = new TreeMap<>(this.unconfirmedMessages); + Map messagesToFail = + new ConcurrentSkipListMap<>(this.unconfirmedMessages); this.unconfirmedMessages.clear(); for (AccumulatedEntity accumulatedEntity : messagesToFail.values()) { try {