Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 9 additions & 13 deletions src/main/java/com/rabbitmq/stream/impl/StreamProducer.java
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -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;
Expand All @@ -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<Long, AccumulatedEntity> unconfirmedMessages;
private final ConcurrentNavigableMap<Long, AccumulatedEntity> unconfirmedMessages;
private final int batchSize;
private final String name;
private final String stream;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -290,10 +287,8 @@ public int fragmentLength(Object entity) {
private Runnable confirmTimeoutTask(Duration confirmTimeout) {
return () -> {
long limit = this.environment.clock().time() - confirmTimeout.toNanos();
SortedMap<Long, AccumulatedEntity> unconfirmedSnapshot =
new TreeMap<>(this.unconfirmedMessages);
int count = 0;
for (Entry<Long, AccumulatedEntity> unconfirmedEntry : unconfirmedSnapshot.entrySet()) {
for (Entry<Long, AccumulatedEntity> unconfirmedEntry : this.unconfirmedMessages.entrySet()) {
if (unconfirmedEntry.getValue().time() < limit) {
if (Thread.currentThread().isInterrupted()) {
return;
Expand Down Expand Up @@ -522,7 +517,7 @@ void running() {
"Re-publishing {} unconfirmed message(s)", this.unconfirmedMessages.size());
if (!this.unconfirmedMessages.isEmpty()) {
Map<Long, AccumulatedEntity> messagesToResend =
new TreeMap<>(this.unconfirmedMessages);
new ConcurrentSkipListMap<>(this.unconfirmedMessages);
this.unconfirmedMessages.clear();
Iterator<Entry<Long, AccumulatedEntity>> resendIterator =
messagesToResend.entrySet().iterator();
Expand Down Expand Up @@ -550,7 +545,8 @@ void running() {
LOGGER.debug(
"Skipping republishing of {} unconfirmed messages",
this.unconfirmedMessages.size());
Map<Long, AccumulatedEntity> messagesToFail = new TreeMap<>(this.unconfirmedMessages);
Map<Long, AccumulatedEntity> messagesToFail =
new ConcurrentSkipListMap<>(this.unconfirmedMessages);
this.unconfirmedMessages.clear();
for (AccumulatedEntity accumulatedEntity : messagesToFail.values()) {
try {
Expand Down
Loading