From cd605fef193800bc4e6dd26f9f74c21f6f19ed3c Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 6 Feb 2025 16:10:08 +0800 Subject: [PATCH 1/5] [fix][ml] Don't switch thread in asyncAddEntry --- .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 8 +++----- 1 file changed, 3 insertions(+), 5 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 4f45fc67b6377..455d9320ec8d5 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 @@ -802,11 +802,9 @@ public void asyncAddEntry(ByteBuf buffer, int numberOfMessages, AddEntryCallback buffer.retain(); // Jump to specific thread to avoid contention from writers writing from different threads - executor.execute(() -> { - OpAddEntry addOperation = OpAddEntry.createNoRetainBuffer(this, buffer, numberOfMessages, callback, ctx, - currentLedgerTimeoutTriggered); - internalAsyncAddEntry(addOperation); - }); + final var addOperation = OpAddEntry.createNoRetainBuffer(this, buffer, numberOfMessages, callback, ctx, + currentLedgerTimeoutTriggered); + internalAsyncAddEntry(addOperation); } protected synchronized void internalAsyncAddEntry(OpAddEntry addOperation) { From bd71aa40c2af5c876779f316633da8d698a90a19 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 6 Feb 2025 19:36:57 +0800 Subject: [PATCH 2/5] Refactor asyncAddEntry --- .../mledger/impl/ManagedLedgerImpl.java | 46 +++++++++++-------- .../mledger/impl/ShadowManagedLedgerImpl.java | 15 +++--- 2 files changed, 35 insertions(+), 26 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 455d9320ec8d5..f9a0ff2620814 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 @@ -804,29 +804,39 @@ public void asyncAddEntry(ByteBuf buffer, int numberOfMessages, AddEntryCallback // Jump to specific thread to avoid contention from writers writing from different threads final var addOperation = OpAddEntry.createNoRetainBuffer(this, buffer, numberOfMessages, callback, ctx, currentLedgerTimeoutTriggered); - internalAsyncAddEntry(addOperation); + var added = false; + try { + // Use synchronized to ensure if `addOperation` is added to queue and fails later, it will be the first + // element in `pendingAddEntries`. + synchronized (this) { + if (managedLedgerInterceptor != null) { + managedLedgerInterceptor.beforeAddEntry(addOperation, addOperation.getNumberOfMessages()); + } + final var state = STATE_UPDATER.get(this); + beforeAddEntryToQueue(state); + pendingAddEntries.add(addOperation); + added = true; + afterAddEntryToQueue(state, addOperation); + } + } catch (Throwable throwable) { + if (!added) { + addOperation.failed(ManagedLedgerException.getManagedLedgerException(throwable)); + } // else: all elements of `pendingAddEntries` will fail in another thread + } } - protected synchronized void internalAsyncAddEntry(OpAddEntry addOperation) { - if (!beforeAddEntry(addOperation)) { - return; - } - final State state = STATE_UPDATER.get(this); + protected void beforeAddEntryToQueue(State state) throws ManagedLedgerException { if (state.isFenced()) { - addOperation.failed(new ManagedLedgerFencedException()); - return; - } else if (state == State.Terminated) { - addOperation.failed(new ManagedLedgerTerminatedException("Managed ledger was already terminated")); - return; - } else if (state == State.Closed) { - addOperation.failed(new ManagedLedgerAlreadyClosedException("Managed ledger was already closed")); - return; - } else if (state == State.WriteFailed) { - addOperation.failed(new ManagedLedgerAlreadyClosedException("Waiting to recover from failure")); - return; + throw new ManagedLedgerFencedException(); } - pendingAddEntries.add(addOperation); + switch (state) { + case Terminated -> throw new ManagedLedgerTerminatedException("Managed ledger was already terminated"); + case Closed -> throw new ManagedLedgerAlreadyClosedException("Managed ledger was already closed"); + case WriteFailed -> throw new ManagedLedgerAlreadyClosedException("Waiting to recover from failure"); + } + } + protected void afterAddEntryToQueue(State state, OpAddEntry addOperation) throws ManagedLedgerException { if (state == State.ClosingLedger || state == State.CreatingLedger) { // We don't have a ready ledger to write into // We are waiting for a new ledger to be created diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java index 4b03cad8e0a1d..84b4bab343507 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java @@ -223,18 +223,17 @@ private void initLastConfirmedEntry() { } @Override - protected synchronized void internalAsyncAddEntry(OpAddEntry addOperation) { - if (!beforeAddEntry(addOperation)) { - return; - } + protected void beforeAddEntryToQueue(State state) throws ManagedLedgerException { if (state != State.LedgerOpened) { - addOperation.failed(new ManagedLedgerException("Managed ledger is not opened")); - return; + throw new ManagedLedgerException("Managed ledger is not opened"); } + } + @Override + protected void afterAddEntryToQueue(State state, OpAddEntry addOperation) throws ManagedLedgerException { if (addOperation.getCtx() == null || !(addOperation.getCtx() instanceof Position position)) { - addOperation.failed(new ManagedLedgerException("Illegal addOperation context object.")); - return; + pendingAddEntries.poll(); + throw new ManagedLedgerException("Illegal addOperation context object."); } if (log.isDebugEnabled()) { From 0009959e85eeab7a861c1499e87a3473b52a9fd2 Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 6 Feb 2025 20:24:07 +0800 Subject: [PATCH 3/5] Add tests for beforeAddEntry --- .../ManagedLedgerInterceptorImplTest.java | 50 +++++++++++++++++++ 1 file changed, 50 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java index b57b5ce94be42..35869b912ee84 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java @@ -21,6 +21,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; @@ -29,9 +30,12 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; +import java.util.Random; import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Predicate; import lombok.Cleanup; @@ -499,4 +503,50 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { ledger.close(); } + @Test + public void testBeforeAddEntry() throws Exception { + final var interceptor = new ManagedLedgerInterceptorImpl(getBrokerEntryMetadataInterceptors(), null); + final var config = new ManagedLedgerConfig(); + final var numEntries = 10; + config.setMaxEntriesPerLedger(numEntries); + config.setManagedLedgerInterceptor(interceptor); + @Cleanup final var ml = (ManagedLedgerImpl) factory.open("test_concurrent_add_entry", config); + + final var indexesBeforeAdd = Collections.synchronizedList(new ArrayList()); + final var batchSizes = Collections.synchronizedList(new ArrayList()); + final var random = new Random(); + final var latch = new CountDownLatch(numEntries); + final var executor = Executors.newFixedThreadPool(2); + for (int i = 0; i < numEntries; i++) { + final var batchSize = random.nextInt(0, 100); + final var msg = "msg-" + i; + final var callback = new AsyncCallbacks.AddEntryCallback() { + + @Override + public void addComplete(Position position, ByteBuf entryData, Object ctx) { + latch.countDown(); + } + + @Override + public void addFailed(ManagedLedgerException exception, Object ctx) { + log.error("Failed to add {}", msg, exception); + latch.countDown(); + } + }; + executor.execute(() -> { + synchronized (ml) { + batchSizes.add((long) batchSize); + indexesBeforeAdd.add(interceptor.getIndex() + 1); // index is updated in each asyncAddEntry call + ml.asyncAddEntry(Unpooled.wrappedBuffer(msg.getBytes()), batchSize, callback, null); + } + }); + } + assertTrue(latch.await(3, TimeUnit.SECONDS)); + for (int i = 1; i < numEntries; i++) { + final var sum = batchSizes.get(i) + batchSizes.get(i - 1); + batchSizes.set(i, sum); + } + assertEquals(indexesBeforeAdd.subList(1, numEntries), batchSizes.subList(0, numEntries - 1)); + executor.shutdown(); + } } From 0ab44048e2d5ec6b49d4ca82dce48fdff730603b Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Thu, 6 Feb 2025 20:54:32 +0800 Subject: [PATCH 4/5] Improve the test --- .../ManagedLedgerInterceptorImplTest.java | 21 +++++++++++-------- 1 file changed, 12 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java index 35869b912ee84..8663019efb8c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/ManagedLedgerInterceptorImplTest.java @@ -507,16 +507,17 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { public void testBeforeAddEntry() throws Exception { final var interceptor = new ManagedLedgerInterceptorImpl(getBrokerEntryMetadataInterceptors(), null); final var config = new ManagedLedgerConfig(); - final var numEntries = 10; + final var numEntries = 100; config.setMaxEntriesPerLedger(numEntries); config.setManagedLedgerInterceptor(interceptor); @Cleanup final var ml = (ManagedLedgerImpl) factory.open("test_concurrent_add_entry", config); - final var indexesBeforeAdd = Collections.synchronizedList(new ArrayList()); - final var batchSizes = Collections.synchronizedList(new ArrayList()); + final var indexesBeforeAdd = new ArrayList(); + final var batchSizes = new ArrayList(); final var random = new Random(); final var latch = new CountDownLatch(numEntries); - final var executor = Executors.newFixedThreadPool(2); + final var executor = Executors.newFixedThreadPool(3); + final var lock = new Object(); // make sure `asyncAddEntry` are called in order for (int i = 0; i < numEntries; i++) { final var batchSize = random.nextInt(0, 100); final var msg = "msg-" + i; @@ -534,7 +535,7 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { } }; executor.execute(() -> { - synchronized (ml) { + synchronized (lock) { batchSizes.add((long) batchSize); indexesBeforeAdd.add(interceptor.getIndex() + 1); // index is updated in each asyncAddEntry call ml.asyncAddEntry(Unpooled.wrappedBuffer(msg.getBytes()), batchSize, callback, null); @@ -542,11 +543,13 @@ public void addFailed(ManagedLedgerException exception, Object ctx) { }); } assertTrue(latch.await(3, TimeUnit.SECONDS)); - for (int i = 1; i < numEntries; i++) { - final var sum = batchSizes.get(i) + batchSizes.get(i - 1); - batchSizes.set(i, sum); + synchronized (lock) { + for (int i = 1; i < numEntries; i++) { + final var sum = batchSizes.get(i) + batchSizes.get(i - 1); + batchSizes.set(i, sum); + } + assertEquals(indexesBeforeAdd.subList(1, numEntries), batchSizes.subList(0, numEntries - 1)); } - assertEquals(indexesBeforeAdd.subList(1, numEntries), batchSizes.subList(0, numEntries - 1)); executor.shutdown(); } } From 4f284835bea9ffbae0c757fc11a832699b72834c Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Sat, 8 Feb 2025 15:47:23 +0800 Subject: [PATCH 5/5] Fix shadow ML --- .../apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java | 1 - 1 file changed, 1 deletion(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java index 84b4bab343507..bae6cd66d2825 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ShadowManagedLedgerImpl.java @@ -240,7 +240,6 @@ protected void afterAddEntryToQueue(State state, OpAddEntry addOperation) throws log.debug("[{}] Add entry into shadow ledger lh={} entries={}, pos=({},{})", name, currentLedger.getId(), currentLedgerEntries, position.getLedgerId(), position.getEntryId()); } - pendingAddEntries.add(addOperation); if (position.getLedgerId() <= currentLedger.getId()) { // Write into lastLedger if (position.getLedgerId() == currentLedger.getId()) {