From 50829c2c29d09d3bd66a73c510091a209b331d46 Mon Sep 17 00:00:00 2001 From: Matteo Merli Date: Tue, 10 Apr 2018 21:38:11 -0700 Subject: [PATCH] Fixed race condition intruduced managed ledger addEntry introduced in #1521 --- .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index bf2cd8c587697..c60bd49827580 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -92,7 +92,6 @@ import org.apache.pulsar.common.api.proto.PulsarApi.CommandSubscribe.InitialPosition; import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata; import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; -import org.apache.pulsar.common.util.collections.GrowableArrayBlockingQueue; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -202,7 +201,7 @@ enum PositionBound { * Queue of pending entries to be added to the managed ledger. Typically entries are queued when a new ledger is * created asynchronously and hence there is no ready ledger to write into. */ - final GrowableArrayBlockingQueue pendingAddEntries = new GrowableArrayBlockingQueue<>(); + final ConcurrentLinkedQueue pendingAddEntries = new ConcurrentLinkedQueue<>(); // ////////////////////////////////////////////////////////////////////// @@ -488,10 +487,11 @@ public void asyncAddEntry(ByteBuf buffer, AddEntryCallback callback, Object ctx) } OpAddEntry addOperation = OpAddEntry.create(this, buffer, callback, ctx); - pendingAddEntries.add(addOperation); // Jump to specific thread to avoid contention from writers writing from different threads executor.executeOrdered(name, safeRun(() -> { + pendingAddEntries.add(addOperation); + internalAsyncAddEntry(addOperation); })); } @@ -1197,7 +1197,7 @@ public synchronized void updateLedgersIdsComplete(Stat stat) { } // Process all the pending addEntry requests - for (OpAddEntry op : pendingAddEntries.toList()) { + for (OpAddEntry op : pendingAddEntries) { op.setLedger(currentLedger); ++currentLedgerEntries; currentLedgerSize += op.data.readableBytes();