From f1a7bd767f344b4c1a1f59d100b7e019abe3b45a Mon Sep 17 00:00:00 2001 From: coderzc Date: Tue, 3 Jan 2023 23:01:48 +0800 Subject: [PATCH 1/5] fix cursor skip read --- .../mledger/impl/ManagedLedgerImpl.java | 2 +- .../mledger/impl/ManagedCursorTest.java | 54 ++++++++++++++++++- 2 files changed, 53 insertions(+), 3 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 3a80d12227224..7be1a54229ed6 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 @@ -2060,7 +2060,7 @@ 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))) { + if (!opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { if (firstValidEntry == -1L) { firstValidEntry = entryId; } 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 7e1731e573a66..df062dfc5230a 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 @@ -4275,7 +4275,14 @@ public void testReadEntriesWithFilterOut() throws ManagedLedgerException, Interr cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { @Override public void readEntriesComplete(List entries, Object ctx) { - completableFuture.complete(entries.size()); + try { + entries.forEach(entry -> { + assertNotEquals(entry.getEntryId() % 2, 0); + }); + completableFuture.complete(entries.size()); + } catch (Exception e) { + completableFuture.completeExceptionally(e); + } } @Override @@ -4295,7 +4302,14 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { @Override public void readEntriesComplete(List entries, Object ctx) { - completableFuture2.complete(entries.size()); + try { + entries.forEach(entry -> { + assertEquals(entry.getEntryId() % 2, 0); + }); + completableFuture2.complete(entries.size()); + } catch (Exception e) { + completableFuture2.completeExceptionally(e); + } } @Override @@ -4311,6 +4325,42 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { assertEquals(cursor.getReadPosition().getEntryId(), 20); + cursor.seek(PositionImpl.EARLIEST); + CompletableFuture completableFuture3 = new CompletableFuture<>(); + cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + completableFuture3.complete(entries.size()); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + completableFuture3.completeExceptionally(exception); + } + }, null, (PositionImpl) maxCanReadPosition, pos -> false); + + int number3 = completableFuture3.get(); + assertEquals(number3, sendNumber); + assertEquals(cursor.getReadPosition().getEntryId(), 20); + + cursor.seek(PositionImpl.EARLIEST); + CompletableFuture completableFuture4 = new CompletableFuture<>(); + cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { + @Override + public void readEntriesComplete(List entries, Object ctx) { + completableFuture4.complete(entries.size()); + } + + @Override + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + completableFuture4.completeExceptionally(exception); + } + }, null, (PositionImpl) maxCanReadPosition, pos -> true); + + int number4 = completableFuture4.get(); + assertEquals(number4, 0); + assertEquals(cursor.getReadPosition().getEntryId(), 20); + cursor.close(); ledger.close(); } From 70382171d4a52096547bcbe0c5adc9ad69616d27 Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 4 Jan 2023 11:32:16 +0800 Subject: [PATCH 2/5] fix test --- .../mledger/impl/ManagedCursorTest.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) 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 df062dfc5230a..35ef00db06b94 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 @@ -4277,10 +4277,10 @@ public void testReadEntriesWithFilterOut() throws ManagedLedgerException, Interr public void readEntriesComplete(List entries, Object ctx) { try { entries.forEach(entry -> { - assertNotEquals(entry.getEntryId() % 2, 0); + assertEquals(entry.getEntryId() % 2, 0L); }); completableFuture.complete(entries.size()); - } catch (Exception e) { + } catch (Throwable e) { completableFuture.completeExceptionally(e); } } @@ -4296,7 +4296,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { int number = completableFuture.get(); assertEquals(number, readMaxNumber / 2); - assertEquals(cursor.getReadPosition().getEntryId(), 10); + assertEquals(cursor.getReadPosition().getEntryId(), 10L); CompletableFuture completableFuture2 = new CompletableFuture<>(); cursor.asyncReadEntriesWithSkipOrWait(sendNumber, new ReadEntriesCallback() { @@ -4304,10 +4304,10 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { public void readEntriesComplete(List entries, Object ctx) { try { entries.forEach(entry -> { - assertEquals(entry.getEntryId() % 2, 0); + assertEquals(entry.getEntryId() % 2, 0L); }); completableFuture2.complete(entries.size()); - } catch (Exception e) { + } catch (Throwable e) { completableFuture2.completeExceptionally(e); } } @@ -4323,7 +4323,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { int number2 = completableFuture2.get(); assertEquals(number2, readMaxNumber / 2); - assertEquals(cursor.getReadPosition().getEntryId(), 20); + assertEquals(cursor.getReadPosition().getEntryId(), 20L); cursor.seek(PositionImpl.EARLIEST); CompletableFuture completableFuture3 = new CompletableFuture<>(); @@ -4341,7 +4341,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { int number3 = completableFuture3.get(); assertEquals(number3, sendNumber); - assertEquals(cursor.getReadPosition().getEntryId(), 20); + assertEquals(cursor.getReadPosition().getEntryId(), 20L); cursor.seek(PositionImpl.EARLIEST); CompletableFuture completableFuture4 = new CompletableFuture<>(); @@ -4359,7 +4359,7 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { int number4 = completableFuture4.get(); assertEquals(number4, 0); - assertEquals(cursor.getReadPosition().getEntryId(), 20); + assertEquals(cursor.getReadPosition().getEntryId(), 20L); cursor.close(); ledger.close(); From 07187a3fc2df5c32939ed13e562726e0e7dba791 Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 4 Jan 2023 11:42:58 +0800 Subject: [PATCH 3/5] fix test name --- .../org/apache/bookkeeper/mledger/impl/ManagedCursorTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 35ef00db06b94..8dc726c249efc 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 @@ -4254,10 +4254,10 @@ public void testLazyCursorLedgerCreationForSubscriptionCreation() throws Excepti } @Test - public void testReadEntriesWithFilterOut() throws ManagedLedgerException, InterruptedException, ExecutionException { + public void testReadEntriesWithSkip() throws ManagedLedgerException, InterruptedException, ExecutionException { int readMaxNumber = 10; int sendNumber = 20; - ManagedLedger ledger = factory.open("testReadEntriesWithFilter"); + ManagedLedger ledger = factory.open("testReadEntriesWithSkip"); ManagedCursor cursor = ledger.openCursor("c"); Position position = PositionImpl.EARLIEST; Position maxCanReadPosition = PositionImpl.EARLIEST; From e67f6a42d435baa17382ed6fffd42f9088bf3d7c Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 4 Jan 2023 16:12:50 +0800 Subject: [PATCH 4/5] Invert if condition --- .../bookkeeper/mledger/impl/ManagedLedgerImpl.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index 7be1a54229ed6..c2b6bf0b3b865 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 @@ -2060,14 +2060,14 @@ 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))) { - if (firstValidEntry == -1L) { - firstValidEntry = entryId; - } - } else { + if (opReadEntry.skipCondition.test(PositionImpl.get(ledger.getId(), entryId))) { if (firstValidEntry != -1L) { break; } + } else { + if (firstValidEntry == -1L) { + firstValidEntry = entryId; + } } if (firstValidEntry != -1L) { From a32996f0e0888b621adc168f1b15682647b451c2 Mon Sep 17 00:00:00 2001 From: coderzc Date: Wed, 4 Jan 2023 19:17:51 +0800 Subject: [PATCH 5/5] improve code --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index c2b6bf0b3b865..7b8fdf1f816fa 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 @@ -2068,9 +2068,7 @@ private void internalReadFromLedger(ReadHandle ledger, OpReadEntry opReadEntry) if (firstValidEntry == -1L) { firstValidEntry = entryId; } - } - if (firstValidEntry != -1L) { lastValidEntry = entryId; } }