From c6132f4de8fced6c8db66a154cc061fa529c59f6 Mon Sep 17 00:00:00 2001 From: coderzc Date: Tue, 30 Aug 2022 12:16:48 +0800 Subject: [PATCH 1/4] Filter out and skip unnecessary entry --- .../bookkeeper/mledger/ManagedCursor.java | 17 ++-- .../bookkeeper/mledger/ReadOnlyCursor.java | 4 +- .../mledger/impl/ManagedCursorImpl.java | 23 ++--- .../mledger/impl/ManagedLedgerImpl.java | 31 +++++++ .../bookkeeper/mledger/impl/OpReadEntry.java | 27 ++++-- .../bookkeeper/mledger/impl/OpScan.java | 2 +- .../impl/ManagedCursorContainerTest.java | 10 +-- .../mledger/impl/ManagedCursorTest.java | 86 ++++++++++++++++--- .../mledger/impl/ManagedLedgerTest.java | 6 +- .../persistent/MessageDeduplication.java | 2 +- ...PersistentDispatcherMultipleConsumers.java | 17 +++- ...sistentDispatcherSingleActiveConsumer.java | 2 +- .../persistent/PersistentReplicator.java | 2 +- .../buffer/impl/TopicTransactionBuffer.java | 2 +- .../pendingack/impl/MLPendingAckStore.java | 3 +- .../pulsar/compaction/CompactedTopicImpl.java | 5 +- .../broker/delayed/MockManagedCursor.java | 10 ++- ...ckyKeyDispatcherMultipleConsumersTest.java | 4 +- .../broker/transaction/TransactionTest.java | 26 +++--- .../pulsar/sql/presto/PulsarRecordCursor.java | 2 +- .../sql/presto/TestPulsarConnector.java | 2 +- .../sql/presto/TestPulsarRecordCursor.java | 56 ++++++------ .../impl/MLTransactionLogImpl.java | 3 +- 23 files changed, 235 insertions(+), 107 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java index 4167c040ed61c..1a49608740d48 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java @@ -150,9 +150,11 @@ enum IndividualDeletedEntries { * opaque context * @param maxPosition * max position can read + * @param skipCondition + * predicate of read filter out */ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition); + PositionImpl maxPosition, Predicate skipCondition); /** @@ -163,9 +165,10 @@ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, O * @param callback callback object * @param ctx opaque context * @param maxPosition max position can read + * @param skipCondition predicate of read filter out */ void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition); + Object ctx, PositionImpl maxPosition, Predicate skipCondition); /** * Get 'N'th entry from the mark delete position in the cursor without updating any cursor positions. @@ -239,9 +242,11 @@ List readEntriesOrWait(int maxEntries, long maxSizeBytes) * opaque context * @param maxPosition * max position can read + * @param skipCondition + * predicate of read filter out */ void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition); + PositionImpl maxPosition, Predicate skipCondition); /** * Asynchronously read entries from the ManagedLedger, up to the specified number and size. @@ -260,14 +265,16 @@ void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callb * opaque context * @param maxPosition * max position can read + * @param skipCondition + * predicate of read filter out */ void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition); + PositionImpl maxPosition, Predicate skipCondition); /** * Cancel a previously scheduled asyncReadEntriesOrWait operation. * - * @see #asyncReadEntriesOrWait(int, ReadEntriesCallback, Object, PositionImpl) + * @see #asyncReadEntriesOrWait(int, ReadEntriesCallback, Object, PositionImpl, Predicate) * @return true if the read operation was canceled or false if there was no pending operation */ boolean cancelPendingReadRequest(); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java index 18d412f893152..5d67d8a5885cc 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java @@ -48,7 +48,7 @@ public interface ReadOnlyCursor { * @see #readEntries(int) */ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition); + Object ctx, PositionImpl maxPosition, Predicate skipCondition); /** * Asynchronously read entries from the ManagedLedger. @@ -60,7 +60,7 @@ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, * @param maxPosition max position can read */ void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition); + Object ctx, PositionImpl maxPosition, Predicate skipCondition); /** * Get the read position. This points to the next message to be read from the cursor. diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 3af2e568d86f3..ba87eaff23af8 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -749,7 +749,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter.await(); @@ -762,13 +762,13 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { @Override public void asyncReadEntries(final int numberOfEntriesToRead, final ReadEntriesCallback callback, - final Object ctx, PositionImpl maxPosition) { - asyncReadEntries(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition); + final Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + asyncReadEntries(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, skipCondition); } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition) { + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { checkArgument(numberOfEntriesToRead > 0); if (isClosed()) { callback.readEntriesFailed(new ManagedLedgerException @@ -779,7 +779,8 @@ public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadE int numOfEntriesToRead = applyMaxSizeCap(numberOfEntriesToRead, maxSizeBytes); PENDING_READ_OPS_UPDATER.incrementAndGet(this); - OpReadEntry op = OpReadEntry.create(this, readPosition, numOfEntriesToRead, callback, ctx, maxPosition); + OpReadEntry op = + OpReadEntry.create(this, readPosition, numOfEntriesToRead, callback, ctx, maxPosition, skipCondition); ledger.asyncReadEntries(op); } @@ -881,7 +882,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter.await(); @@ -894,13 +895,13 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition) { - asyncReadEntriesOrWait(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition); + PositionImpl maxPosition, Predicate skipCondition) { + asyncReadEntriesOrWait(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, skipCondition); } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition) { + PositionImpl maxPosition, Predicate skipCondition) { checkArgument(maxEntries > 0); if (isClosed()) { callback.readEntriesFailed(new CursorAlreadyClosedException("Cursor was already closed"), ctx); @@ -914,10 +915,10 @@ public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntrie if (log.isDebugEnabled()) { log.debug("[{}] [{}] Read entries immediately", ledger.getName(), name); } - asyncReadEntries(numberOfEntriesToRead, callback, ctx, maxPosition); + asyncReadEntries(numberOfEntriesToRead, callback, ctx, maxPosition, skipCondition); } else { OpReadEntry op = OpReadEntry.create(this, readPosition, numberOfEntriesToRead, callback, - ctx, maxPosition); + ctx, maxPosition, skipCondition); if (!WAITING_READ_OP_UPDATER.compareAndSet(this, null, op)) { op.recycle(); 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 c849a347c7d8e..bdece03aa2432 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 @@ -2056,6 +2056,37 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) long lastEntry = min(firstEntry + opReadEntry.getNumberOfEntriesToRead() - 1, lastEntryInLedger); + // Filer out and skip unnecessary read entry + if (opReadEntry.skipCondition != null) { + long firstValidEntry = -1L; + long lastValidEntry = -1L; + long entryId = firstEntry; + for (; entryId <= lastEntry; entryId++) { + if (!opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { + firstValidEntry = entryId; + break; + } + } + + // If all messages in [firstEntry...lastEntry] are filter out, + // then manual call internalReadEntriesComplete to advance read position. + if (firstValidEntry == -1L) { + opReadEntry.internalReadEntriesComplete(Collections.emptyList(), opReadEntry.ctx, + PositionImpl.get(ledger.getId(), lastEntry)); + return; + } + + for (; entryId <= lastEntry; entryId++) { + if (opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { + break; + } + lastValidEntry = entryId; + } + + firstEntry = firstValidEntry; + lastEntry = lastValidEntry; + } + if (log.isDebugEnabled()) { log.debug("[{}] Reading entries from ledger {} - first={} last={}", name, ledger.getId(), firstEntry, lastEntry); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java index 86290d8bc8f5f..eff4101368ffe 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java @@ -22,7 +22,9 @@ import io.netty.util.Recycler; import io.netty.util.Recycler.Handle; import java.util.ArrayList; +import java.util.Collections; import java.util.List; +import java.util.function.Predicate; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedgerException; @@ -45,8 +47,10 @@ class OpReadEntry implements ReadEntriesCallback { private PositionImpl nextReadPosition; PositionImpl maxPosition; + Predicate skipCondition; + public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPositionRef, int count, - ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { + ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition, Predicate skipCondition) { OpReadEntry op = RECYCLER.get(); op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef); op.cursor = cursor; @@ -57,13 +61,13 @@ public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPosi maxPosition = PositionImpl.LATEST; } op.maxPosition = maxPosition; + op.skipCondition = skipCondition; op.ctx = ctx; op.nextReadPosition = PositionImpl.get(op.readPosition); return op; } - @Override - public void readEntriesComplete(List returnedEntries, Object ctx) { + void internalReadEntriesComplete(List returnedEntries, Object ctx, PositionImpl lastPosition) { // Filter the returned entries for individual deleted messages int entriesCount = returnedEntries.size(); long entriesSize = 0; @@ -72,13 +76,19 @@ public void readEntriesComplete(List returnedEntries, Object ctx) { } cursor.updateReadStats(entriesCount, entriesSize); - final PositionImpl lastPosition = (PositionImpl) returnedEntries.get(entriesCount - 1).getPosition(); + if (lastPosition == null || entriesCount != 0) { + lastPosition = (PositionImpl) returnedEntries.get(entriesCount - 1).getPosition(); + } if (log.isDebugEnabled()) { log.debug("[{}][{}] Read entries succeeded batch_size={} cumulative_size={} requested_count={}", cursor.ledger.getName(), cursor.getName(), returnedEntries.size(), entries.size(), count); } - List filteredEntries = cursor.filterReadEntries(returnedEntries); - entries.addAll(filteredEntries); + + List filteredEntries = Collections.emptyList(); + if (entriesCount != 0) { + filteredEntries = cursor.filterReadEntries(returnedEntries); + entries.addAll(filteredEntries); + } // if entries have been filtered out then try to skip reading of already deletedMessages in that range final Position nexReadPosition = entriesCount != filteredEntries.size() @@ -87,6 +97,11 @@ public void readEntriesComplete(List returnedEntries, Object ctx) { checkReadCompletion(); } + @Override + public void readEntriesComplete(List returnedEntries, Object ctx) { + internalReadEntriesComplete(returnedEntries, ctx, null); + } + @Override public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { cursor.readOperationCompleted(); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpScan.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpScan.java index 9f1cd310f92e8..6d68b042a7ad6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpScan.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpScan.java @@ -128,7 +128,7 @@ public void find() { } if (cursor.hasMoreEntries(searchPosition)) { OpReadEntry opReadEntry = OpReadEntry.create(cursor, searchPosition, batchSize, - this, OpScan.this.ctx, null); + this, OpScan.this.ctx, null, null); ledger.asyncReadEntries(opReadEntry); } else { callback.scanComplete(lastSeenPosition, ScanOutcome.COMPLETED, OpScan.this.ctx); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java index 54ebc4c293a94..0b55f71f75d06 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java @@ -33,8 +33,8 @@ import java.util.List; import java.util.Map; import java.util.Set; -import java.util.function.Predicate; import java.util.concurrent.CompletableFuture; +import java.util.function.Predicate; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.AsyncCallbacks.ClearBacklogCallback; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteCallback; @@ -111,13 +111,13 @@ public List readEntries(int numberOfEntriesToRead) throws ManagedLedgerEx @Override public void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition) { + PositionImpl maxPosition, Predicate skipCondition) { callback.readEntriesComplete(null, ctx); } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition) { + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { callback.readEntriesComplete(null, ctx); } @@ -302,12 +302,12 @@ public List readEntriesOrWait(int numberOfEntriesToRead) @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition) { + PositionImpl maxPosition, Predicate skipCondition) { } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition) { + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index fa3d6ec6c9035..c6a8c2097aaa7 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -55,6 +55,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.Future; @@ -453,7 +454,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter.await(); } @@ -481,7 +482,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail("async-call should not have failed"); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter.await(); @@ -503,7 +504,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter2.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter2.await(); } @@ -530,7 +531,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); counter.await(); } @@ -567,7 +568,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, null, null); + }, null, null, null); counter.await(); } @@ -1826,7 +1827,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { log.error("Error reading", exception); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, position -> false); } ledger.addEntry("test".getBytes()); @@ -2830,7 +2831,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); assertTrue(c1.cancelPendingReadRequest()); @@ -2846,7 +2847,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter2.countDown(); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); ledger.addEntry("entry-1".getBytes(Encoding)); @@ -3753,7 +3754,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { completableFuture.completeExceptionally(exception); } - }, null, (PositionImpl) position); + }, null, (PositionImpl) position, null); int number = completableFuture.get(); assertEquals(number, readMaxNumber); @@ -3768,7 +3769,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { completableFuture.completeExceptionally(exception); } - }, null, (PositionImpl) maxCanReadPosition); + }, null, (PositionImpl) maxCanReadPosition, null); assertEquals(number, sendNumber - readMaxNumber); @@ -4126,7 +4127,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { // op readPosition is bigger than maxReadPosition OpReadEntry opReadEntry = OpReadEntry.create(cursor, ledger.lastConfirmedEntry, 10, callback, - null, PositionImpl.get(lastPosition.getLedgerId(), -1)); + null, PositionImpl.get(lastPosition.getLedgerId(), -1), null); Field field = ManagedCursorImpl.class.getDeclaredField("readPosition"); field.setAccessible(true); field.set(cursor, PositionImpl.EARLIEST); @@ -4148,7 +4149,7 @@ public void testOpReadEntryRecycle() throws Exception { }; @Cleanup final MockedStatic mockedStaticOpReadEntry = Mockito.mockStatic(OpReadEntry.class); - mockedStaticOpReadEntry.when(() -> OpReadEntry.create(any(), any(), anyInt(), any(), any(), any())) + mockedStaticOpReadEntry.when(() -> OpReadEntry.create(any(), any(), anyInt(), any(), any(), any(), any())) .thenAnswer(__ -> createOpReadEntry.get()); final ManagedLedgerConfig ledgerConfig = new ManagedLedgerConfig(); @@ -4171,7 +4172,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { final int numReadRequests = 3; for (int i = 0; i < numReadRequests; i++) { - cursor.asyncReadEntriesOrWait(1, callback, null, new PositionImpl(0, 0)); + cursor.asyncReadEntriesOrWait(1, callback, null, new PositionImpl(0, 0), null); } Awaitility.await().atMost(Duration.ofSeconds(1)) .untilAsserted(() -> assertEquals(ledger.waitingCursors.size(), 1)); @@ -4252,5 +4253,64 @@ public void testLazyCursorLedgerCreationForSubscriptionCreation() throws Excepti factory2.shutdown(); } + @Test + public void testReadEntriesWithFilterOut() throws ManagedLedgerException, InterruptedException, ExecutionException { + int readMaxNumber = 10; + int sendNumber = 20; + ManagedLedger ledger = factory.open("testReadEntriesWithFilter"); + ManagedCursor cursor = ledger.openCursor("c"); + Position position = PositionImpl.EARLIEST; + Position maxCanReadPosition = PositionImpl.EARLIEST; + for (int i = 0; i < sendNumber; i++) { + if (i == readMaxNumber - 1) { + position = ledger.addEntry(new byte[1024]); + } else if (i == sendNumber - 1) { + maxCanReadPosition = ledger.addEntry(new byte[1024]); + } else { + ledger.addEntry(new byte[1024]); + } + + } + CompletableFuture completableFuture = new CompletableFuture<>(); + cursor.asyncReadEntriesOrWait(sendNumber, new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + completableFuture.complete(entries.size()); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + completableFuture.completeExceptionally(exception); + } + }, null, (PositionImpl) position, pos -> { + return pos.getEntryId() % 2 != 0; + }); + + int number = completableFuture.get(); + assertEquals(number, readMaxNumber / 2); + + assertEquals(cursor.getReadPosition().getEntryId(), 10); + + CompletableFuture completableFuture2 = new CompletableFuture<>(); + cursor.asyncReadEntriesOrWait(sendNumber, new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + completableFuture2.complete(entries.size()); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + completableFuture2.completeExceptionally(exception); + } + }, null, (PositionImpl) maxCanReadPosition, pos -> { + return pos.getEntryId() % 2 != 0; + }); + + int number2 = completableFuture2.get(); + assertEquals(number2, readMaxNumber / 2); + + assertEquals(cursor.getReadPosition().getEntryId(), 20); + } + private static final Logger log = LoggerFactory.getLogger(ManagedCursorTest.class); } 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 29e30a958bd02..d85370bf92bf5 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 @@ -470,7 +470,7 @@ public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, cursor, PositionImpl.LATEST); + }, cursor, PositionImpl.LATEST, null); } @Override @@ -555,7 +555,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { } - }, null, maxPosition); + }, null, maxPosition, null); Assert.assertEquals(opReadEntry.readPosition, position); } @@ -3030,7 +3030,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { responseException2.set(exception); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); ledger.asyncReadEntry(ledgerHandle, PositionImpl.EARLIEST.getEntryId(), PositionImpl.EARLIEST.getEntryId(), opReadEntry, ctxStr); retryStrategically((test) -> { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 030d74a014f29..2be52903f36c3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -193,7 +193,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { future.completeExceptionally(exception); } - }, null, PositionImpl.LATEST); + }, null, PositionImpl.LATEST, null); } public Status getStatus() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 9df30095131c3..1f36bba528780 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -37,6 +37,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; +import java.util.function.Predicate; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -48,6 +49,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.delayed.DelayedDeliveryTracker; import org.apache.pulsar.broker.delayed.InMemoryDelayedDeliveryTracker; +import org.apache.pulsar.broker.delayed.bucket.BucketDelayedDeliveryTracker; import org.apache.pulsar.broker.service.AbstractDispatcherMultipleConsumers; import org.apache.pulsar.broker.service.BrokerServiceException; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerBusyException; @@ -318,8 +320,17 @@ public synchronized void readMoreEntries() { minReplayedPosition = null; } - cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, - ReadType.Normal, topic.getMaxReadPosition()); + Predicate skipCondition = null; + // Filter out and skip read delayed messages exist in DelayedDeliveryTracker + if (delayedDeliveryTracker.isPresent()) { + final DelayedDeliveryTracker deliveryTracker = delayedDeliveryTracker.get(); + if (deliveryTracker instanceof BucketDelayedDeliveryTracker) { + skipCondition = position -> deliveryTracker + .containsMessage(position.getLedgerId(), position.getEntryId()); + } + } + cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, + topic.getMaxReadPosition(), skipCondition); } else { log.debug("[{}] Cannot schedule next read until previous one is done", name); } @@ -1045,7 +1056,7 @@ public boolean trackDelayedDelivery(long ledgerId, long entryId, MessageMetadata } } - protected synchronized NavigableSet getMessagesToReplayNow(int maxMessagesToRead) { + protected synchronized NavigableSet getMessagesToReplayNow(int maxMessagesToRead) { if (!redeliveryMessages.isEmpty()) { return redeliveryMessages.getMessagesToReplayNow(maxMessagesToRead); } else if (delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().hasMessageAvailable()) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index 03dce58af3a92..d41d2a831c6ac 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -360,7 +360,7 @@ protected void readMoreEntries(Consumer consumer) { ReadEntriesCtx readEntriesCtx = ReadEntriesCtx.create(consumer, consumer.getConsumerEpoch()); cursor.asyncReadEntriesOrWait(messagesToRead, - bytesToRead, this, readEntriesCtx, topic.getMaxReadPosition()); + bytesToRead, this, readEntriesCtx, topic.getMaxReadPosition(), null); } } } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index 2d341016f089e..d9b8332462461 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -249,7 +249,7 @@ protected void readMoreEntries() { log.debug("[{}] Schedule read of {} messages", replicatorId, messagesToRead); } cursor.asyncReadEntriesOrWait(messagesToRead, readMaxSizeBytes, this, - null, topic.getMaxReadPosition()); + null, topic.getMaxReadPosition(), null); } else { if (log.isDebugEnabled()) { log.debug("[{}] Not scheduling read due to pending read. Messages To Read {}", diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java index f3bf4f95923cd..4131c03ec5e04 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java @@ -677,7 +677,7 @@ boolean fillQueue() { if (cursor.hasMoreEntries()) { outstandingReadsRequests.incrementAndGet(); cursor.asyncReadEntries(NUMBER_OF_PER_READ_ENTRY, - this, System.nanoTime(), PositionImpl.LATEST); + this, System.nanoTime(), PositionImpl.LATEST, null); } else { if (entryQueue.size() == 0) { isReadable = false; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java index 4dce8b9a0fcf4..3825e5775844c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java @@ -147,7 +147,8 @@ public void replayAsync(PendingAckHandleImpl pendingAckHandle, ExecutorService t //TODO can control the number of entry to read private void readAsync(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback readEntriesCallback) { - cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST); + cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST, + null); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java index c8114f9adb652..f008ae454df29 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java @@ -102,7 +102,8 @@ public void asyncReadEntriesOrWait(ManagedCursor cursor, ReadEntriesCtx readEntriesCtx = ReadEntriesCtx.create(consumer, DEFAULT_CONSUMER_EPOCH); if (compactionHorizon == null || compactionHorizon.compareTo(cursorPosition) < 0) { - cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, PositionImpl.LATEST); + cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, PositionImpl.LATEST, + null); } else { compactedTopicContext.thenCompose( (context) -> findStartPoint(cursorPosition, context.ledger.getLastAddConfirmed(), context.cache) @@ -116,7 +117,7 @@ public void asyncReadEntriesOrWait(ManagedCursor cursor, } if (startPoint == NEWER_THAN_COMPACTED && compactionHorizon.compareTo(cursorPosition) < 0) { cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, - PositionImpl.LATEST); + PositionImpl.LATEST, null); return CompletableFuture.completedFuture(null); } else { long endPoint = Math.min(context.ledger.getLastAddConfirmed(), diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java index efb0fa7ab7ba2..9e3697174ddec 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java @@ -24,6 +24,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Predicate; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -102,13 +103,14 @@ public List readEntries(int numberOfEntriesToRead) throws InterruptedExce @Override public void asyncReadEntries(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition) { + PositionImpl maxPosition, Predicate skipCondition) { } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, - AsyncCallbacks.ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { + AsyncCallbacks.ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { } @@ -138,13 +140,13 @@ public List readEntriesOrWait(int maxEntries, long maxSizeBytes) @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition) { + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, AsyncCallbacks.ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition) { + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java index 1c23afd957a64..eba18936eaab6 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java @@ -252,7 +252,7 @@ public void testSkipRedeliverTemporally() { return null; }).when(cursorMock).asyncReadEntriesOrWait( anyInt(), anyLong(), any(PersistentStickyKeyDispatcherMultipleConsumers.class), - eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any()); + eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any(), any()); } catch (Exception e) { fail("Failed to set to field", e); } @@ -421,7 +421,7 @@ public void testMessageRedelivery() throws Exception { return null; }).when(cursorMock).asyncReadEntriesOrWait(anyInt(), anyLong(), any(PersistentStickyKeyDispatcherMultipleConsumers.class), - eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any()); + eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any(), any()); // (1) Run sendMessagesToConsumers // (2) Attempts to send message1 to consumer1 but skipped because availablePermits is 0 diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 0fcf9ff9b33d7..50cf0b3168684 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -42,12 +42,12 @@ import io.netty.util.concurrent.DefaultThreadFactory; import java.lang.reflect.Field; import java.lang.reflect.Method; -import java.util.ArrayList; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; -import java.util.Map; import java.util.List; +import java.util.Map; import java.util.Optional; import java.util.UUID; import java.util.concurrent.CompletableFuture; @@ -121,9 +121,9 @@ import org.apache.pulsar.client.impl.ConsumerBase; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.MessagesImpl; +import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.common.api.proto.CommandSubscribe; -import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.events.EventType; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.SystemTopicNames; @@ -141,8 +141,8 @@ import org.apache.pulsar.transaction.coordinator.TransactionRecoverTracker; import org.apache.pulsar.transaction.coordinator.TransactionTimeoutTracker; import org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl; -import org.apache.pulsar.transaction.coordinator.impl.MLTransactionSequenceIdGenerator; import org.apache.pulsar.transaction.coordinator.impl.MLTransactionMetadataStore; +import org.apache.pulsar.transaction.coordinator.impl.MLTransactionSequenceIdGenerator; import org.apache.pulsar.transaction.coordinator.impl.TxnLogBufferedWriterConfig; import org.awaitility.Awaitility; import org.mockito.invocation.InvocationOnMock; @@ -701,7 +701,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); Class managedLedgerClass = ManagedLedgerImpl.class; Field field = managedLedgerClass.getDeclaredField("cursors"); field.setAccessible(true); @@ -713,7 +713,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); TransactionBuffer buffer2 = new TopicTransactionBuffer(persistentTopic); Awaitility.await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> @@ -724,7 +724,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); managedCursors.add(managedCursor, managedCursor.getMarkDeletedPosition()); TransactionBuffer buffer3 = new TopicTransactionBuffer(persistentTopic); @@ -769,7 +769,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); TransactionPendingAckStoreProvider pendingAckStoreProvider = mock(TransactionPendingAckStoreProvider.class); doReturn(CompletableFuture.completedFuture( @@ -794,7 +794,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); PendingAckHandleImpl pendingAckHandle2 = new PendingAckHandleImpl(persistentSubscription); Awaitility.await().untilAsserted(() -> @@ -804,7 +804,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); PendingAckHandleImpl pendingAckHandle3 = new PendingAckHandleImpl(persistentSubscription); @@ -838,7 +838,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); MLTransactionSequenceIdGenerator mlTransactionSequenceIdGenerator = new MLTransactionSequenceIdGenerator(); persistentTopic.getManagedLedger().getConfig().setManagedLedgerInterceptor(mlTransactionSequenceIdGenerator); MLTransactionLogImpl mlTransactionLog = @@ -869,7 +869,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); MLTransactionMetadataStore metadataStore2 = new MLTransactionMetadataStore(new TransactionCoordinatorID(1), @@ -883,7 +883,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); MLTransactionMetadataStore metadataStore3 = new MLTransactionMetadataStore(new TransactionCoordinatorID(1), diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java index f86335ae3780b..df3927e6c2089 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java @@ -410,7 +410,7 @@ public void run() { // if the available size is invalid and the entry queue size is 0, read one entry outstandingReadsRequests.decrementAndGet(); cursor.asyncReadEntries(batchSize, entryQueueCacheSizeAllocator.getAvailableCacheSize(), - this, System.nanoTime(), PositionImpl.LATEST); + this, System.nanoTime(), PositionImpl.LATEST, null); } // stats for successful read request diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java index 61f6edd38530c..a54d9511b49a1 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java @@ -650,7 +650,7 @@ public void run() { return null; } - }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any()); + }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any(), any()); when(readOnlyCursor.hasMoreEntries()).thenAnswer(new Answer() { @Override diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java index 7eaa2da498f45..c5e089d81b0df 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java @@ -18,6 +18,22 @@ */ package org.apache.pulsar.sql.presto; +import static java.util.concurrent.CompletableFuture.completedFuture; +import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertTrue; +import static org.testng.Assert.fail; import com.fasterxml.jackson.databind.ObjectMapper; import io.airlift.log.Logger; import io.netty.buffer.ByteBuf; @@ -26,7 +42,17 @@ import io.trino.spi.type.RowType; import io.trino.spi.type.Type; import io.trino.spi.type.VarcharType; +import java.lang.reflect.Field; +import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Method; import java.math.BigDecimal; +import java.nio.charset.Charset; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; import lombok.Data; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; @@ -57,34 +83,6 @@ import org.mockito.stubbing.Answer; import org.testng.annotations.Test; -import java.lang.reflect.Field; -import java.lang.reflect.InvocationTargetException; -import java.lang.reflect.Method; -import java.nio.charset.Charset; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.HashMap; -import java.util.LinkedList; -import java.util.List; -import java.util.Map; - -import static java.util.concurrent.CompletableFuture.completedFuture; -import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyInt; -import static org.mockito.ArgumentMatchers.anyLong; -import static org.mockito.ArgumentMatchers.anyString; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.doAnswer; -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.when; -import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertNotNull; -import static org.testng.Assert.assertTrue; -import static org.testng.Assert.fail; - public class TestPulsarRecordCursor extends TestPulsarConnector { private static final Logger log = Logger.get(TestPulsarRecordCursor.class); @@ -385,7 +383,7 @@ public void run() { return null; } - }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any()); + }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any(), any()); when(readOnlyCursor.hasMoreEntries()).thenAnswer(new Answer() { @Override diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java index f2e1f60663d28..1485d6aa15e72 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java @@ -156,7 +156,8 @@ public void replayAsync(TransactionLogReplayCallback transactionLogReplayCallbac private void readAsync(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback readEntriesCallback) { - cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST); + cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST, + null); } @Override From e2efb0dbfe7b903244741d36064131d4b9cd256d Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 28 Dec 2022 15:26:19 +0800 Subject: [PATCH 2/4] add a new method --- .../bookkeeper/mledger/ManagedCursor.java | 69 ++++++++++++++++--- .../bookkeeper/mledger/ReadOnlyCursor.java | 4 +- .../mledger/impl/ManagedCursorImpl.java | 37 +++++++--- .../impl/ManagedCursorContainerTest.java | 28 ++++++-- .../mledger/impl/ManagedCursorTest.java | 26 +++---- .../mledger/impl/ManagedLedgerTest.java | 2 +- .../persistent/MessageDeduplication.java | 2 +- ...PersistentDispatcherMultipleConsumers.java | 2 +- ...sistentDispatcherSingleActiveConsumer.java | 2 +- .../persistent/PersistentReplicator.java | 2 +- .../buffer/impl/TopicTransactionBuffer.java | 2 +- .../pendingack/impl/MLPendingAckStore.java | 3 +- .../pulsar/compaction/CompactedTopicImpl.java | 5 +- .../broker/delayed/MockManagedCursor.java | 29 ++++++-- ...ckyKeyDispatcherMultipleConsumersTest.java | 4 +- .../broker/transaction/TransactionTest.java | 26 +++---- .../pulsar/sql/presto/PulsarRecordCursor.java | 2 +- .../sql/presto/TestPulsarConnector.java | 2 +- .../sql/presto/TestPulsarRecordCursor.java | 56 +++++++-------- .../impl/MLTransactionLogImpl.java | 3 +- 20 files changed, 207 insertions(+), 99 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java index 1a49608740d48..4f46abe863fb7 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java @@ -150,11 +150,9 @@ enum IndividualDeletedEntries { * opaque context * @param maxPosition * max position can read - * @param skipCondition - * predicate of read filter out */ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition); + PositionImpl maxPosition); /** @@ -165,10 +163,21 @@ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, O * @param callback callback object * @param ctx opaque context * @param maxPosition max position can read - * @param skipCondition predicate of read filter out */ void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition); + Object ctx, PositionImpl maxPosition); + + /** + * Asynchronously read entries from the ManagedLedger. + * + * @param numberOfEntriesToRead maximum number of entries to return + * @param maxSizeBytes max size in bytes of the entries to return + * @param callback callback object + * @param ctx opaque context + * @param maxPosition max position can read + */ + void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, Predicate skipCondition); /** * Get 'N'th entry from the mark delete position in the cursor without updating any cursor positions. @@ -242,11 +251,51 @@ List readEntriesOrWait(int maxEntries, long maxSizeBytes) * opaque context * @param maxPosition * max position can read + */ + void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition); + + /** + * Asynchronously read entries from the ManagedLedger, up to the specified number and size. + * + *

If no entries are available, the callback will not be triggered. Instead it will be registered to wait until + * a new message will be persisted into the managed ledger + * + * @see #readEntriesOrWait(int, long) + * @param maxEntries + * maximum number of entries to return + * @param maxSizeBytes + * max size in bytes of the entries to return + * @param callback + * callback object + * @param ctx + * opaque context + * @param maxPosition + * max position can read + */ + void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition); + + /** + * Asynchronously read entries from the ManagedLedger, up to the specified number and size. + * + *

If no entries are available, the callback will not be triggered. Instead it will be registered to wait until + * a new message will be persisted into the managed ledger + * + * @see #readEntriesOrWait(int, long) + * @param maxEntries + * maximum number of entries to return + * @param callback + * callback object + * @param ctx + * opaque context + * @param maxPosition + * max position can read * @param skipCondition * predicate of read filter out */ - void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition); + void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition); /** * Asynchronously read entries from the ManagedLedger, up to the specified number and size. @@ -268,13 +317,13 @@ void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callb * @param skipCondition * predicate of read filter out */ - void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition); + void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition); /** * Cancel a previously scheduled asyncReadEntriesOrWait operation. * - * @see #asyncReadEntriesOrWait(int, ReadEntriesCallback, Object, PositionImpl, Predicate) + * @see #asyncReadEntriesOrWait(int, ReadEntriesCallback, Object, PositionImpl) * @return true if the read operation was canceled or false if there was no pending operation */ boolean cancelPendingReadRequest(); diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java index 5d67d8a5885cc..18d412f893152 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ReadOnlyCursor.java @@ -48,7 +48,7 @@ public interface ReadOnlyCursor { * @see #readEntries(int) */ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition); + Object ctx, PositionImpl maxPosition); /** * Asynchronously read entries from the ManagedLedger. @@ -60,7 +60,7 @@ void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, * @param maxPosition max position can read */ void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition); + Object ctx, PositionImpl maxPosition); /** * Get the read position. This points to the next message to be read from the cursor. diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index ba87eaff23af8..0307c4847fc16 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -749,7 +749,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter.await(); @@ -762,12 +762,18 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { @Override public void asyncReadEntries(final int numberOfEntriesToRead, final ReadEntriesCallback callback, - final Object ctx, PositionImpl maxPosition, Predicate skipCondition) { - asyncReadEntries(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, skipCondition); + final Object ctx, PositionImpl maxPosition) { + asyncReadEntries(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition); } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition) { + asyncReadEntriesWithSkip(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, null); + } + + @Override + public void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition, Predicate skipCondition) { checkArgument(numberOfEntriesToRead > 0); if (isClosed()) { @@ -882,7 +888,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter.await(); @@ -895,13 +901,27 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { - asyncReadEntriesOrWait(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, skipCondition); + PositionImpl maxPosition) { + asyncReadEntriesOrWait(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition); } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { + PositionImpl maxPosition) { + asyncReadEntriesWithSkipOrWait(maxEntries, maxSizeBytes, callback, ctx, maxPosition, null); + } + + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { + asyncReadEntriesWithSkipOrWait(maxEntries, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, skipCondition); + } + + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { checkArgument(maxEntries > 0); if (isClosed()) { callback.readEntriesFailed(new CursorAlreadyClosedException("Cursor was already closed"), ctx); @@ -915,7 +935,8 @@ public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntrie if (log.isDebugEnabled()) { log.debug("[{}] [{}] Read entries immediately", ledger.getName(), name); } - asyncReadEntries(numberOfEntriesToRead, callback, ctx, maxPosition, skipCondition); + asyncReadEntriesWithSkip(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, + maxPosition, skipCondition); } else { OpReadEntry op = OpReadEntry.create(this, readPosition, numberOfEntriesToRead, callback, ctx, maxPosition, skipCondition); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java index 0b55f71f75d06..cdba480892535 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java @@ -111,16 +111,23 @@ public List readEntries(int numberOfEntriesToRead) throws ManagedLedgerEx @Override public void asyncReadEntries(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { + PositionImpl maxPosition) { callback.readEntriesComplete(null, ctx); } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + Object ctx, PositionImpl maxPosition) { callback.readEntriesComplete(null, ctx); } + @Override + public void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { + + } + @Override public boolean hasMoreEntries() { return true; @@ -302,12 +309,12 @@ public List readEntriesOrWait(int numberOfEntriesToRead) @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { + PositionImpl maxPosition) { } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + Object ctx, PositionImpl maxPosition) { } @@ -404,6 +411,19 @@ public List readEntriesOrWait(int maxEntries, long maxSizeBytes) return null; } + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition) { + + } + + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { + + } + @Override public boolean checkAndUpdateReadPositionChanged() { return false; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index c6a8c2097aaa7..d83975c18f120 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -454,7 +454,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter.await(); } @@ -482,7 +482,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail("async-call should not have failed"); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter.await(); @@ -504,7 +504,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter2.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter2.await(); } @@ -531,7 +531,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); counter.await(); } @@ -568,7 +568,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, null, null, null); + }, null, null); counter.await(); } @@ -1827,7 +1827,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { log.error("Error reading", exception); } - }, null, PositionImpl.LATEST, position -> false); + }, null, PositionImpl.LATEST); } ledger.addEntry("test".getBytes()); @@ -2831,7 +2831,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); assertTrue(c1.cancelPendingReadRequest()); @@ -2847,7 +2847,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { counter2.countDown(); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); ledger.addEntry("entry-1".getBytes(Encoding)); @@ -3754,7 +3754,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { completableFuture.completeExceptionally(exception); } - }, null, (PositionImpl) position, null); + }, null, (PositionImpl) position); int number = completableFuture.get(); assertEquals(number, readMaxNumber); @@ -3769,7 +3769,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { completableFuture.completeExceptionally(exception); } - }, null, (PositionImpl) maxCanReadPosition, null); + }, null, (PositionImpl) maxCanReadPosition); assertEquals(number, sendNumber - readMaxNumber); @@ -4172,7 +4172,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { final int numReadRequests = 3; for (int i = 0; i < numReadRequests; i++) { - cursor.asyncReadEntriesOrWait(1, callback, null, new PositionImpl(0, 0), null); + cursor.asyncReadEntriesOrWait(1, callback, null, new PositionImpl(0, 0)); } Awaitility.await().atMost(Duration.ofSeconds(1)) .untilAsserted(() -> assertEquals(ledger.waitingCursors.size(), 1)); @@ -4272,7 +4272,7 @@ public void testReadEntriesWithFilterOut() throws ManagedLedgerException, Interr } CompletableFuture completableFuture = new CompletableFuture<>(); - cursor.asyncReadEntriesOrWait(sendNumber, new ReadEntriesCallback() { + cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { @Override public void readEntriesComplete(List entries, Object ctx) { completableFuture.complete(entries.size()); @@ -4292,7 +4292,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { assertEquals(cursor.getReadPosition().getEntryId(), 10); CompletableFuture completableFuture2 = new CompletableFuture<>(); - cursor.asyncReadEntriesOrWait(sendNumber, new ReadEntriesCallback() { + cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { @Override public void readEntriesComplete(List entries, Object ctx) { completableFuture2.complete(entries.size()); 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 d85370bf92bf5..896d7312462e4 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 @@ -470,7 +470,7 @@ public void markDeleteFailed(ManagedLedgerException exception, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { fail(exception.getMessage()); } - }, cursor, PositionImpl.LATEST, null); + }, cursor, PositionImpl.LATEST); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java index 2be52903f36c3..030d74a014f29 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/MessageDeduplication.java @@ -193,7 +193,7 @@ public void readEntriesComplete(List entries, Object ctx) { public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { future.completeExceptionally(exception); } - }, null, PositionImpl.LATEST, null); + }, null, PositionImpl.LATEST); } public Status getStatus() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 1f36bba528780..c8e18e2255333 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -329,7 +329,7 @@ public synchronized void readMoreEntries() { .containsMessage(position.getLedgerId(), position.getEntryId()); } } - cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, + cursor.asyncReadEntriesWithSkipOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, topic.getMaxReadPosition(), skipCondition); } else { log.debug("[{}] Cannot schedule next read until previous one is done", name); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java index d41d2a831c6ac..03dce58af3a92 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherSingleActiveConsumer.java @@ -360,7 +360,7 @@ protected void readMoreEntries(Consumer consumer) { ReadEntriesCtx readEntriesCtx = ReadEntriesCtx.create(consumer, consumer.getConsumerEpoch()); cursor.asyncReadEntriesOrWait(messagesToRead, - bytesToRead, this, readEntriesCtx, topic.getMaxReadPosition(), null); + bytesToRead, this, readEntriesCtx, topic.getMaxReadPosition()); } } } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java index d9b8332462461..2d341016f089e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentReplicator.java @@ -249,7 +249,7 @@ protected void readMoreEntries() { log.debug("[{}] Schedule read of {} messages", replicatorId, messagesToRead); } cursor.asyncReadEntriesOrWait(messagesToRead, readMaxSizeBytes, this, - null, topic.getMaxReadPosition(), null); + null, topic.getMaxReadPosition()); } else { if (log.isDebugEnabled()) { log.debug("[{}] Not scheduling read due to pending read. Messages To Read {}", diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java index 4131c03ec5e04..f3bf4f95923cd 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java @@ -677,7 +677,7 @@ boolean fillQueue() { if (cursor.hasMoreEntries()) { outstandingReadsRequests.incrementAndGet(); cursor.asyncReadEntries(NUMBER_OF_PER_READ_ENTRY, - this, System.nanoTime(), PositionImpl.LATEST, null); + this, System.nanoTime(), PositionImpl.LATEST); } else { if (entryQueue.size() == 0) { isReadable = false; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java index 3825e5775844c..4dce8b9a0fcf4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/pendingack/impl/MLPendingAckStore.java @@ -147,8 +147,7 @@ public void replayAsync(PendingAckHandleImpl pendingAckHandle, ExecutorService t //TODO can control the number of entry to read private void readAsync(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback readEntriesCallback) { - cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST, - null); + cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java index f008ae454df29..c8114f9adb652 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactedTopicImpl.java @@ -102,8 +102,7 @@ public void asyncReadEntriesOrWait(ManagedCursor cursor, ReadEntriesCtx readEntriesCtx = ReadEntriesCtx.create(consumer, DEFAULT_CONSUMER_EPOCH); if (compactionHorizon == null || compactionHorizon.compareTo(cursorPosition) < 0) { - cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, PositionImpl.LATEST, - null); + cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, PositionImpl.LATEST); } else { compactedTopicContext.thenCompose( (context) -> findStartPoint(cursorPosition, context.ledger.getLastAddConfirmed(), context.cache) @@ -117,7 +116,7 @@ public void asyncReadEntriesOrWait(ManagedCursor cursor, } if (startPoint == NEWER_THAN_COMPACTED && compactionHorizon.compareTo(cursorPosition) < 0) { cursor.asyncReadEntriesOrWait(numberOfEntriesToRead, callback, readEntriesCtx, - PositionImpl.LATEST, null); + PositionImpl.LATEST); return CompletableFuture.completedFuture(null); } else { long endPoint = Math.min(context.ledger.getLastAddConfirmed(), diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java index 9e3697174ddec..8cb29c0da5541 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java @@ -103,14 +103,20 @@ public List readEntries(int numberOfEntriesToRead) throws InterruptedExce @Override public void asyncReadEntries(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { + PositionImpl maxPosition) { } @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, - AsyncCallbacks.ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition, - Predicate skipCondition) { + AsyncCallbacks.ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { + + } + + @Override + public void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, + AsyncCallbacks.ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition) { } @@ -140,13 +146,26 @@ public List readEntriesOrWait(int maxEntries, long maxSizeBytes) @Override public void asyncReadEntriesOrWait(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + Object ctx, PositionImpl maxPosition) { } @Override public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, AsyncCallbacks.ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + Object ctx, PositionImpl maxPosition) { + + } + + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, AsyncCallbacks.ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition) { + + } + + @Override + public void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, + AsyncCallbacks.ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition) { } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java index eba18936eaab6..1c23afd957a64 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumersTest.java @@ -252,7 +252,7 @@ public void testSkipRedeliverTemporally() { return null; }).when(cursorMock).asyncReadEntriesOrWait( anyInt(), anyLong(), any(PersistentStickyKeyDispatcherMultipleConsumers.class), - eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any(), any()); + eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any()); } catch (Exception e) { fail("Failed to set to field", e); } @@ -421,7 +421,7 @@ public void testMessageRedelivery() throws Exception { return null; }).when(cursorMock).asyncReadEntriesOrWait(anyInt(), anyLong(), any(PersistentStickyKeyDispatcherMultipleConsumers.class), - eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any(), any()); + eq(PersistentStickyKeyDispatcherMultipleConsumers.ReadType.Normal), any()); // (1) Run sendMessagesToConsumers // (2) Attempts to send message1 to consumer1 but skipped because availablePermits is 0 diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 50cf0b3168684..0fcf9ff9b33d7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -42,12 +42,12 @@ import io.netty.util.concurrent.DefaultThreadFactory; import java.lang.reflect.Field; import java.lang.reflect.Method; -import java.nio.charset.StandardCharsets; import java.util.ArrayList; +import java.nio.charset.StandardCharsets; import java.util.Collections; import java.util.HashMap; -import java.util.List; import java.util.Map; +import java.util.List; import java.util.Optional; import java.util.UUID; import java.util.concurrent.CompletableFuture; @@ -121,9 +121,9 @@ import org.apache.pulsar.client.impl.ConsumerBase; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.MessagesImpl; -import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.client.util.ExecutorProvider; import org.apache.pulsar.common.api.proto.CommandSubscribe; +import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.events.EventType; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.SystemTopicNames; @@ -141,8 +141,8 @@ import org.apache.pulsar.transaction.coordinator.TransactionRecoverTracker; import org.apache.pulsar.transaction.coordinator.TransactionTimeoutTracker; import org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl; -import org.apache.pulsar.transaction.coordinator.impl.MLTransactionMetadataStore; import org.apache.pulsar.transaction.coordinator.impl.MLTransactionSequenceIdGenerator; +import org.apache.pulsar.transaction.coordinator.impl.MLTransactionMetadataStore; import org.apache.pulsar.transaction.coordinator.impl.TxnLogBufferedWriterConfig; import org.awaitility.Awaitility; import org.mockito.invocation.InvocationOnMock; @@ -701,7 +701,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); Class managedLedgerClass = ManagedLedgerImpl.class; Field field = managedLedgerClass.getDeclaredField("cursors"); field.setAccessible(true); @@ -713,7 +713,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); TransactionBuffer buffer2 = new TopicTransactionBuffer(persistentTopic); Awaitility.await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> @@ -724,7 +724,7 @@ public void testEndTBRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); managedCursors.add(managedCursor, managedCursor.getMarkDeletedPosition()); TransactionBuffer buffer3 = new TopicTransactionBuffer(persistentTopic); @@ -769,7 +769,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); TransactionPendingAckStoreProvider pendingAckStoreProvider = mock(TransactionPendingAckStoreProvider.class); doReturn(CompletableFuture.completedFuture( @@ -794,7 +794,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); PendingAckHandleImpl pendingAckHandle2 = new PendingAckHandleImpl(persistentSubscription); Awaitility.await().untilAsserted(() -> @@ -804,7 +804,7 @@ public void testEndTPRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); PendingAckHandleImpl pendingAckHandle3 = new PendingAckHandleImpl(persistentSubscription); @@ -838,7 +838,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ callback.readEntriesFailed(new ManagedLedgerException.NonRecoverableLedgerException("No ledger exist"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); MLTransactionSequenceIdGenerator mlTransactionSequenceIdGenerator = new MLTransactionSequenceIdGenerator(); persistentTopic.getManagedLedger().getConfig().setManagedLedgerInterceptor(mlTransactionSequenceIdGenerator); MLTransactionLogImpl mlTransactionLog = @@ -869,7 +869,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.ManagedLedgerFencedException(), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); MLTransactionMetadataStore metadataStore2 = new MLTransactionMetadataStore(new TransactionCoordinatorID(1), @@ -883,7 +883,7 @@ public void testEndTCRecoveringWhenManagerLedgerDisReadable() throws Exception{ AsyncCallbacks.ReadEntriesCallback callback = invocation.getArgument(1); callback.readEntriesFailed(new ManagedLedgerException.CursorAlreadyClosedException("test"), null); return null; - }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any(), any()); + }).when(managedCursor).asyncReadEntries(anyInt(), any(), any(), any()); MLTransactionMetadataStore metadataStore3 = new MLTransactionMetadataStore(new TransactionCoordinatorID(1), diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java index df3927e6c2089..f86335ae3780b 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java @@ -410,7 +410,7 @@ public void run() { // if the available size is invalid and the entry queue size is 0, read one entry outstandingReadsRequests.decrementAndGet(); cursor.asyncReadEntries(batchSize, entryQueueCacheSizeAllocator.getAvailableCacheSize(), - this, System.nanoTime(), PositionImpl.LATEST, null); + this, System.nanoTime(), PositionImpl.LATEST); } // stats for successful read request diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java index a54d9511b49a1..61f6edd38530c 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java @@ -650,7 +650,7 @@ public void run() { return null; } - }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any(), any()); + }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any()); when(readOnlyCursor.hasMoreEntries()).thenAnswer(new Answer() { @Override diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java index c5e089d81b0df..7eaa2da498f45 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarRecordCursor.java @@ -18,22 +18,6 @@ */ package org.apache.pulsar.sql.presto; -import static java.util.concurrent.CompletableFuture.completedFuture; -import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyInt; -import static org.mockito.ArgumentMatchers.anyLong; -import static org.mockito.ArgumentMatchers.anyString; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.doAnswer; -import static org.mockito.Mockito.doReturn; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.when; -import static org.testng.Assert.assertEquals; -import static org.testng.Assert.assertNotNull; -import static org.testng.Assert.assertTrue; -import static org.testng.Assert.fail; import com.fasterxml.jackson.databind.ObjectMapper; import io.airlift.log.Logger; import io.netty.buffer.ByteBuf; @@ -42,17 +26,7 @@ import io.trino.spi.type.RowType; import io.trino.spi.type.Type; import io.trino.spi.type.VarcharType; -import java.lang.reflect.Field; -import java.lang.reflect.InvocationTargetException; -import java.lang.reflect.Method; import java.math.BigDecimal; -import java.nio.charset.Charset; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.HashMap; -import java.util.LinkedList; -import java.util.List; -import java.util.Map; import lombok.Data; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; @@ -83,6 +57,34 @@ import org.mockito.stubbing.Answer; import org.testng.annotations.Test; +import java.lang.reflect.Field; +import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Method; +import java.nio.charset.Charset; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; + +import static java.util.concurrent.CompletableFuture.completedFuture; +import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.assertTrue; +import static org.testng.Assert.fail; + public class TestPulsarRecordCursor extends TestPulsarConnector { private static final Logger log = Logger.get(TestPulsarRecordCursor.class); @@ -383,7 +385,7 @@ public void run() { return null; } - }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any(), any()); + }).when(readOnlyCursor).asyncReadEntries(anyInt(), anyLong(), any(), any(), any()); when(readOnlyCursor.hasMoreEntries()).thenAnswer(new Answer() { @Override diff --git a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java index 1485d6aa15e72..f2e1f60663d28 100644 --- a/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java +++ b/pulsar-transaction/coordinator/src/main/java/org/apache/pulsar/transaction/coordinator/impl/MLTransactionLogImpl.java @@ -156,8 +156,7 @@ public void replayAsync(TransactionLogReplayCallback transactionLogReplayCallbac private void readAsync(int numberOfEntriesToRead, AsyncCallbacks.ReadEntriesCallback readEntriesCallback) { - cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST, - null); + cursor.asyncReadEntries(numberOfEntriesToRead, readEntriesCallback, System.nanoTime(), PositionImpl.LATEST); } @Override From 95012e73a5ed219296c6b09387bdefd5c265ead0 Mon Sep 17 00:00:00 2001 From: coderzc Date: Thu, 29 Dec 2022 16:48:53 +0800 Subject: [PATCH 3/4] Apply comment --- .../bookkeeper/mledger/ManagedCursor.java | 20 ++++++++++++----- .../mledger/impl/ManagedLedgerImpl.java | 22 ++++++++++--------- .../bookkeeper/mledger/impl/OpReadEntry.java | 3 ++- .../impl/ManagedCursorContainerTest.java | 20 ----------------- .../mledger/impl/ManagedCursorTest.java | 3 +++ ...PersistentDispatcherMultipleConsumers.java | 2 +- .../broker/delayed/MockManagedCursor.java | 21 ------------------ 7 files changed, 32 insertions(+), 59 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java index 4f46abe863fb7..7802ed07781ba 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java @@ -175,9 +175,12 @@ void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesC * @param callback callback object * @param ctx opaque context * @param maxPosition max position can read + * @param skipCondition predicate of read filter out */ - void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, Predicate skipCondition); + default void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, Predicate skipCondition) { + asyncReadEntries(numberOfEntriesToRead, maxSizeBytes, callback, ctx, maxPosition); + } /** * Get 'N'th entry from the mark delete position in the cursor without updating any cursor positions. @@ -294,8 +297,10 @@ void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallba * @param skipCondition * predicate of read filter out */ - void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition); + default void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, Object ctx, + PositionImpl maxPosition, Predicate skipCondition) { + asyncReadEntriesOrWait(maxEntries, callback, ctx, maxPosition); + } /** * Asynchronously read entries from the ManagedLedger, up to the specified number and size. @@ -317,8 +322,11 @@ void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback * @param skipCondition * predicate of read filter out */ - void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition); + default void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, + Object ctx, PositionImpl maxPosition, + Predicate skipCondition) { + asyncReadEntriesOrWait(maxEntries, maxSizeBytes, callback, ctx, maxPosition); + } /** * Cancel a previously scheduled asyncReadEntriesOrWait operation. 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 bdece03aa2432..ba91f605acf8c 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 @@ -2062,9 +2062,18 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) long lastValidEntry = -1L; long entryId = firstEntry; for (; entryId <= lastEntry; entryId++) { - if (!opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { - firstValidEntry = entryId; - break; + if (opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { + if (firstValidEntry == -1L) { + firstValidEntry = entryId; + } + } else { + if (firstValidEntry != -1L) { + break; + } + } + + if (firstValidEntry != -1L) { + lastValidEntry = entryId; } } @@ -2076,13 +2085,6 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) return; } - for (; entryId <= lastEntry; entryId++) { - if (opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { - break; - } - lastValidEntry = entryId; - } - firstEntry = firstValidEntry; lastEntry = lastValidEntry; } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java index eff4101368ffe..04437b5320d78 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpReadEntry.java @@ -76,7 +76,7 @@ void internalReadEntriesComplete(List returnedEntries, Object ctx, Positi } cursor.updateReadStats(entriesCount, entriesSize); - if (lastPosition == null || entriesCount != 0) { + if (entriesCount != 0) { lastPosition = (PositionImpl) returnedEntries.get(entriesCount - 1).getPosition(); } if (log.isDebugEnabled()) { @@ -205,6 +205,7 @@ public void recycle() { nextReadPosition = null; maxPosition = null; recyclerHandle.recycle(this); + skipCondition = null; } private static final Logger log = LoggerFactory.getLogger(OpReadEntry.class); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java index cdba480892535..2c01b778caf6b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorContainerTest.java @@ -121,13 +121,6 @@ public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadE callback.readEntriesComplete(null, ctx); } - @Override - public void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, - Predicate skipCondition) { - - } - @Override public boolean hasMoreEntries() { return true; @@ -411,19 +404,6 @@ public List readEntriesOrWait(int maxEntries, long maxSizeBytes) return null; } - @Override - public void asyncReadEntriesWithSkipOrWait(int maxEntries, ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { - - } - - @Override - public void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, ReadEntriesCallback callback, - Object ctx, PositionImpl maxPosition, - Predicate skipCondition) { - - } - @Override public boolean checkAndUpdateReadPositionChanged() { return false; diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java index d83975c18f120..7e1731e573a66 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java @@ -4310,6 +4310,9 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { assertEquals(number2, readMaxNumber / 2); assertEquals(cursor.getReadPosition().getEntryId(), 20); + + cursor.close(); + ledger.close(); } private static final Logger log = LoggerFactory.getLogger(ManagedCursorTest.class); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index c8e18e2255333..d7a1b0b89fa27 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -1056,7 +1056,7 @@ public boolean trackDelayedDelivery(long ledgerId, long entryId, MessageMetadata } } - protected synchronized NavigableSet getMessagesToReplayNow(int maxMessagesToRead) { + protected synchronized NavigableSet getMessagesToReplayNow(int maxMessagesToRead) { if (!redeliveryMessages.isEmpty()) { return redeliveryMessages.getMessagesToReplayNow(maxMessagesToRead); } else if (delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().hasMessageAvailable()) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java index 8cb29c0da5541..efb0fa7ab7ba2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/MockManagedCursor.java @@ -24,7 +24,6 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; -import java.util.function.Predicate; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedCursor; @@ -113,13 +112,6 @@ public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, } - @Override - public void asyncReadEntriesWithSkip(int numberOfEntriesToRead, long maxSizeBytes, - AsyncCallbacks.ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { - - } - @Override public Entry getNthEntry(int n, IndividualDeletedEntries deletedEntries) throws InterruptedException, ManagedLedgerException { @@ -156,19 +148,6 @@ public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, AsyncCallb } - @Override - public void asyncReadEntriesWithSkipOrWait(int maxEntries, AsyncCallbacks.ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { - - } - - @Override - public void asyncReadEntriesWithSkipOrWait(int maxEntries, long maxSizeBytes, - AsyncCallbacks.ReadEntriesCallback callback, Object ctx, - PositionImpl maxPosition, Predicate skipCondition) { - - } - @Override public boolean cancelPendingReadRequest() { return false; From fadff0bf388c277d0fb8ffd9e38af454f1939594 Mon Sep 17 00:00:00 2001 From: coderzc Date: Fri, 30 Dec 2022 11:27:31 +0800 Subject: [PATCH 4/4] rebase master & fix test --- .../bookkeeper/mledger/impl/ManagedCursorImpl.java | 2 +- .../PersistentDispatcherMultipleConsumers.java | 9 ++++++--- 2 files changed, 7 insertions(+), 4 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 0307c4847fc16..5b351c99649ed 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -769,7 +769,7 @@ public void asyncReadEntries(final int numberOfEntriesToRead, final ReadEntriesC @Override public void asyncReadEntries(int numberOfEntriesToRead, long maxSizeBytes, ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { - asyncReadEntriesWithSkip(numberOfEntriesToRead, NO_MAX_SIZE_LIMIT, callback, ctx, maxPosition, null); + asyncReadEntriesWithSkip(numberOfEntriesToRead, maxSizeBytes, callback, ctx, maxPosition, null); } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index d7a1b0b89fa27..9d07d4278fc9f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -320,17 +320,20 @@ public synchronized void readMoreEntries() { minReplayedPosition = null; } - Predicate skipCondition = null; // Filter out and skip read delayed messages exist in DelayedDeliveryTracker if (delayedDeliveryTracker.isPresent()) { + Predicate skipCondition = null; final DelayedDeliveryTracker deliveryTracker = delayedDeliveryTracker.get(); if (deliveryTracker instanceof BucketDelayedDeliveryTracker) { skipCondition = position -> deliveryTracker .containsMessage(position.getLedgerId(), position.getEntryId()); } + cursor.asyncReadEntriesWithSkipOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, + topic.getMaxReadPosition(), skipCondition); + } else { + cursor.asyncReadEntriesOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, + topic.getMaxReadPosition()); } - cursor.asyncReadEntriesWithSkipOrWait(messagesToRead, bytesToRead, this, ReadType.Normal, - topic.getMaxReadPosition(), skipCondition); } else { log.debug("[{}] Cannot schedule next read until previous one is done", name); }