From 8499ade5d098a1a5340b294124d0703e4b9f0c3e Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Sat, 16 May 2026 10:40:52 +0800 Subject: [PATCH 1/2] [improve] Coalesce automatic managed ledger offload triggers Automatic offload can be triggered repeatedly while a previous offload is still running. Each trigger may read offload policies, scan the ledger list, fail to acquire the offload mutex, and schedule another 100ms retry. With slow offload or many ledgers, this can create unnecessary scheduler and executor pressure. Coalesce automatic triggers so there is at most one in-flight automatic offload and one pending rerun. If another trigger arrives during an in-flight offload, run one follow-up pass after the current offload completes. This keeps the final offload progression behavior while avoiding repeated retry loops, policy reads, and ledger scans. Keep explicit offload requests unchanged, and rename the automatic sentinel to make it clear that its Position value is not consumed. Add tests for trigger coalescing, coalesced reruns, and automatic state release when offload thresholds are disabled. --- .../impl/ManagedLedgerFactoryImpl.java | 4 +- .../mledger/impl/ManagedLedgerImpl.java | 76 +++++++++-- .../mledger/impl/OffloadPrefixTest.java | 122 ++++++++++++++++++ 3 files changed, 189 insertions(+), 13 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index 60e972e561be2..269cbafc20799 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -19,7 +19,7 @@ package org.apache.bookkeeper.mledger.impl; import static org.apache.bookkeeper.mledger.ManagedLedgerException.getManagedLedgerException; -import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.NULL_OFFLOAD_PROMISE; +import static org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER; import static org.apache.pulsar.common.util.Runnables.catchingAndLoggingThrowables; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Predicates; @@ -496,7 +496,7 @@ public void initializeComplete() { future.complete(newledger); // May need to trigger offloading if (config.isTriggerOffloadOnTopicLoad()) { - newledger.maybeOffloadInBackground(NULL_OFFLOAD_PROMISE); + newledger.maybeOffloadInBackground(AUTOMATIC_OFFLOAD_TRIGGER); } }); } 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 4a1a3d12ab075..dc8a83db4cdb4 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 @@ -231,8 +231,19 @@ public Logger getLogger() { protected final CallbackMutex trimmerMutex = new CallbackMutex(); protected final CallbackMutex offloadMutex = new CallbackMutex(); - public static final CompletableFuture NULL_OFFLOAD_PROMISE = CompletableFuture + // Automatic offload has no caller-visible future. Coalesce concurrent automatic triggers into at most one + // running offload and one follow-up run. + private final AtomicBoolean automaticOffloadInProgress = new AtomicBoolean(false); + private final AtomicBoolean automaticOffloadRerunRequested = new AtomicBoolean(false); + // Identity sentinel for automatic offload requests. The completed Position value is not used. + public static final CompletableFuture AUTOMATIC_OFFLOAD_TRIGGER = CompletableFuture .completedFuture(PositionFactory.LATEST); + + private enum OffloadRequestSource { + AUTOMATIC, + EXPLICIT + } + @VisibleForTesting @Getter protected volatile LedgerHandle currentLedger; @@ -1968,7 +1979,7 @@ synchronized void ledgerClosed(final LedgerHandle lh, Long lastAddConfirmed) { trimConsumedLedgersInBackground(); - maybeOffloadInBackground(NULL_OFFLOAD_PROMISE); + maybeOffloadInBackground(AUTOMATIC_OFFLOAD_TRIGGER); createLedgerAfterClosed(); } @@ -2804,22 +2815,66 @@ private void scheduleDeferredTrimming(boolean isTruncate, CompletableFuture p } public void maybeOffloadInBackground(CompletableFuture promise) { - if (getOffloadPoliciesIfAppendable().isEmpty()) { + if (promise == AUTOMATIC_OFFLOAD_TRIGGER) { + if (!automaticOffloadInProgress.compareAndSet(false, true)) { + automaticOffloadRerunRequested.set(true); + return; + } + CompletableFuture automaticOffloadCompletion = new CompletableFuture<>(); + automaticOffloadCompletion.whenComplete((res, ex) -> finishAutomaticOffload(ex)); + maybeOffloadInBackground(automaticOffloadCompletion, OffloadRequestSource.AUTOMATIC); + return; + } + + maybeOffloadInBackground(promise, OffloadRequestSource.EXPLICIT); + } + + private void maybeOffloadInBackground(CompletableFuture promise, OffloadRequestSource source) { + Optional> offloadThresholds = getOffloadThresholds(); + if (offloadThresholds.isEmpty()) { + // Explicit callers keep the previous no-completion behavior. The internal automatic completion must be + // finished so automaticOffloadInProgress can be cleared. + if (source == OffloadRequestSource.AUTOMATIC) { + promise.complete(PositionFactory.LATEST); + } return; } - final OffloadPolicies policies = config.getLedgerOffloader().getOffloadPolicies(); + Pair thresholds = offloadThresholds.get(); + executor.execute(() -> maybeOffload(thresholds.getLeft(), thresholds.getRight(), promise, + source)); + } + + private Optional> getOffloadThresholds() { + Optional optionalOffloadPolicies = getOffloadPoliciesIfAppendable(); + if (optionalOffloadPolicies.isEmpty()) { + return Optional.empty(); + } + + final OffloadPolicies policies = optionalOffloadPolicies.get(); final long offloadThresholdInBytes = Optional.ofNullable(policies.getManagedLedgerOffloadThresholdInBytes()).orElse(-1L); final long offloadThresholdInSeconds = Optional.ofNullable(policies.getManagedLedgerOffloadThresholdInSeconds()).orElse(-1L); if (offloadThresholdInBytes >= 0 || offloadThresholdInSeconds >= 0) { - executor.execute(() -> maybeOffload(offloadThresholdInBytes, offloadThresholdInSeconds, promise)); + return Optional.of(Pair.of(offloadThresholdInBytes, offloadThresholdInSeconds)); + } + + return Optional.empty(); + } + + private void finishAutomaticOffload(Throwable exception) { + if (exception != null) { + log.warn().exception(exception).log("Failed to automatically offload ledgers"); + } + automaticOffloadInProgress.set(false); + if (automaticOffloadRerunRequested.getAndSet(false)) { + maybeOffloadInBackground(AUTOMATIC_OFFLOAD_TRIGGER); } } private void maybeOffload(long offloadThresholdInBytes, long offloadThresholdInSeconds, - CompletableFuture finalPromise) { + CompletableFuture finalPromise, OffloadRequestSource source) { if (getOffloadPoliciesIfAppendable().isEmpty()) { String msg = String.format("[%s] Nothing to offload due to offloader or offloadPolicies is NULL", name); finalPromise.completeExceptionally(new IllegalArgumentException(msg)); @@ -2834,7 +2889,7 @@ private void maybeOffload(long offloadThresholdInBytes, long offloadThresholdInS } if (!offloadMutex.tryLock()) { - scheduledExecutor.schedule(() -> maybeOffloadInBackground(finalPromise), + scheduledExecutor.schedule(() -> maybeOffloadInBackground(finalPromise, source), 100, TimeUnit.MILLISECONDS); return; } @@ -2926,12 +2981,11 @@ void internalTrimConsumedLedgers(CompletableFuture promise) { private Optional getOffloadPoliciesIfAppendable() { LedgerOffloader ledgerOffloader = config.getLedgerOffloader(); - if (ledgerOffloader == null - || !ledgerOffloader.isAppendable() - || ledgerOffloader.getOffloadPolicies() == null) { + if (ledgerOffloader == null || !ledgerOffloader.isAppendable()) { return Optional.empty(); } - return Optional.ofNullable(ledgerOffloader.getOffloadPolicies()); + OffloadPolicies offloadPolicies = ledgerOffloader.getOffloadPolicies(); + return Optional.ofNullable(offloadPolicies); } @VisibleForTesting diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java index e991ceabf456a..5385f642c5713 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java @@ -1184,6 +1184,128 @@ public CompletableFuture offload(ReadHandle ledger, } } + @Test + public void automaticOffloadTriggersAreCoalescedWhileOffloadInProgress() throws Exception { + CompletableFuture slowOffload = new CompletableFuture<>(); + CountDownLatch offloadRunning = new CountDownLatch(1); + AtomicInteger offloadPolicyCalls = new AtomicInteger(); + MockLedgerOffloader offloader = new MockLedgerOffloader() { + @Override + public CompletableFuture offload(ReadHandle ledger, + UUID uuid, + Map extraMetadata) { + offloadRunning.countDown(); + return slowOffload.thenCompose((res) -> super.offload(ledger, uuid, extraMetadata)); + } + + @Override + public OffloadPoliciesImpl getOffloadPolicies() { + offloadPolicyCalls.incrementAndGet(); + return super.getOffloadPolicies(); + } + }; + + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setMaxEntriesPerLedger(10); + config.setRetentionTime(10, TimeUnit.MINUTES); + config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(0L); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); + config.setLedgerOffloader(offloader); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + + for (int i = 0; i < 25; i++) { + ledger.addEntry(buildEntry(10, "entry-" + i)); + } + assertTrue(offloadRunning.await(5, TimeUnit.SECONDS)); + + int callsBeforeRepeatedTriggers = offloadPolicyCalls.get(); + for (int i = 0; i < 20; i++) { + ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER); + } + + Thread.sleep(300); + assertTrue(offloadPolicyCalls.get() < callsBeforeRepeatedTriggers + 5, + "Repeated automatic triggers should not create independent retry loops"); + + slowOffload.complete(null); + + assertEventuallyTrue(() -> offloader.offloadedLedgers().size() == 2); + List allLedgerIds = ledger.getLedgersInfoAsList().stream().map(LedgerInfo::getLedgerId).toList(); + assertEquals(offloader.offloadedLedgers(), Set.of(allLedgerIds.get(0), allLedgerIds.get(1))); + } + + @Test + public void automaticOffloadRunsAgainForCoalescedTrigger() throws Exception { + CompletableFuture slowOffload = new CompletableFuture<>(); + CountDownLatch offloadRunning = new CountDownLatch(1); + MockLedgerOffloader offloader = new MockLedgerOffloader() { + @Override + public CompletableFuture offload(ReadHandle ledger, + UUID uuid, + Map extraMetadata) { + offloadRunning.countDown(); + return slowOffload.thenCompose((res) -> super.offload(ledger, uuid, extraMetadata)); + } + }; + + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setMaxEntriesPerLedger(10); + config.setRetentionTime(10, TimeUnit.MINUTES); + config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(0L); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); + config.setLedgerOffloader(offloader); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + + for (int i = 0; i < 11; i++) { + ledger.addEntry(buildEntry(10, "entry-" + i)); + } + assertTrue(offloadRunning.await(5, TimeUnit.SECONDS)); + + // The next ledger closes after the first automatic scan, so it depends on the coalesced rerun. + for (int i = 11; i < 21; i++) { + ledger.addEntry(buildEntry(10, "entry-" + i)); + } + assertEquals(offloader.offloadedLedgers().size(), 0); + + slowOffload.complete(null); + + assertEventuallyTrue(() -> offloader.offloadedLedgers().size() == 2); + List allLedgerIds = ledger.getLedgersInfoAsList().stream().map(LedgerInfo::getLedgerId).toList(); + assertEquals(offloader.offloadedLedgers(), Set.of(allLedgerIds.get(0), allLedgerIds.get(1))); + } + + @Test + public void automaticOffloadWithoutThresholdDoesNotBlockLaterTriggers() throws Exception { + MockLedgerOffloader offloader = new MockLedgerOffloader(); + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setMaxEntriesPerLedger(10); + config.setRetentionTime(10, TimeUnit.MINUTES); + config.setRetentionSizeInMB(10); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(-1L); + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); + config.setLedgerOffloader(offloader); + + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + + for (int i = 0; i < 25; i++) { + ledger.addEntry(buildEntry(10, "entry-" + i)); + } + ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER); + assertEquals(offloader.offloadedLedgers().size(), 0); + + // A disabled automatic trigger must complete internally so a later enabled trigger can run. + offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(0L); + ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER); + + assertEventuallyTrue(() -> offloader.offloadedLedgers().size() == 2); + List allLedgerIds = ledger.getLedgersInfoAsList().stream().map(LedgerInfo::getLedgerId).toList(); + assertEquals(offloader.offloadedLedgers(), Set.of(allLedgerIds.get(0), allLedgerIds.get(1))); + } + @DataProvider(name = "offloadAsSoonAsClosed") public Object[][] offloadAsSoonAsClosedProvider() { return new Object[][]{ From 915a73362d42e11c69446b25d6fa411c8b1f51c0 Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Wed, 3 Jun 2026 22:29:08 +0800 Subject: [PATCH 2/2] [fix][ml] Fix automatic offload trigger coalescing race --- .../AutomaticOffloadTriggerController.java | 86 +++++++++++++++++++ .../mledger/impl/ManagedLedgerImpl.java | 60 +++++++------ ...AutomaticOffloadTriggerControllerTest.java | 85 ++++++++++++++++++ .../mledger/impl/OffloadPrefixTest.java | 16 ++-- 4 files changed, 216 insertions(+), 31 deletions(-) create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerController.java create mode 100644 managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerControllerTest.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerController.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerController.java new file mode 100644 index 0000000000000..35f5a26706371 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerController.java @@ -0,0 +1,86 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.impl; + +import java.util.concurrent.atomic.AtomicInteger; + +/** + * Coalesces repeated automatic offload triggers into at most one active run and one follow-up run. + */ +final class AutomaticOffloadTriggerController { + private static final int IDLE = 0; + private static final int RUNNING = 1; + private static final int RUNNING_WITH_PENDING_TRIGGER = 2; + + private final AtomicInteger state = new AtomicInteger(IDLE); + + /** + * Records an automatic offload trigger. + * + * @return true when the caller must start a new automatic offload run + */ + boolean requestRun() { + while (true) { + int current = state.get(); + switch (current) { + case IDLE: + if (state.compareAndSet(IDLE, RUNNING)) { + return true; + } + break; + case RUNNING: + if (state.compareAndSet(RUNNING, RUNNING_WITH_PENDING_TRIGGER)) { + return false; + } + break; + case RUNNING_WITH_PENDING_TRIGGER: + return false; + default: + throw new IllegalStateException("Unknown automatic offload trigger state: " + current); + } + } + } + + /** + * Records completion of the current automatic offload run. + * + * @return true when the caller must immediately start one coalesced follow-up run + */ + boolean completeRun() { + while (true) { + int current = state.get(); + switch (current) { + case IDLE: + return false; + case RUNNING: + if (state.compareAndSet(RUNNING, IDLE)) { + return false; + } + break; + case RUNNING_WITH_PENDING_TRIGGER: + if (state.compareAndSet(RUNNING_WITH_PENDING_TRIGGER, RUNNING)) { + return true; + } + break; + default: + throw new IllegalStateException("Unknown automatic offload trigger state: " + current); + } + } + } +} 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 dc8a83db4cdb4..a06d0bf553f2d 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 @@ -231,10 +231,8 @@ public Logger getLogger() { protected final CallbackMutex trimmerMutex = new CallbackMutex(); protected final CallbackMutex offloadMutex = new CallbackMutex(); - // Automatic offload has no caller-visible future. Coalesce concurrent automatic triggers into at most one - // running offload and one follow-up run. - private final AtomicBoolean automaticOffloadInProgress = new AtomicBoolean(false); - private final AtomicBoolean automaticOffloadRerunRequested = new AtomicBoolean(false); + private final AutomaticOffloadTriggerController automaticOffloadTriggerController = + new AutomaticOffloadTriggerController(); // Identity sentinel for automatic offload requests. The completed Position value is not used. public static final CompletableFuture AUTOMATIC_OFFLOAD_TRIGGER = CompletableFuture .completedFuture(PositionFactory.LATEST); @@ -244,6 +242,9 @@ private enum OffloadRequestSource { EXPLICIT } + private record OffloadThresholds(long thresholdInBytes, long thresholdInSeconds) { + } + @VisibleForTesting @Getter protected volatile LedgerHandle currentLedger; @@ -2816,36 +2817,44 @@ private void scheduleDeferredTrimming(boolean isTruncate, CompletableFuture p public void maybeOffloadInBackground(CompletableFuture promise) { if (promise == AUTOMATIC_OFFLOAD_TRIGGER) { - if (!automaticOffloadInProgress.compareAndSet(false, true)) { - automaticOffloadRerunRequested.set(true); - return; + if (automaticOffloadTriggerController.requestRun()) { + startAutomaticOffload(); } - CompletableFuture automaticOffloadCompletion = new CompletableFuture<>(); - automaticOffloadCompletion.whenComplete((res, ex) -> finishAutomaticOffload(ex)); - maybeOffloadInBackground(automaticOffloadCompletion, OffloadRequestSource.AUTOMATIC); return; } maybeOffloadInBackground(promise, OffloadRequestSource.EXPLICIT); } + private void startAutomaticOffload() { + CompletableFuture automaticOffloadCompletion = new CompletableFuture<>(); + automaticOffloadCompletion.whenComplete((res, ex) -> finishAutomaticOffload(ex)); + try { + maybeOffloadInBackground(automaticOffloadCompletion, OffloadRequestSource.AUTOMATIC); + } catch (RuntimeException e) { + automaticOffloadCompletion.completeExceptionally(e); + } + } + private void maybeOffloadInBackground(CompletableFuture promise, OffloadRequestSource source) { - Optional> offloadThresholds = getOffloadThresholds(); + Optional offloadThresholds = getOffloadThresholds(); if (offloadThresholds.isEmpty()) { - // Explicit callers keep the previous no-completion behavior. The internal automatic completion must be - // finished so automaticOffloadInProgress can be cleared. if (source == OffloadRequestSource.AUTOMATIC) { promise.complete(PositionFactory.LATEST); } return; } - Pair thresholds = offloadThresholds.get(); - executor.execute(() -> maybeOffload(thresholds.getLeft(), thresholds.getRight(), promise, - source)); + OffloadThresholds thresholds = offloadThresholds.get(); + try { + executor.execute(() -> maybeOffload(thresholds.thresholdInBytes(), thresholds.thresholdInSeconds(), + promise, source)); + } catch (RuntimeException e) { + promise.completeExceptionally(e); + } } - private Optional> getOffloadThresholds() { + private Optional getOffloadThresholds() { Optional optionalOffloadPolicies = getOffloadPoliciesIfAppendable(); if (optionalOffloadPolicies.isEmpty()) { return Optional.empty(); @@ -2857,7 +2866,7 @@ private Optional> getOffloadThresholds() { final long offloadThresholdInSeconds = Optional.ofNullable(policies.getManagedLedgerOffloadThresholdInSeconds()).orElse(-1L); if (offloadThresholdInBytes >= 0 || offloadThresholdInSeconds >= 0) { - return Optional.of(Pair.of(offloadThresholdInBytes, offloadThresholdInSeconds)); + return Optional.of(new OffloadThresholds(offloadThresholdInBytes, offloadThresholdInSeconds)); } return Optional.empty(); @@ -2865,11 +2874,10 @@ private Optional> getOffloadThresholds() { private void finishAutomaticOffload(Throwable exception) { if (exception != null) { - log.warn().exception(exception).log("Failed to automatically offload ledgers"); + log.debug().exception(exception).log("Failed to automatically offload ledgers"); } - automaticOffloadInProgress.set(false); - if (automaticOffloadRerunRequested.getAndSet(false)) { - maybeOffloadInBackground(AUTOMATIC_OFFLOAD_TRIGGER); + if (automaticOffloadTriggerController.completeRun()) { + startAutomaticOffload(); } } @@ -2889,8 +2897,12 @@ private void maybeOffload(long offloadThresholdInBytes, long offloadThresholdInS } if (!offloadMutex.tryLock()) { - scheduledExecutor.schedule(() -> maybeOffloadInBackground(finalPromise, source), - 100, TimeUnit.MILLISECONDS); + try { + scheduledExecutor.schedule(() -> maybeOffloadInBackground(finalPromise, source), + 100, TimeUnit.MILLISECONDS); + } catch (RuntimeException e) { + finalPromise.completeExceptionally(e); + } return; } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerControllerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerControllerTest.java new file mode 100644 index 0000000000000..9aefc7e3e1e1f --- /dev/null +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/AutomaticOffloadTriggerControllerTest.java @@ -0,0 +1,85 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.bookkeeper.mledger.impl; + +import static org.assertj.core.api.Assertions.assertThat; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import org.testng.annotations.Test; + +public class AutomaticOffloadTriggerControllerTest { + + @Test + public void triggersCoalesceWhileRunIsActive() { + AutomaticOffloadTriggerController controller = new AutomaticOffloadTriggerController(); + + assertThat(controller.requestRun()).isTrue(); + assertThat(controller.requestRun()).isFalse(); + assertThat(controller.requestRun()).isFalse(); + } + + @Test + public void pendingTriggerSchedulesOneFollowUpRun() { + AutomaticOffloadTriggerController controller = new AutomaticOffloadTriggerController(); + + assertThat(controller.requestRun()).isTrue(); + assertThat(controller.requestRun()).isFalse(); + + assertThat(controller.completeRun()).isTrue(); + assertThat(controller.completeRun()).isFalse(); + assertThat(controller.requestRun()).isTrue(); + } + + @Test(timeOut = 30000) + public void concurrentTriggerAndCompletionAlwaysReserveOneFollowUpRun() throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + for (int i = 0; i < 1000; i++) { + AutomaticOffloadTriggerController controller = new AutomaticOffloadTriggerController(); + assertThat(controller.requestRun()).isTrue(); + + // Completion and a new trigger can race; exactly one side must reserve the follow-up run. + CyclicBarrier barrier = new CyclicBarrier(3); + Future completeResult = executor.submit(() -> { + barrier.await(5, TimeUnit.SECONDS); + return controller.completeRun(); + }); + Future triggerResult = executor.submit(() -> { + barrier.await(5, TimeUnit.SECONDS); + return controller.requestRun(); + }); + + barrier.await(5, TimeUnit.SECONDS); + boolean followUpReservedByComplete = completeResult.get(5, TimeUnit.SECONDS); + boolean followUpReservedByTrigger = triggerResult.get(5, TimeUnit.SECONDS); + + assertThat(followUpReservedByComplete) + .as("iteration %s must reserve exactly one follow-up run", i) + .isNotEqualTo(followUpReservedByTrigger); + assertThat(controller.completeRun()).isFalse(); + } + } finally { + executor.shutdownNow(); + executor.awaitTermination(5, TimeUnit.SECONDS); + } + } +} diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java index 5385f642c5713..10f72cb8d53b7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java @@ -1213,21 +1213,21 @@ public OffloadPoliciesImpl getOffloadPolicies() { offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); config.setLedgerOffloader(offloader); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + ManagedLedgerImpl ledger = + (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); for (int i = 0; i < 25; i++) { ledger.addEntry(buildEntry(10, "entry-" + i)); } assertTrue(offloadRunning.await(5, TimeUnit.SECONDS)); + // Repeated automatic triggers should stop at the controller and avoid another policy lookup. int callsBeforeRepeatedTriggers = offloadPolicyCalls.get(); for (int i = 0; i < 20; i++) { ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER); } - Thread.sleep(300); - assertTrue(offloadPolicyCalls.get() < callsBeforeRepeatedTriggers + 5, - "Repeated automatic triggers should not create independent retry loops"); + assertEquals(offloadPolicyCalls.get(), callsBeforeRepeatedTriggers); slowOffload.complete(null); @@ -1258,7 +1258,8 @@ public CompletableFuture offload(ReadHandle ledger, offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); config.setLedgerOffloader(offloader); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + ManagedLedgerImpl ledger = + (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); for (int i = 0; i < 11; i++) { ledger.addEntry(buildEntry(10, "entry-" + i)); @@ -1289,7 +1290,8 @@ public void automaticOffloadWithoutThresholdDoesNotBlockLaterTriggers() throws E offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInSeconds(null); config.setLedgerOffloader(offloader); - ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); + ManagedLedgerImpl ledger = + (ManagedLedgerImpl) factory.open("my_test_ledger" + UUID.randomUUID(), config); for (int i = 0; i < 25; i++) { ledger.addEntry(buildEntry(10, "entry-" + i)); @@ -1297,7 +1299,7 @@ public void automaticOffloadWithoutThresholdDoesNotBlockLaterTriggers() throws E ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER); assertEquals(offloader.offloadedLedgers().size(), 0); - // A disabled automatic trigger must complete internally so a later enabled trigger can run. + // A disabled automatic trigger must complete internally so a later valid trigger can run. offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(0L); ledger.maybeOffloadInBackground(ManagedLedgerImpl.AUTOMATIC_OFFLOAD_TRIGGER);