From f6acd6dc39168d27edcce9a06d3f2efe493326e5 Mon Sep 17 00:00:00 2001 From: dao-jun Date: Tue, 19 May 2026 02:35:19 +0800 Subject: [PATCH 1/2] Add cursorless RandomReader to ManagedLedger. --- .../bookkeeper/mledger/ManagedLedger.java | 8 + .../bookkeeper/mledger/RandomReader.java | 79 +++ .../mledger/impl/ManagedLedgerImpl.java | 48 +- .../bookkeeper/mledger/impl/OpAddEntry.java | 3 + .../mledger/impl/OpReadEntries.java | 261 ++++++++++ .../mledger/impl/RandomReaderImpl.java | 105 ++++ .../mledger/impl/RandomReaders.java | 62 +++ .../mledger/impl/ShadowManagedLedgerImpl.java | 1 + .../mledger/impl/cache/EntryCache.java | 3 +- .../mledger/impl/ManagedLedgerTest.java | 459 ++++++++++++++++++ 10 files changed, 1026 insertions(+), 3 deletions(-) create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 0455f0efa8bb6..67f6901cebb07 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -764,6 +764,14 @@ default ManagedLedgerAttributes getManagedLedgerAttributes() { void asyncReadEntry(Position position, AsyncCallbacks.ReadEntryCallback callback, Object ctx); + /** + * Create a standalone cursorless reader. + * @throws UnsupportedOperationException when the managed-ledger implementation does not support random reads + */ + default RandomReader newRandomReader() { + throw new UnsupportedOperationException("RandomReader is not supported by this ManagedLedger implementation"); + } + /** * Get all the managed ledgers. */ diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java new file mode 100644 index 0000000000000..e2c03983ce549 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java @@ -0,0 +1,79 @@ +/* + * 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; + +import java.io.Closeable; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import org.apache.bookkeeper.common.annotation.InterfaceAudience; +import org.apache.bookkeeper.common.annotation.InterfaceStability; +import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; + +/** + * A cursorless, stateless reader for entries that are currently available in a managed ledger. + * + *

A random reader does not maintain a read position, acknowledge entries, contribute to backlog, or prevent ledger + * trimming. A read can therefore fail, or start at the next retained ledger, when ledgers are trimmed concurrently. + * Reads do not wait for future entries. + * + *

Callers must release every returned {@link Entry}. Completion can run on a BookKeeper, Netty, or managed-ledger + * thread; callers that mutate thread-confined state must explicitly select an appropriate executor. + */ +@InterfaceAudience.LimitedPrivate +@InterfaceStability.Evolving +public interface RandomReader extends Closeable { + + /** + * Read up to {@code numberOfEntries} starting at {@code startPosition}, inclusive. + */ + CompletableFuture> read(Position startPosition, int numberOfEntries); + + /** + * Read up to {@code maxPosition}, inclusive, without a size limit. + */ + default CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition) { + return read(startPosition, numberOfEntries, maxPosition, ManagedLedgerUtils.NO_MAX_SIZE_LIMIT); + } + + /** + * Read using an estimated-size limit and no position limit. + */ + default CompletableFuture> read(Position startPosition, int numberOfEntries, long maxSizeBytes) { + return read(startPosition, numberOfEntries, PositionFactory.LATEST, maxSizeBytes); + } + + /** + * Read entries subject to count, position, and estimated-size limits. + * + *

{@code maxPosition} is inclusive. A null value is equivalent to {@link PositionFactory#LATEST}. + * {@code maxSizeBytes} uses the same estimate-based cap as {@link ManagedCursor}; at least one entry can be + * returned even when that entry exceeds the requested size. + * + *

If a storage error occurs after entries have been collected, the future completes successfully with that + * partial list and the read stops. An error before the first entry completes the future exceptionally. + */ + CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition, + long maxSizeBytes); + + /** + * Unregister this reader. Closing does not cancel reads already in progress. + */ + @Override + void close(); +} 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..cb61e093e4486 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 @@ -122,6 +122,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionBound; import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.bookkeeper.mledger.RandomReader; import org.apache.bookkeeper.mledger.WaitingEntryCallBack; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; @@ -190,6 +191,7 @@ public Logger getLogger() { // ordered by read position (when cacheEvictionByMarkDeletedPosition=false) or by mark delete position // (when cacheEvictionByMarkDeletedPosition=true) private final ActiveManagedCursorContainer activeCursors; + private final RandomReaders randomReaders; // Ever-increasing counter of entries added @@ -381,6 +383,7 @@ public ManagedLedgerImpl(ManagedLedgerFactoryImpl factory, BookKeeper bookKeeper } else { activeCursors = new ManagedCursorContainerImpl(); } + randomReaders = new RandomReaders(this); this.factory = factory; this.bookKeeper = bookKeeper; this.config = config; @@ -1638,11 +1641,13 @@ public void closeFailed(ManagedLedgerException exception, Object ctx) { public synchronized void asyncClose(final CloseCallback callback, final Object ctx) { State state = STATE_UPDATER.get(this); if (state.isFenced()) { + randomReaders.closeAll(); cancelScheduledTasks(); factory.close(this); callback.closeFailed(new ManagedLedgerFencedException(), ctx); return; } else if (state == State.Closed) { + randomReaders.closeAll(); log.debug("Ignoring request to close a closed managed ledger"); callback.closeComplete(ctx); return; @@ -1652,6 +1657,7 @@ public synchronized void asyncClose(final CloseCallback callback, final Object c factory.close(this); STATE_UPDATER.set(this, State.Closed); + randomReaders.closeAll(); clearPendingAddEntries(new ManagedLedgerAlreadyClosedException("Managed ledger is closed")); cancelScheduledTasks(); @@ -2330,6 +2336,11 @@ public void asyncReadEntry(Position position, ReadEntryCallback callback, Object } + @Override + public RandomReader newRandomReader() { + return randomReaders.create(); + } + private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) { if (opReadEntry.readPosition.compareTo(opReadEntry.maxPosition) > 0) { @@ -2451,6 +2462,32 @@ protected void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry } } + protected void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry, ReadEntriesCallback callback, + Object ctx) { + asyncReadEntry(ledger, firstEntry, lastEntry, () -> 0, callback, ctx); + } + + void asyncReadEntryForRandomReader(ReadHandle ledger, long firstEntry, long lastEntry, + ReadEntriesCallback callback, Object ctx) { + // Random-reader misses must be admitted to the shared cache for reuse by later positional reads. + asyncReadEntry(ledger, firstEntry, lastEntry, () -> 1, callback, ctx); + } + + private void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry, + IntSupplier expectedReadCount, ReadEntriesCallback callback, Object ctx) { + if (config.getReadEntryTimeoutSeconds() > 0) { + // set readOpCount to uniquely validate if ReadEntryCallbackWrapper is already recycled + long readOpCount = READ_OP_COUNT_UPDATER.incrementAndGet(this); + long createdTime = System.nanoTime(); + ReadEntryCallbackWrapper readCallback = ReadEntryCallbackWrapper.create(name, ledger.getId(), firstEntry, + callback, readOpCount, createdTime, ctx); + lastReadCallback = readCallback; + entryCache.asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, readCallback, readOpCount); + } else { + entryCache.asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, callback, ctx); + } + } + static final class ReadEntryCallbackWrapper implements ReadEntryCallback, ReadEntriesCallback { volatile ReadEntryCallback readEntryCallback; @@ -4441,6 +4478,7 @@ public synchronized void setFenced() { log.info().log("Moving to Fenced state"); State prev = STATE_UPDATER.getAndSet(this, State.Fenced); if (prev != State.Fenced) { + randomReaders.closeAll(); clearPendingAddEntries(new ManagedLedgerFencedException("ManagedLedger " + name + " is fenced")); } @@ -4450,6 +4488,7 @@ synchronized void setFencedForDeletion() { log.info().log("Moving to FencedForDeletion state"); State prev = STATE_UPDATER.getAndSet(this, State.FencedForDeletion); if (prev != State.FencedForDeletion) { + randomReaders.closeAll(); clearPendingAddEntries(new ManagedLedgerFencedException("ManagedLedger " + name + " is fenced")); } @@ -5185,7 +5224,12 @@ public void waitForPendingCacheEvictions() { } boolean shouldCacheAddedEntry() { - // Avoid caching entries if no cursor has been created - return getActiveCursors().shouldCacheAddedEntry(); + // Random readers are deliberately not active cursors, but still need add-path cache seeding. + return getActiveCursors().shouldCacheAddedEntry() || randomReaders.hasActiveReaders(); + } + + @VisibleForTesting + int getActiveRandomReaderCount() { + return randomReaders.size(); } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java index 2079caf08a36a..dfb1354f9cc94 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java @@ -268,6 +268,9 @@ public void run() { long ledgerId = ledger != null ? ledger.getId() : ((Position) ctx).getLedgerId(); // Handle caching for tailing reads + // Cursorless readers are deliberately not active cursors. Their explicit admission signal keeps add-path + // seeding independent from cursor accounting. EntryCache.insert must remain unconditional with respect to the + // expected-read count: on a zero-cursor topic this entry intentionally has a null read-count handler. if (ml.shouldCacheAddedEntry()) { int expectedReadCount = 0; // only use expectedReadCount if cache eviction is enabled by expected read count diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java new file mode 100644 index 0000000000000..a0978229945cf --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java @@ -0,0 +1,261 @@ +/* + * 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 java.lang.Math.min; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicBoolean; +import lombok.CustomLog; +import org.apache.bookkeeper.client.LedgerHandle; +import org.apache.bookkeeper.client.api.ReadHandle; +import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedgerException; +import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerFencedException; +import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.State; +import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo; +import org.apache.pulsar.common.util.FutureUtil; + +@CustomLog +class OpReadEntries implements ReadEntriesCallback { + private final ManagedLedgerImpl ledger; + private final Position maxPosition; + private final int count; + private final CompletableFuture> promise = new CompletableFuture<>(); + private final List entries = new ArrayList<>(); + private final AtomicBoolean terminal = new AtomicBoolean(); + private Position readPosition; + private Position nextReadPosition; + + private OpReadEntries(ManagedLedgerImpl ledger, Position readPosition, int count, Position maxPosition) { + this.ledger = ledger; + this.readPosition = ledger.startReadOperationOnLedger(readPosition); + this.count = count; + this.maxPosition = maxPosition; + this.nextReadPosition = this.readPosition; + } + + static CompletableFuture> read(ManagedLedgerImpl ledger, Position readPosition, int count, + Position maxPosition) { + OpReadEntries op = new OpReadEntries(ledger, readPosition, count, maxPosition); + op.readEntries(); + return op.promise; + } + + void readEntries() { + if (terminal.get()) { + return; + } + final State state = ManagedLedgerImpl.STATE_UPDATER.get(ledger); + if (state.isFenced() || state == State.Closed) { + readEntriesFailed(new ManagedLedgerFencedException(), null); + return; + } + + if (readPosition.compareTo(maxPosition) > 0) { + checkReadCompletion(); + return; + } + + long ledgerId = readPosition.getLedgerId(); + LedgerHandle currentLedger = ledger.currentLedger; + + if (currentLedger != null && ledgerId == currentLedger.getId()) { + // Current writing ledger is not in the cache (since we don't want + // it to be automatically evicted), and we cannot use 2 different + // ledger handles (read & write)for the same ledger. + internalReadFromLedger(currentLedger); + } else { + LedgerInfo ledgerInfo = ledger.ledgers.get(ledgerId); + if (ledgerInfo == null || ledgerInfo.getEntries() == 0) { + updateReadPosition(getNextLedgerPosition(ledgerId)); + checkReadCompletion(); + return; + } + + ledger.getLedgerHandle(ledgerId).thenAccept(this::internalReadFromLedger) + .exceptionally(ex -> { + ledger.log.error().attr("position", readPosition).exceptionMessage(ex) + .log("Error opening ledger for reading"); + readEntriesFailed(ManagedLedgerException.getManagedLedgerException( + FutureUtil.unwrapCompletionException(ex)), null); + return null; + }); + } + } + + private void internalReadFromLedger(ReadHandle readHandle) { + long firstEntry = readPosition.getEntryId(); + long lastEntryInLedger; + + Position lastPosition = ledger.lastConfirmedEntry; + + if (readHandle.getId() == lastPosition.getLedgerId()) { + // For the current ledger, we only give read visibility to the last entry we have received a confirmation in + // the managed ledger layer + lastEntryInLedger = lastPosition.getEntryId(); + } else { + // For other ledgers, already closed the BK lastAddConfirmed is appropriate + lastEntryInLedger = readHandle.getLastAddConfirmed(); + } + + if (readHandle.getId() == maxPosition.getLedgerId()) { + lastEntryInLedger = min(maxPosition.getEntryId(), lastEntryInLedger); + } + + if (firstEntry > lastEntryInLedger) { + log.debug().attr("ledgerId", readHandle.getId()) + .attr("lastEntry", lastEntryInLedger) + .attr("readEntry", firstEntry) + .log("No more messages to read from ledger"); + + LedgerHandle currentLedger = ledger.currentLedger; + if (currentLedger == null || readHandle.getId() != currentLedger.getId()) { + updateReadPosition(getNextLedgerPosition(readHandle.getId())); + } else { + updateReadPosition(readPosition); + } + + checkReadCompletion(); + return; + } + + long lastEntry = min(firstEntry + getNumberOfEntriesToRead() - 1, lastEntryInLedger); + + log.debug().attr("ledgerId", readHandle.getId()) + .attr("firstEntry", firstEntry) + .attr("lastEntry", lastEntry) + .log("Reading entries from ledger"); + ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry, this, null); + } + + private Position getNextLedgerPosition(long ledgerId) { + Long nextLedgerId = ledger.ledgers.ceilingKey(ledgerId + 1); + return PositionFactory.create(nextLedgerId != null ? nextLedgerId : ledgerId + 1, 0); + } + + @Override + public void readEntriesComplete(List returnedEntries, Object ctx) { + if (terminal.get()) { + returnedEntries.forEach(Entry::release); + return; + } + try { + internalReadEntriesComplete(returnedEntries); + } catch (Throwable throwable) { + log.error().attr("op", this).exception(throwable) + .log("Fallback to readEntriesFailed for exception in readEntriesComplete"); + readEntriesFailed(ManagedLedgerException.getManagedLedgerException(throwable), ctx); + } + } + + private void internalReadEntriesComplete(List returnedEntries) { + if (returnedEntries.isEmpty()) { + log.warn().attr("op", this).log("Read no entries unexpectedly"); + checkReadCompletion(); + return; + } + + log.debug() + .attr("managedLedger", ledger.getName()) + .attr("batchSize", returnedEntries.size()) + .attr("cumulativeSize", entries.size()) + .attr("requestedCount", count) + .log("Read entries succeeded"); + + entries.addAll(returnedEntries); + updateReadPosition(returnedEntries.get(returnedEntries.size() - 1).getPosition().getNext()); + checkReadCompletion(); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + try { + internalReadEntriesFailed(exception); + } catch (Throwable throwable) { + failPromise(ManagedLedgerException.getManagedLedgerException(throwable)); + } + } + + private void internalReadEntriesFailed(ManagedLedgerException exception) { + if (!entries.isEmpty()) { + log.warn() + .attr("managedLedger", ledger.getName()) + .attr("readPosition", readPosition) + .attr("partialEntries", entries.size()) + .exception(exception) + .log("Read failed after partial result"); + completeWithEntries(); + return; + } + + log.warn() + .attr("managedLedger", ledger.getName()) + .attr("readPosition", readPosition) + .exception(exception) + .log("Read failed from ledger"); + failPromise(exception); + } + + private void updateReadPosition(Position newReadPosition) { + nextReadPosition = newReadPosition; + } + + private void checkReadCompletion() { + if (entries.size() < count + && ledger.hasMoreEntries(nextReadPosition) + && maxPosition.compareTo(readPosition) > 0) { + ledger.getExecutor().execute(() -> { + readPosition = ledger.startReadOperationOnLedger(nextReadPosition); + readEntries(); + }); + } else { + completeWithEntries(); + } + } + + private void completeWithEntries() { + if (terminal.compareAndSet(false, true) && !promise.complete(entries)) { + // CompletableFuture is mutable: cancellation or external completion must not leak owned entries. + entries.forEach(Entry::release); + } + } + + private void failPromise(ManagedLedgerException exception) { + if (terminal.compareAndSet(false, true)) { + promise.completeExceptionally(exception); + } + } + + private int getNumberOfEntriesToRead() { + return count - entries.size(); + } + + + @Override + public String toString() { + return ledger.getName() + "{ readPosition: " + readPosition + ", maxPosition: " + maxPosition + + ", nextReadPosition: " + nextReadPosition + ", entries count: " + entries.size() + + ", count: " + count + " }"; + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java new file mode 100644 index 0000000000000..1a416874146df --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java @@ -0,0 +1,105 @@ +/* + * 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.apache.bookkeeper.mledger.impl.EntryCountEstimator.estimateEntryCountByBytesSize; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicBoolean; +import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerAlreadyClosedException; +import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.bookkeeper.mledger.RandomReader; +import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo; +import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; + +final class RandomReaderImpl implements RandomReader { + private final ManagedLedgerImpl ledger; + private final RandomReaders owner; + private final AtomicBoolean closed = new AtomicBoolean(); + + RandomReaderImpl(ManagedLedgerImpl ledger, RandomReaders owner) { + this.ledger = ledger; + this.owner = owner; + } + + @Override + public CompletableFuture> read(Position startPosition, int numberOfEntries) { + return read(startPosition, numberOfEntries, PositionFactory.LATEST, + ManagedLedgerUtils.NO_MAX_SIZE_LIMIT); + } + + @Override + public CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition, + long maxSizeBytes) { + if (closed.get() || owner.isClosed()) { + return CompletableFuture.failedFuture( + new ManagedLedgerAlreadyClosedException("Random reader is already closed")); + } + if (startPosition == null || numberOfEntries <= 0) { + return CompletableFuture.failedFuture(new IllegalArgumentException("Invalid parameters")); + } + + Position normalizedMaxPosition = maxPosition != null ? maxPosition : PositionFactory.LATEST; + Position normalizedStartPosition = normalizeStartPosition(startPosition); + if (normalizedStartPosition == null || normalizedMaxPosition.compareTo(normalizedStartPosition) < 0) { + return CompletableFuture.completedFuture(Collections.emptyList()); + } + + int effectiveCount = numberOfEntries; + if (maxSizeBytes != ManagedLedgerUtils.NO_MAX_SIZE_LIMIT) { + effectiveCount = Math.min(numberOfEntries, + estimateEntryCountByBytesSize(numberOfEntries, maxSizeBytes, normalizedStartPosition, ledger)); + } + + return OpReadEntries.read(ledger, normalizedStartPosition, effectiveCount, normalizedMaxPosition); + } + + private Position normalizeStartPosition(Position startPosition) { + if (PositionFactory.EARLIEST.equals(startPosition)) { + Map.Entry firstLedger = ledger.getLedgersInfo().firstEntry(); + return firstLedger != null ? PositionFactory.create(firstLedger.getKey(), 0) : null; + } + + Position lastPosition = ledger.getLastConfirmedEntry(); + if (PositionFactory.LATEST.equals(startPosition)) { + return lastPosition != null ? lastPosition.getNext() : null; + } + if (lastPosition == null || startPosition.compareTo(lastPosition) > 0) { + return null; + } + return ledger.isValidPosition(startPosition) + ? startPosition + : ledger.getNextValidPosition(startPosition); + } + + @Override + public void close() { + if (closed.compareAndSet(false, true)) { + owner.unregister(); + } + } + + boolean isClosed() { + return closed.get() || owner.isClosed(); + } +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java new file mode 100644 index 0000000000000..0742f2d355d30 --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java @@ -0,0 +1,62 @@ +/* + * 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; + +final class RandomReaders { + private final ManagedLedgerImpl ledger; + // Lifecycle mutations are synchronized; volatile reads keep add and read paths lock-free. + private volatile int activeReaders; + private volatile boolean closed; + + RandomReaders(ManagedLedgerImpl ledger) { + this.ledger = ledger; + } + + synchronized RandomReaderImpl create() { + if (closed) { + throw new IllegalStateException("Managed ledger is already closed"); + } + RandomReaderImpl reader = new RandomReaderImpl(ledger, this); + activeReaders++; + return reader; + } + + synchronized void unregister() { + if (activeReaders > 0) { + activeReaders--; + } + } + + boolean hasActiveReaders() { + return activeReaders > 0; + } + + int size() { + return activeReaders; + } + + boolean isClosed() { + return closed; + } + + synchronized void closeAll() { + closed = true; + activeReaders = 0; + } +} 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 05c96e27aed5c..01a2e638ac5b9 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 @@ -430,6 +430,7 @@ protected void updateLastLedgerCreatedTimeAndScheduleRolloverTask() { @Override boolean shouldCacheAddedEntry() { + // Shadow ledgers intentionally do not seed mirrored writes. Random reads still use the read-through cache. return false; } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCache.java index b2ebf7430560c..428318b262b32 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/EntryCache.java @@ -37,7 +37,8 @@ public interface EntryCache { String getName(); /** - * Insert an entry in the cache. + * Insert an entry in the cache. This method is an explicit cache-admission decision and must not reject an entry + * based on its expected read count. Entries without an expected-read handler are valid transient cache entries. * *

If the overall limit have been reached, this will trigger the eviction of other entries, possibly from * other EntryCache instances diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index bef94ecec5e64..055b5ad2dbb6f 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -76,6 +76,7 @@ import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.FutureTask; @@ -133,6 +134,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionBound; import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.bookkeeper.mledger.RandomReader; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; import org.apache.bookkeeper.mledger.impl.cache.EntryCache; @@ -202,6 +204,37 @@ public static void makeReadEntryProbFail(ManagedLedgerImpl ml, Supplier future) throws Exception { + try { + future.get(5, TimeUnit.SECONDS); + throw new AssertionError("Expected future to fail"); + } catch (ExecutionException e) { + return e.getCause(); + } + } + + private static void assertEntryPositionsAndRelease(List entries, Position... expectedPositions) { + try { + assertEquals(entries.size(), expectedPositions.length); + for (int i = 0; i < expectedPositions.length; i++) { + assertEquals(entries.get(i).getPosition(), expectedPositions[i]); + } + } finally { + entries.forEach(Entry::release); + } + } + + private static void assertEntryDataAndRelease(List entries, String... expectedData) { + try { + assertEquals(entries.size(), expectedData.length); + for (int i = 0; i < expectedData.length; i++) { + assertEquals(new String(entries.get(i).getData(), Encoding), expectedData[i]); + } + } finally { + entries.forEach(Entry::release); + } + } + @Data private static class DeleteLedgerInfo{ volatile boolean hasCalled; @@ -333,6 +366,432 @@ public void managedLedgerApi() throws Exception { ledger.close(); } + @Test(timeOut = 20000) + public void testRandomReader() throws Exception { + ManagedLedger ledger = factory.open("testRandomReader", + new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + assertEquals(reader.read(PositionFactory.EARLIEST, 10).get(5, TimeUnit.SECONDS), + Collections.emptyList()); + assertEquals(reader.read(PositionFactory.LATEST, 10).get(5, TimeUnit.SECONDS), + Collections.emptyList()); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + ledger.addEntry("entry-3".getBytes(Encoding)); + Position p4 = ledger.addEntry("entry-4".getBytes(Encoding)); + + assertEntryPositionsAndRelease(reader.read(PositionFactory.EARLIEST, 3).get(5, TimeUnit.SECONDS), + p0, p1, p2); + assertEntryPositionsAndRelease(reader.read(p1, 10).get(5, TimeUnit.SECONDS), p1, p2, + PositionFactory.create(p2.getLedgerId(), p2.getEntryId() + 1), p4); + assertEntryPositionsAndRelease( + reader.read(PositionFactory.create(p1.getLedgerId(), p1.getEntryId() + 100), 10) + .get(5, TimeUnit.SECONDS), + p2, PositionFactory.create(p2.getLedgerId(), p2.getEntryId() + 1), p4); + + assertEquals(reader.read(p4.getNext(), 10).get(5, TimeUnit.SECONDS), Collections.emptyList()); + assertEquals(reader.read(PositionFactory.LATEST, 10).get(5, TimeUnit.SECONDS), + Collections.emptyList()); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderRejectsInvalidArguments() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderRejectsInvalidArguments"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + assertTrue(expectFutureFailure(reader.read(null, 1)) instanceof IllegalArgumentException); + assertTrue(expectFutureFailure(reader.read(PositionFactory.EARLIEST, 0)) + instanceof IllegalArgumentException); + assertTrue(expectFutureFailure(reader.read(PositionFactory.EARLIEST, -1)) + instanceof IllegalArgumentException); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderPositionBoundaryValidation() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderPositionBoundaryValidation", + new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + + assertNotEquals(p1.getLedgerId(), p2.getLedgerId()); + + assertEntryPositionsAndRelease( + reader.read(PositionFactory.create(p0.getLedgerId(), -1), 1).get(5, TimeUnit.SECONDS), + p0); + assertEntryPositionsAndRelease(reader.read(PositionFactory.create(-2, 0), 1) + .get(5, TimeUnit.SECONDS), p0); + assertEntryPositionsAndRelease( + reader.read(PositionFactory.create(p0.getLedgerId(), Long.MAX_VALUE), 1) + .get(5, TimeUnit.SECONDS), + p2); + assertEquals(reader.read(PositionFactory.create(Long.MAX_VALUE, 0), 1).get(5, TimeUnit.SECONDS), + Collections.emptyList()); + assertEquals(reader.read(PositionFactory.create(p2.getLedgerId(), Long.MAX_VALUE), 1) + .get(5, TimeUnit.SECONDS), Collections.emptyList()); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderDoesNotWaitForFutureWrites() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderDoesNotWaitForFutureWrites"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + CompletableFuture> readAfterLast = reader.read(p0.getNext(), 10); + CompletableFuture> readLatest = reader.read(PositionFactory.LATEST, 10); + + assertEquals(readAfterLast.get(5, TimeUnit.SECONDS), Collections.emptyList()); + assertEquals(readLatest.get(5, TimeUnit.SECONDS), Collections.emptyList()); + + ledger.addEntry("entry-1".getBytes(Encoding)); + assertEquals(readAfterLast.get(5, TimeUnit.SECONDS), Collections.emptyList()); + assertEquals(readLatest.get(5, TimeUnit.SECONDS), Collections.emptyList()); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderExactCountBoundaries() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderExactCountBoundaries", + new ManagedLedgerConfig().setMaxEntriesPerLedger(3)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + Position p3 = ledger.addEntry("entry-3".getBytes(Encoding)); + Position p4 = ledger.addEntry("entry-4".getBytes(Encoding)); + + assertEquals(p0.getLedgerId(), p2.getLedgerId()); + assertNotEquals(p2.getLedgerId(), p3.getLedgerId()); + + assertEntryPositionsAndRelease(reader.read(p1, 1).get(5, TimeUnit.SECONDS), p1); + assertEntryPositionsAndRelease(reader.read(p1, 2).get(5, TimeUnit.SECONDS), p1, p2); + assertEntryPositionsAndRelease(reader.read(p1, 10).get(5, TimeUnit.SECONDS), p1, p2, p3, p4); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderSkipsEmptyLedgers() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderSkipsEmptyLedgers", + new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + ledger.ledgerClosed(ledger.currentLedger, 0L); + Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger), + ManagedLedgerImpl.State.LedgerOpened)); + LedgerHandle emptyLedger = ledger.currentLedger; + ledger.ledgerClosed(emptyLedger, -1L); + Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger), + ManagedLedgerImpl.State.LedgerOpened)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + + assertNotEquals(p0.getLedgerId(), p1.getLedgerId()); + assertFalse(ledger.getLedgersInfo().containsKey(emptyLedger.getId())); + assertEntryPositionsAndRelease( + reader.read(PositionFactory.create(emptyLedger.getId(), 0), 1).get(5, TimeUnit.SECONDS), + p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStopsOnErrorAndReturnsPartialEntries() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderStopsOnErrorAndReturnsPartialEntries", + new ManagedLedgerConfig().setMaxEntriesPerLedger(1)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + + assertNotEquals(p0.getLedgerId(), p1.getLedgerId()); + assertNotEquals(p1.getLedgerId(), p2.getLedgerId()); + + bkc.deleteLedger(p1.getLedgerId()); + + assertEntryPositionsAndRelease(reader.read(p0, 3).get(5, TimeUnit.SECONDS), p0); + + assertTrue(expectFutureFailure(reader.read(p1, 3)) instanceof ManagedLedgerException); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderFailsWhenClosed() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderFailsWhenClosed"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + Position position = ledger.addEntry("entry-0".getBytes(Encoding)); + + ledger.close(); + + assertTrue(expectFutureFailure(reader.read(position, 1)) + instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException); + } + + @Test(timeOut = 20000) + public void testRandomReaderIgnoresCursorAckState() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderIgnoresCursorAckState"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedCursor cursor = ledger.openCursor("c1"); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + + List cursorEntries = cursor.readEntries(2); + cursor.markDelete(cursorEntries.get(cursorEntries.size() - 1).getPosition()); + cursorEntries.forEach(Entry::release); + + assertEntryDataAndRelease(reader.read(p0, 2).get(5, TimeUnit.SECONDS), "entry-0", "entry-1"); + assertEntryPositionsAndRelease(reader.read(p1, 1).get(5, TimeUnit.SECONDS), p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxPositionBeforeStartReturnsEmpty() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxPositionBeforeStartReturnsEmpty"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + + // maxPosition is p0, startPosition is p1 -- nothing to read + assertEquals(reader.read(p1, 10, p0).get(5, TimeUnit.SECONDS), + Collections.emptyList()); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxPositionCapsSameLedger() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxPositionCapsSameLedger", + new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + Position p3 = ledger.addEntry("entry-3".getBytes(Encoding)); + + // Request 10 entries starting at p0 but capped at p1 -- should only get p0, p1 + assertEntryPositionsAndRelease( + reader.read(p0, 10, p1).get(5, TimeUnit.SECONDS), + p0, p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxPositionInclusive() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxPositionInclusive", + new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + + // maxPosition == p1 means p1 is the last readable entry + assertEntryPositionsAndRelease( + reader.read(p0, 10, p1).get(5, TimeUnit.SECONDS), + p0, p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxPositionCrossesLedger() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxPositionCrossesLedger", + new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + // ledger roll + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + Position p3 = ledger.addEntry("entry-3".getBytes(Encoding)); + + assertNotEquals(p1.getLedgerId(), p2.getLedgerId()); + + // maxPosition in the second ledger, should only read entries up to maxPosition + assertEntryPositionsAndRelease( + reader.read(p0, 10, p2).get(5, TimeUnit.SECONDS), + p0, p1, p2); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxPositionRespectsBothConstraints() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxPositionRespectsBothConstraints", + new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + Position p3 = ledger.addEntry("entry-3".getBytes(Encoding)); + Position p4 = ledger.addEntry("entry-4".getBytes(Encoding)); + + // Count=2 is tighter than maxPosition, should get only 2 entries + assertEntryPositionsAndRelease( + reader.read(p0, 2, p4).get(5, TimeUnit.SECONDS), + p0, p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMaxSizeBytesAndNullMaxPosition() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderMaxSizeBytesAndNullMaxPosition"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + byte[] data = new byte[100]; + Position p0 = ledger.addEntry(data); + Position p1 = ledger.addEntry(data); + Position p2 = ledger.addEntry(data); + Position p3 = ledger.addEntry(data); + + assertEntryPositionsAndRelease(reader.read(p0, 10, null, -1L).get(5, TimeUnit.SECONDS), p0, p1, p2, p3); + assertEntryPositionsAndRelease(reader.read(p0, 10, null, 1L).get(5, TimeUnit.SECONDS), p0); + // Average payload (100) plus BookKeeper overhead (64) gives an estimated two-entry budget. + assertEntryPositionsAndRelease(reader.read(p0, 10, null, 328L).get(5, TimeUnit.SECONDS), p0, p1); + + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderSeedsCacheOnZeroCursorTopic() throws Exception { + for (boolean expectedReadCountEviction : List.of(true, false)) { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setCacheEvictionByExpectedReadCount(expectedReadCountEviction); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderSeedsCacheOnZeroCursorTopic-" + expectedReadCountEviction, config); + + assertEquals(ledger.getActiveCursors().size(), 0); + assertFalse(ledger.shouldCacheAddedEntry()); + RandomReader reader = ledger.newRandomReader(); + assertEquals(ledger.getActiveRandomReaderCount(), 1); + assertTrue(ledger.shouldCacheAddedEntry()); + + Position position = ledger.addEntry("seeded".getBytes(Encoding)); + assertTrue(ledger.getCacheSize() > 0); + long hitsBefore = factory.getMbean().getCacheHitsTotal(); + long missesBefore = factory.getMbean().getCacheMissesTotal(); + assertEntryPositionsAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), position); + assertEquals(factory.getMbean().getCacheHitsTotal(), hitsBefore + 1); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore); + + reader.close(); + reader.close(); + assertEquals(ledger.getActiveRandomReaderCount(), 0); + ledger.entryCache.clear(); + ledger.addEntry("not-seeded".getBytes(Encoding)); + assertEquals(ledger.getCacheSize(), 0); + ledger.close(); + } + } + + @Test(timeOut = 20000) + public void testRandomReaderMissPopulatesReadThroughCache() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMissPopulatesReadThroughCache"); + Position position = ledger.addEntry("entry".getBytes(Encoding)); + assertEquals(ledger.getCacheSize(), 0); + @Cleanup RandomReader reader = ledger.newRandomReader(); + + long missesBefore = factory.getMbean().getCacheMissesTotal(); + assertEntryPositionsAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), position); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore + 1); + assertTrue(ledger.getCacheSize() > 0); + + long hitsBefore = factory.getMbean().getCacheHitsTotal(); + assertEntryPositionsAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), position); + assertEquals(factory.getMbean().getCacheHitsTotal(), hitsBefore + 1); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderLifecycleTracksManagedLedgerCloseAndFence() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderLifecycleClose"); + RandomReader first = ledger.newRandomReader(); + RandomReader second = ledger.newRandomReader(); + assertEquals(ledger.getActiveRandomReaderCount(), 2); + + first.close(); + first.close(); + assertEquals(ledger.getActiveRandomReaderCount(), 1); + ledger.close(); + assertEquals(ledger.getActiveRandomReaderCount(), 0); + assertTrue(((RandomReaderImpl) second).isClosed()); + assertTrue(expectFutureFailure(second.read(PositionFactory.EARLIEST, 1)) + instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException); + try { + ledger.newRandomReader(); + fail("Expected opening a reader on a closed ledger to fail"); + } catch (IllegalStateException expected) { + // expected + } + + ManagedLedgerImpl fencedLedger = (ManagedLedgerImpl) factory.open("testRandomReaderLifecycleFence"); + RandomReader fencedReader = fencedLedger.newRandomReader(); + fencedLedger.setFenced(); + assertEquals(fencedLedger.getActiveRandomReaderCount(), 0); + assertTrue(((RandomReaderImpl) fencedReader).isClosed()); + } + + @Test(timeOut = 20000) + public void testRandomReaderConcurrentReadsAreStateless() throws Exception { + ManagedLedger ledger = factory.open("testRandomReaderConcurrentReadsAreStateless"); + @Cleanup RandomReader reader = ledger.newRandomReader(); + Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); + Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); + Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); + Position p3 = ledger.addEntry("entry-3".getBytes(Encoding)); + + CompletableFuture> first = reader.read(p0, 2); + CompletableFuture> second = reader.read(p2, 2); + assertEntryPositionsAndRelease(first.get(5, TimeUnit.SECONDS), p0, p1); + assertEntryPositionsAndRelease(second.get(5, TimeUnit.SECONDS), p2, p3); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testOpReadEntriesSkipsEmptyLedger() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testOpReadEntriesSkipsEmptyLedger", + new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); + ledger.addEntry("entry-0".getBytes(Encoding)); + ledger.ledgerClosed(ledger.currentLedger, 0L); + Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger), + ManagedLedgerImpl.State.LedgerOpened)); + LedgerHandle emptyLedger = ledger.currentLedger; + ledger.ledgerClosed(emptyLedger, -1L); + Awaitility.await().untilAsserted(() -> assertEquals(ManagedLedgerImpl.STATE_UPDATER.get(ledger), + ManagedLedgerImpl.State.LedgerOpened)); + Position nextPosition = ledger.addEntry("entry-1".getBytes(Encoding)); + ledger.ledgers.put(emptyLedger.getId(), new LedgerInfo() + .setLedgerId(emptyLedger.getId()).setEntries(0).setSize(0)); + + CompletableFuture> promise = OpReadEntries.read(ledger, + PositionFactory.create(emptyLedger.getId(), 0), 1, PositionFactory.LATEST); + assertEntryPositionsAndRelease(promise.get(5, TimeUnit.SECONDS), nextPosition); + ledger.close(); + } + @Test(timeOut = 20000) public void simple() throws Exception { ManagedLedger ledger = factory.open("my_test_ledger"); From 7789c55886e3c6fd3545542f8fc9d9cbe6b28253 Mon Sep 17 00:00:00 2001 From: dao-jun Date: Fri, 17 Jul 2026 19:32:34 +0800 Subject: [PATCH 2/2] improve RandomReader --- .../bookkeeper/mledger/ManagedLedger.java | 8 - .../bookkeeper/mledger/RandomReader.java | 79 ----- .../mledger/impl/ManagedLedgerImpl.java | 32 +- .../bookkeeper/mledger/impl/OpAddEntry.java | 9 +- ...dEntries.java => OpRandomReadEntries.java} | 25 +- ...andomReaderImpl.java => RandomReader.java} | 22 +- .../mledger/impl/RandomReaders.java | 38 ++- .../mledger/impl/ManagedLedgerTest.java | 313 +++++++++++++++--- 8 files changed, 355 insertions(+), 171 deletions(-) delete mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java rename managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/{OpReadEntries.java => OpRandomReadEntries.java} (89%) rename managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/{RandomReaderImpl.java => RandomReader.java} (82%) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 67f6901cebb07..0455f0efa8bb6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -764,14 +764,6 @@ default ManagedLedgerAttributes getManagedLedgerAttributes() { void asyncReadEntry(Position position, AsyncCallbacks.ReadEntryCallback callback, Object ctx); - /** - * Create a standalone cursorless reader. - * @throws UnsupportedOperationException when the managed-ledger implementation does not support random reads - */ - default RandomReader newRandomReader() { - throw new UnsupportedOperationException("RandomReader is not supported by this ManagedLedger implementation"); - } - /** * Get all the managed ledgers. */ diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java deleted file mode 100644 index e2c03983ce549..0000000000000 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/RandomReader.java +++ /dev/null @@ -1,79 +0,0 @@ -/* - * 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; - -import java.io.Closeable; -import java.util.List; -import java.util.concurrent.CompletableFuture; -import org.apache.bookkeeper.common.annotation.InterfaceAudience; -import org.apache.bookkeeper.common.annotation.InterfaceStability; -import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; - -/** - * A cursorless, stateless reader for entries that are currently available in a managed ledger. - * - *

A random reader does not maintain a read position, acknowledge entries, contribute to backlog, or prevent ledger - * trimming. A read can therefore fail, or start at the next retained ledger, when ledgers are trimmed concurrently. - * Reads do not wait for future entries. - * - *

Callers must release every returned {@link Entry}. Completion can run on a BookKeeper, Netty, or managed-ledger - * thread; callers that mutate thread-confined state must explicitly select an appropriate executor. - */ -@InterfaceAudience.LimitedPrivate -@InterfaceStability.Evolving -public interface RandomReader extends Closeable { - - /** - * Read up to {@code numberOfEntries} starting at {@code startPosition}, inclusive. - */ - CompletableFuture> read(Position startPosition, int numberOfEntries); - - /** - * Read up to {@code maxPosition}, inclusive, without a size limit. - */ - default CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition) { - return read(startPosition, numberOfEntries, maxPosition, ManagedLedgerUtils.NO_MAX_SIZE_LIMIT); - } - - /** - * Read using an estimated-size limit and no position limit. - */ - default CompletableFuture> read(Position startPosition, int numberOfEntries, long maxSizeBytes) { - return read(startPosition, numberOfEntries, PositionFactory.LATEST, maxSizeBytes); - } - - /** - * Read entries subject to count, position, and estimated-size limits. - * - *

{@code maxPosition} is inclusive. A null value is equivalent to {@link PositionFactory#LATEST}. - * {@code maxSizeBytes} uses the same estimate-based cap as {@link ManagedCursor}; at least one entry can be - * returned even when that entry exceeds the requested size. - * - *

If a storage error occurs after entries have been collected, the future completes successfully with that - * partial list and the read stops. An error before the first entry completes the future exceptionally. - */ - CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition, - long maxSizeBytes); - - /** - * Unregister this reader. Closing does not cancel reads already in progress. - */ - @Override - void close(); -} 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 30e448c28aa27..ac54f3ae1e7aa 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 @@ -122,7 +122,6 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionBound; import org.apache.bookkeeper.mledger.PositionFactory; -import org.apache.bookkeeper.mledger.RandomReader; import org.apache.bookkeeper.mledger.WaitingEntryCallBack; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; @@ -2350,9 +2349,13 @@ public void asyncReadEntry(Position position, ReadEntryCallback callback, Object } - @Override - public RandomReader newRandomReader() { - return randomReaders.create(); + /** + * Returns the registry for creating RandomReader instances bound to this managed ledger's entry cache. + * RandomReader support is specific to this implementation and is intentionally not part of the + * {@link ManagedLedger} contract. + */ + public RandomReaders randomReaders() { + return randomReaders; } private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) { @@ -2476,15 +2479,11 @@ protected void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry } } - protected void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry, ReadEntriesCallback callback, - Object ctx) { - asyncReadEntry(ledger, firstEntry, lastEntry, () -> 0, callback, ctx); - } - void asyncReadEntryForRandomReader(ReadHandle ledger, long firstEntry, long lastEntry, - ReadEntriesCallback callback, Object ctx) { - // Random-reader misses must be admitted to the shared cache for reuse by later positional reads. - asyncReadEntry(ledger, firstEntry, lastEntry, () -> 1, callback, ctx); + IntSupplier expectedReadCount, ReadEntriesCallback callback) { + // expectedReadCount is resolved by the caller: cache-populating readers weight misses (>0), + // bypass readers use (0). + asyncReadEntry(ledger, firstEntry, lastEntry, expectedReadCount, callback, null); } private void asyncReadEntry(ReadHandle ledger, long firstEntry, long lastEntry, @@ -5309,12 +5308,17 @@ public void waitForPendingCacheEvictions() { } boolean shouldCacheAddedEntry() { - // Random readers are deliberately not active cursors, but still need add-path cache seeding. - return getActiveCursors().shouldCacheAddedEntry() || randomReaders.hasActiveReaders(); + // Only cache-populating random readers seed the tail; bypass readers do not. + return getActiveCursors().shouldCacheAddedEntry() || randomReaders.hasCachePopulatingReaders(); } @VisibleForTesting int getActiveRandomReaderCount() { return randomReaders.size(); } + + // Package-private: cache-populating readers contribute to expectedReadCount on the read and add paths. + int getActiveCachePopulatingRandomReaderCount() { + return randomReaders.cachePopulatingCount(); + } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java index dfb1354f9cc94..57adf624463bd 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java @@ -268,15 +268,14 @@ public void run() { long ledgerId = ledger != null ? ledger.getId() : ((Position) ctx).getLedgerId(); // Handle caching for tailing reads - // Cursorless readers are deliberately not active cursors. Their explicit admission signal keeps add-path - // seeding independent from cursor accounting. EntryCache.insert must remain unconditional with respect to the - // expected-read count: on a zero-cursor topic this entry intentionally has a null read-count handler. + // Streaming readers (like cursors) consume the tail, so they count toward expectedReadCount and eviction + // priority. With cacheEvictionByExpectedReadCount disabled the handler is intentionally null. if (ml.shouldCacheAddedEntry()) { int expectedReadCount = 0; // only use expectedReadCount if cache eviction is enabled by expected read count if (ml.getConfig().isCacheEvictionByExpectedReadCount()) { - // use the number of active cursors as the expected read count - expectedReadCount = ml.getActiveCursors().size(); + // active cursors + cache-populating random readers all read the tail entry + expectedReadCount = ml.getActiveCursors().size() + ml.getActiveCachePopulatingRandomReaderCount(); } EntryImpl entry = EntryImpl.create(ledgerId, entryId, data, expectedReadCount); entry.setDecreaseReadCountOnRelease(false); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpRandomReadEntries.java similarity index 89% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java rename to managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpRandomReadEntries.java index a0978229945cf..d4e7a06a16187 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntries.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpRandomReadEntries.java @@ -37,7 +37,7 @@ import org.apache.pulsar.common.util.FutureUtil; @CustomLog -class OpReadEntries implements ReadEntriesCallback { +class OpRandomReadEntries implements ReadEntriesCallback { private final ManagedLedgerImpl ledger; private final Position maxPosition; private final int count; @@ -46,18 +46,26 @@ class OpReadEntries implements ReadEntriesCallback { private final AtomicBoolean terminal = new AtomicBoolean(); private Position readPosition; private Position nextReadPosition; + private final boolean populateCache; - private OpReadEntries(ManagedLedgerImpl ledger, Position readPosition, int count, Position maxPosition) { + private OpRandomReadEntries(ManagedLedgerImpl ledger, Position readPosition, int count, Position maxPosition, + boolean populateCache) { this.ledger = ledger; this.readPosition = ledger.startReadOperationOnLedger(readPosition); this.count = count; this.maxPosition = maxPosition; this.nextReadPosition = this.readPosition; + this.populateCache = populateCache; } static CompletableFuture> read(ManagedLedgerImpl ledger, Position readPosition, int count, Position maxPosition) { - OpReadEntries op = new OpReadEntries(ledger, readPosition, count, maxPosition); + return read(ledger, readPosition, count, maxPosition, true); + } + + static CompletableFuture> read(ManagedLedgerImpl ledger, Position readPosition, int count, + Position maxPosition, boolean populateCache) { + OpRandomReadEntries op = new OpRandomReadEntries(ledger, readPosition, count, maxPosition, populateCache); op.readEntries(); return op.promise; } @@ -146,7 +154,13 @@ private void internalReadFromLedger(ReadHandle readHandle) { .attr("firstEntry", firstEntry) .attr("lastEntry", lastEntry) .log("Reading entries from ledger"); - ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry, this, null); + if (populateCache) { + ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry, + ledger::getActiveCachePopulatingRandomReaderCount, this); + } else { + // Cache-bypass: do not write back misses (expectedReadCount = 0). + ledger.asyncReadEntryForRandomReader(readHandle, firstEntry, lastEntry, () -> 0, this); + } } private Position getNextLedgerPosition(long ledgerId) { @@ -156,6 +170,9 @@ private Position getNextLedgerPosition(long ledgerId) { @Override public void readEntriesComplete(List returnedEntries, Object ctx) { + if (!populateCache) { + returnedEntries.forEach(entry -> ((EntryImpl) entry).setDecreaseReadCountOnRelease(false)); + } if (terminal.get()) { returnedEntries.forEach(Entry::release); return; diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReader.java similarity index 82% rename from managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java rename to managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReader.java index 1a416874146df..64cc06a91fe27 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaderImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReader.java @@ -28,27 +28,26 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerAlreadyClosedException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionFactory; -import org.apache.bookkeeper.mledger.RandomReader; import org.apache.bookkeeper.mledger.proto.ManagedLedgerInfo.LedgerInfo; import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; -final class RandomReaderImpl implements RandomReader { +public final class RandomReader implements AutoCloseable { private final ManagedLedgerImpl ledger; private final RandomReaders owner; + private final boolean populateCache; private final AtomicBoolean closed = new AtomicBoolean(); - RandomReaderImpl(ManagedLedgerImpl ledger, RandomReaders owner) { + RandomReader(ManagedLedgerImpl ledger, RandomReaders owner, boolean populateCache) { this.ledger = ledger; this.owner = owner; + this.populateCache = populateCache; } - @Override public CompletableFuture> read(Position startPosition, int numberOfEntries) { return read(startPosition, numberOfEntries, PositionFactory.LATEST, ManagedLedgerUtils.NO_MAX_SIZE_LIMIT); } - @Override public CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition, long maxSizeBytes) { if (closed.get() || owner.isClosed()) { @@ -71,7 +70,16 @@ public CompletableFuture> read(Position startPosition, int numberOfE estimateEntryCountByBytesSize(numberOfEntries, maxSizeBytes, normalizedStartPosition, ledger)); } - return OpReadEntries.read(ledger, normalizedStartPosition, effectiveCount, normalizedMaxPosition); + return OpRandomReadEntries.read( + ledger, normalizedStartPosition, effectiveCount, normalizedMaxPosition, populateCache); + } + + public CompletableFuture> read(Position startPosition, int numberOfEntries, Position maxPosition) { + return read(startPosition, numberOfEntries, maxPosition, ManagedLedgerUtils.NO_MAX_SIZE_LIMIT); + } + + public CompletableFuture> read(Position startPosition, int numberOfEntries, long maxSizeBytes) { + return read(startPosition, numberOfEntries, PositionFactory.LATEST, maxSizeBytes); } private Position normalizeStartPosition(Position startPosition) { @@ -95,7 +103,7 @@ private Position normalizeStartPosition(Position startPosition) { @Override public void close() { if (closed.compareAndSet(false, true)) { - owner.unregister(); + owner.unregister(populateCache); } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java index 0742f2d355d30..f12413ed8f275 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/RandomReaders.java @@ -18,37 +18,48 @@ */ package org.apache.bookkeeper.mledger.impl; -final class RandomReaders { +import java.util.concurrent.atomic.AtomicInteger; + +public final class RandomReaders { private final ManagedLedgerImpl ledger; - // Lifecycle mutations are synchronized; volatile reads keep add and read paths lock-free. - private volatile int activeReaders; + private final AtomicInteger activeReaders = new AtomicInteger(); + private final AtomicInteger cachePopulatingReaders = new AtomicInteger(); private volatile boolean closed; RandomReaders(ManagedLedgerImpl ledger) { this.ledger = ledger; } - synchronized RandomReaderImpl create() { + public synchronized RandomReader create(boolean populateCache) { if (closed) { throw new IllegalStateException("Managed ledger is already closed"); } - RandomReaderImpl reader = new RandomReaderImpl(ledger, this); - activeReaders++; + RandomReader reader = new RandomReader(ledger, this, populateCache); + activeReaders.incrementAndGet(); + if (populateCache) { + cachePopulatingReaders.incrementAndGet(); + } return reader; } - synchronized void unregister() { - if (activeReaders > 0) { - activeReaders--; + synchronized void unregister(boolean populateCache) { + // Guard against negative: closeAll() may have zeroed counters before a late unregister. + activeReaders.updateAndGet(current -> Math.max(0, current - 1)); + if (populateCache) { + cachePopulatingReaders.updateAndGet(current -> Math.max(0, current - 1)); } } - boolean hasActiveReaders() { - return activeReaders > 0; + boolean hasCachePopulatingReaders() { + return cachePopulatingReaders.get() > 0; + } + + int cachePopulatingCount() { + return cachePopulatingReaders.get(); } int size() { - return activeReaders; + return activeReaders.get(); } boolean isClosed() { @@ -57,6 +68,7 @@ boolean isClosed() { synchronized void closeAll() { closed = true; - activeReaders = 0; + activeReaders.set(0); + cachePopulatingReaders.set(0); } } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 79102bcf2e305..49e7ac1764cd4 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -123,6 +123,7 @@ import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntryCallback; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.EntryReadCountHandler; import org.apache.bookkeeper.mledger.LedgerOffloader; import org.apache.bookkeeper.mledger.ManagedCursor; import org.apache.bookkeeper.mledger.ManagedCursor.IndividualDeletedEntries; @@ -138,7 +139,6 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionBound; import org.apache.bookkeeper.mledger.PositionFactory; -import org.apache.bookkeeper.mledger.RandomReader; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.VoidCallback; import org.apache.bookkeeper.mledger.impl.MetaStore.MetaStoreCallback; import org.apache.bookkeeper.mledger.impl.cache.EntryCache; @@ -374,9 +374,9 @@ public void managedLedgerApi() throws Exception { @Test(timeOut = 20000) public void testRandomReader() throws Exception { - ManagedLedger ledger = factory.open("testRandomReader", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReader", new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); assertEquals(reader.read(PositionFactory.EARLIEST, 10).get(5, TimeUnit.SECONDS), Collections.emptyList()); @@ -407,8 +407,8 @@ public void testRandomReader() throws Exception { @Test(timeOut = 20000) public void testRandomReaderRejectsInvalidArguments() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderRejectsInvalidArguments"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderRejectsInvalidArguments"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); assertTrue(expectFutureFailure(reader.read(null, 1)) instanceof IllegalArgumentException); assertTrue(expectFutureFailure(reader.read(PositionFactory.EARLIEST, 0)) @@ -421,9 +421,9 @@ public void testRandomReaderRejectsInvalidArguments() throws Exception { @Test(timeOut = 20000) public void testRandomReaderPositionBoundaryValidation() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderPositionBoundaryValidation", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderPositionBoundaryValidation", new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -450,8 +450,8 @@ public void testRandomReaderPositionBoundaryValidation() throws Exception { @Test(timeOut = 20000) public void testRandomReaderDoesNotWaitForFutureWrites() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderDoesNotWaitForFutureWrites"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderDoesNotWaitForFutureWrites"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); CompletableFuture> readAfterLast = reader.read(p0.getNext(), 10); @@ -469,9 +469,9 @@ public void testRandomReaderDoesNotWaitForFutureWrites() throws Exception { @Test(timeOut = 20000) public void testRandomReaderExactCountBoundaries() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderExactCountBoundaries", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderExactCountBoundaries", new ManagedLedgerConfig().setMaxEntriesPerLedger(3)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -494,7 +494,7 @@ public void testRandomReaderSkipsEmptyLedgers() throws Exception { ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( "testRandomReaderSkipsEmptyLedgers", new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); ledger.ledgerClosed(ledger.currentLedger, 0L); @@ -517,9 +517,10 @@ public void testRandomReaderSkipsEmptyLedgers() throws Exception { @Test(timeOut = 20000) public void testRandomReaderStopsOnErrorAndReturnsPartialEntries() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderStopsOnErrorAndReturnsPartialEntries", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStopsOnErrorAndReturnsPartialEntries", new ManagedLedgerConfig().setMaxEntriesPerLedger(1)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -539,8 +540,8 @@ public void testRandomReaderStopsOnErrorAndReturnsPartialEntries() throws Except @Test(timeOut = 20000) public void testRandomReaderFailsWhenClosed() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderFailsWhenClosed"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderFailsWhenClosed"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position position = ledger.addEntry("entry-0".getBytes(Encoding)); ledger.close(); @@ -551,8 +552,8 @@ public void testRandomReaderFailsWhenClosed() throws Exception { @Test(timeOut = 20000) public void testRandomReaderIgnoresCursorAckState() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderIgnoresCursorAckState"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderIgnoresCursorAckState"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); ManagedCursor cursor = ledger.openCursor("c1"); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); @@ -570,8 +571,9 @@ public void testRandomReaderIgnoresCursorAckState() throws Exception { @Test(timeOut = 20000) public void testRandomReaderMaxPositionBeforeStartReturnsEmpty() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxPositionBeforeStartReturnsEmpty"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderMaxPositionBeforeStartReturnsEmpty"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -585,9 +587,9 @@ public void testRandomReaderMaxPositionBeforeStartReturnsEmpty() throws Exceptio @Test(timeOut = 20000) public void testRandomReaderMaxPositionCapsSameLedger() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxPositionCapsSameLedger", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMaxPositionCapsSameLedger", new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -604,9 +606,9 @@ public void testRandomReaderMaxPositionCapsSameLedger() throws Exception { @Test(timeOut = 20000) public void testRandomReaderMaxPositionInclusive() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxPositionInclusive", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMaxPositionInclusive", new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -622,9 +624,9 @@ public void testRandomReaderMaxPositionInclusive() throws Exception { @Test(timeOut = 20000) public void testRandomReaderMaxPositionCrossesLedger() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxPositionCrossesLedger", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMaxPositionCrossesLedger", new ManagedLedgerConfig().setMaxEntriesPerLedger(2)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -644,9 +646,10 @@ public void testRandomReaderMaxPositionCrossesLedger() throws Exception { @Test(timeOut = 20000) public void testRandomReaderMaxPositionRespectsBothConstraints() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxPositionRespectsBothConstraints", + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderMaxPositionRespectsBothConstraints", new ManagedLedgerConfig().setMaxEntriesPerLedger(10)); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); @@ -664,8 +667,8 @@ public void testRandomReaderMaxPositionRespectsBothConstraints() throws Exceptio @Test(timeOut = 20000) public void testRandomReaderMaxSizeBytesAndNullMaxPosition() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderMaxSizeBytesAndNullMaxPosition"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMaxSizeBytesAndNullMaxPosition"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); byte[] data = new byte[100]; Position p0 = ledger.addEntry(data); @@ -691,7 +694,7 @@ public void testRandomReaderSeedsCacheOnZeroCursorTopic() throws Exception { assertEquals(ledger.getActiveCursors().size(), 0); assertFalse(ledger.shouldCacheAddedEntry()); - RandomReader reader = ledger.newRandomReader(); + RandomReader reader = ledger.randomReaders().create(true); assertEquals(ledger.getActiveRandomReaderCount(), 1); assertTrue(ledger.shouldCacheAddedEntry()); @@ -718,7 +721,7 @@ public void testRandomReaderMissPopulatesReadThroughCache() throws Exception { ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderMissPopulatesReadThroughCache"); Position position = ledger.addEntry("entry".getBytes(Encoding)); assertEquals(ledger.getCacheSize(), 0); - @Cleanup RandomReader reader = ledger.newRandomReader(); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); long missesBefore = factory.getMbean().getCacheMissesTotal(); assertEntryPositionsAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), position); @@ -734,8 +737,8 @@ public void testRandomReaderMissPopulatesReadThroughCache() throws Exception { @Test(timeOut = 20000) public void testRandomReaderLifecycleTracksManagedLedgerCloseAndFence() throws Exception { ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderLifecycleClose"); - RandomReader first = ledger.newRandomReader(); - RandomReader second = ledger.newRandomReader(); + RandomReader first = ledger.randomReaders().create(true); + RandomReader second = ledger.randomReaders().create(true); assertEquals(ledger.getActiveRandomReaderCount(), 2); first.close(); @@ -743,27 +746,27 @@ public void testRandomReaderLifecycleTracksManagedLedgerCloseAndFence() throws E assertEquals(ledger.getActiveRandomReaderCount(), 1); ledger.close(); assertEquals(ledger.getActiveRandomReaderCount(), 0); - assertTrue(((RandomReaderImpl) second).isClosed()); + assertTrue(((RandomReader) second).isClosed()); assertTrue(expectFutureFailure(second.read(PositionFactory.EARLIEST, 1)) instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException); try { - ledger.newRandomReader(); + ledger.randomReaders().create(true); fail("Expected opening a reader on a closed ledger to fail"); } catch (IllegalStateException expected) { // expected } ManagedLedgerImpl fencedLedger = (ManagedLedgerImpl) factory.open("testRandomReaderLifecycleFence"); - RandomReader fencedReader = fencedLedger.newRandomReader(); + RandomReader fencedReader = fencedLedger.randomReaders().create(true); fencedLedger.setFenced(); assertEquals(fencedLedger.getActiveRandomReaderCount(), 0); - assertTrue(((RandomReaderImpl) fencedReader).isClosed()); + assertTrue(((RandomReader) fencedReader).isClosed()); } @Test(timeOut = 20000) public void testRandomReaderConcurrentReadsAreStateless() throws Exception { - ManagedLedger ledger = factory.open("testRandomReaderConcurrentReadsAreStateless"); - @Cleanup RandomReader reader = ledger.newRandomReader(); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderConcurrentReadsAreStateless"); + @Cleanup RandomReader reader = ledger.randomReaders().create(true); Position p0 = ledger.addEntry("entry-0".getBytes(Encoding)); Position p1 = ledger.addEntry("entry-1".getBytes(Encoding)); Position p2 = ledger.addEntry("entry-2".getBytes(Encoding)); @@ -792,12 +795,240 @@ public void testOpReadEntriesSkipsEmptyLedger() throws Exception { ledger.ledgers.put(emptyLedger.getId(), new LedgerInfo() .setLedgerId(emptyLedger.getId()).setEntries(0).setSize(0)); - CompletableFuture> promise = OpReadEntries.read(ledger, + CompletableFuture> promise = OpRandomReadEntries.read(ledger, PositionFactory.create(emptyLedger.getId(), 0), 1, PositionFactory.LATEST); assertEntryPositionsAndRelease(promise.get(5, TimeUnit.SECONDS), nextPosition); ledger.close(); } + @Test(timeOut = 20000) + public void testRandomReaderStreamingFalseDoesNotSeedCache() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testRandomReaderStreamingFalseDoesNotSeedCache"); + assertEquals(ledger.getActiveRandomReaderCount(), 0); + assertFalse(ledger.shouldCacheAddedEntry()); + try (RandomReader reader = ledger.randomReaders().create(false)) { + assertEquals(ledger.getActiveRandomReaderCount(), 1); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 0); + assertFalse(ledger.shouldCacheAddedEntry()); + ledger.addEntry("entry-0".getBytes(Encoding)); + assertEquals(ledger.getCacheSize(), 0); + } + assertEquals(ledger.getActiveRandomReaderCount(), 0); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStreamingFalseDoesNotPopulateCacheOnMiss() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingFalseDoesNotPopulateCacheOnMiss"); + Position position = ledger.addEntry("entry".getBytes(Encoding)); + assertEquals(ledger.getCacheSize(), 0); + try (RandomReader reader = ledger.randomReaders().create(false)) { + long missesBefore = factory.getMbean().getCacheMissesTotal(); + assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore + 1); + assertEquals(ledger.getCacheSize(), 0); + long missesBefore2 = factory.getMbean().getCacheMissesTotal(); + assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore2 + 1); + } + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStreamingTruePopulatesCacheOnMiss() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingTruePopulatesCacheOnMiss"); + Position position = ledger.addEntry("entry".getBytes(Encoding)); + assertEquals(ledger.getCacheSize(), 0); + try (RandomReader reader = ledger.randomReaders().create(true)) { + long missesBefore = factory.getMbean().getCacheMissesTotal(); + assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore + 1); + assertTrue(ledger.getCacheSize() > 0); + long hitsBefore = factory.getMbean().getCacheHitsTotal(); + assertEntryDataAndRelease(reader.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + assertEquals(factory.getMbean().getCacheHitsTotal(), hitsBefore + 1); + } + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStreamingModesCoexistForSeeding() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingModesCoexistForSeeding"); + RandomReader bypass = ledger.randomReaders().create(false); + assertFalse(ledger.shouldCacheAddedEntry()); + RandomReader streaming = ledger.randomReaders().create(true); + assertTrue(ledger.shouldCacheAddedEntry()); + streaming.close(); + assertFalse(ledger.shouldCacheAddedEntry()); + bypass.close(); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testCachePopulatingReadersAreCounted() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testCachePopulatingReadersAreCounted"); + try (RandomReader first = ledger.randomReaders().create(true); + RandomReader second = ledger.randomReaders().create(true)) { + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 2); + } + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 0); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStreamingFalseServesExistingCacheHit() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingFalseServesExistingCacheHit"); + Position position = ledger.addEntry("entry".getBytes(Encoding)); + try (RandomReader seeder = ledger.randomReaders().create(true)) { + assertEntryDataAndRelease(seeder.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + } + assertTrue(ledger.getCacheSize() > 0); + try (RandomReader bypass = ledger.randomReaders().create(false)) { + long missesBefore = factory.getMbean().getCacheMissesTotal(); + assertEntryDataAndRelease(bypass.read(position, 1).get(5, TimeUnit.SECONDS), "entry"); + assertEquals(factory.getMbean().getCacheMissesTotal(), missesBefore); + } + ledger.close(); + } + + @Test(timeOut = 20000) + public void testStreamingReaderCountAccessorWiring() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testStreamingReaderCountAccessorWiring"); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 0); + RandomReader s1 = ledger.randomReaders().create(true); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 1); + RandomReader b1 = ledger.randomReaders().create(false); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 1); + RandomReader s2 = ledger.randomReaders().create(true); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 2); + s1.close(); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 1); + s2.close(); + b1.close(); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 0); + ledger.close(); + } + + @Test(timeOut = 20000) + public void testFenceZeroesStreamingReaders() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("testFenceZeroesStreamingReaders"); + RandomReader streaming = ledger.randomReaders().create(true); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 1); + ledger.setFenced(); + assertEquals(ledger.getActiveCachePopulatingRandomReaderCount(), 0); + assertTrue(((RandomReader) streaming).isClosed()); + // Do not close(): a fenced ledger cannot be closed (throws ManagedLedgerFencedException); the factory + // tears it down. Mirrors testRandomReaderLifecycleTracksManagedLedgerCloseAndFence. + } + + // These tests assert the expectedReadCount value entries carry, not eviction under pressure. + + @Test(timeOut = 20000) + public void testRandomReaderStreamingSeedsTailWithEvictionWeight() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setCacheEvictionByExpectedReadCount(true); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingSeedsTailWithEvictionWeight", config); + try (RandomReader reader = ledger.randomReaders().create(true)) { + Position position = ledger.addEntry("entry-0".getBytes(Encoding)); + List entries = reader.read(position, 1).get(5, TimeUnit.SECONDS); + try { + assertEquals(entries.size(), 1); + EntryReadCountHandler handler = entries.get(0).getReadCountHandler(); + assertNotNull(handler); + // write-path seed: cursors(0) + streaming(1) = 1. Assert BEFORE release (shared handler). + assertEquals(handler.getExpectedReadCount(), 1); + } finally { + entries.forEach(Entry::release); + } + } + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderMultiStreamingFanoutSharedHandler() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setCacheEvictionByExpectedReadCount(true); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderMultiStreamingFanoutSharedHandler", config); + try (RandomReader a = ledger.randomReaders().create(true); + RandomReader b = ledger.randomReaders().create(true)) { + Position position = ledger.addEntry("entry-0".getBytes(Encoding)); + // Seeded weight is cursors(0)+streaming(2)=2; the shared handler decrements per release (A:2, B:1). + Entry entryA = a.read(position, 1).get(5, TimeUnit.SECONDS).get(0); + try { + assertEquals(entryA.getReadCountHandler().getExpectedReadCount(), 2); + } finally { + entryA.release(); // shared handler: 2 -> 1 + } + Entry entryB = b.read(position, 1).get(5, TimeUnit.SECONDS).get(0); + try { + assertEquals(entryB.getReadCountHandler().getExpectedReadCount(), 1); + } finally { + entryB.release(); + } + } + ledger.close(); + } + + @Test(timeOut = 20000) + public void testRandomReaderStreamingFalseWithCursorsContributesNothingToWeight() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setCacheEvictionByExpectedReadCount(true); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderStreamingFalseWithCursorsContributesNothingToWeight", config); + ManagedCursor cursor = ledger.openCursor("c1"); + try (RandomReader streaming = ledger.randomReaders().create(true); + RandomReader bypass = ledger.randomReaders().create(false)) { + // weight = cursors(1) + streaming(1) = 2; the streaming=false (bypass) reader contributes nothing. + // A regression that counted bypass readers would yield 3. + Position position = ledger.addEntry("entry-0".getBytes(Encoding)); + Entry entry = streaming.read(position, 1).get(5, TimeUnit.SECONDS).get(0); + EntryReadCountHandler handler = entry.getReadCountHandler(); + try { + assertEquals(handler.getExpectedReadCount(), 2); + } finally { + entry.release(); + } + assertEquals(handler.getExpectedReadCount(), 1); + + Entry bypassEntry = bypass.read(position, 1).get(5, TimeUnit.SECONDS).get(0); + try { + assertSame(bypassEntry.getReadCountHandler(), handler); + } finally { + bypassEntry.release(); + } + assertEquals(handler.getExpectedReadCount(), 1); + + cursor.readEntries(1).forEach(Entry::release); + assertEquals(handler.getExpectedReadCount(), 0); + } + ledger.close(); + } + + // With cacheEvictionByExpectedReadCount=false the add-path expectedReadCount stays 0, so the handler is null. + @Test(timeOut = 20000) + public void testRandomReaderEvictionWeightDisabledWhenFlagFalse() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setCacheEvictionByExpectedReadCount(false); + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open( + "testRandomReaderEvictionWeightDisabledWhenFlagFalse", config); + try (RandomReader reader = ledger.randomReaders().create(true)) { + Position position = ledger.addEntry("entry-0".getBytes(Encoding)); + Entry entry = reader.read(position, 1).get(5, TimeUnit.SECONDS).get(0); + try { + assertNull(entry.getReadCountHandler()); + } finally { + entry.release(); + } + } + ledger.close(); + } + @Test(timeOut = 20000) public void simple() throws Exception { ManagedLedger ledger = factory.open("my_test_ledger");