From 65cdb4d8da714081400947b27c6793eedada722d Mon Sep 17 00:00:00 2001 From: "xiaolong.ran" Date: Fri, 6 Dec 2019 18:03:37 +0800 Subject: [PATCH 1/3] [Issue:5669] Fix the ledgerID not found cause NPE Signed-off-by: xiaolong.ran --- .../apache/bookkeeper/mledger/ManagedLedgerException.java | 4 ++++ .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 8 ++++++-- .../org/apache/bookkeeper/mledger/impl/OpReadEntry.java | 6 +++--- 3 files changed, 13 insertions(+), 5 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java index 41aa45b77ce20..698acfd33e74d 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java @@ -71,6 +71,10 @@ public static class ManagedLedgerNotFoundException extends ManagedLedgerExceptio public ManagedLedgerNotFoundException(Exception e) { super(e); } + + public ManagedLedgerNotFoundException(String message) { + super(message); + } } public static class ManagedLedgerTerminatedException extends ManagedLedgerException { 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 0193056b84f21..b16569b389996 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 @@ -1790,8 +1790,12 @@ void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { } } - PositionImpl startReadOperationOnLedger(PositionImpl position) { - long ledgerId = ledgers.ceilingKey(position.getLedgerId()); + PositionImpl startReadOperationOnLedger(PositionImpl position, ReadEntriesCallback callback) { + Long ledgerId = ledgers.ceilingKey(position.getLedgerId()); + if (null == ledgerId) { + callback.readEntriesFailed(new ManagedLedgerNotFoundException("read ledger failed, the ledgerID is null"), null); + } + if (ledgerId != position.getLedgerId()) { // The ledger pointed by this position does not exist anymore. It was deleted because it was empty. We need // to skip on the next available ledger 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 dbe5d259d050d..e3bc08025edfe 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 @@ -49,7 +49,7 @@ class OpReadEntry implements ReadEntriesCallback { public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPositionRef, int count, ReadEntriesCallback callback, Object ctx) { OpReadEntry op = RECYCLER.get(); - op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef); + op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, callback); op.cursor = cursor; op.count = count; op.callback = callback; @@ -128,12 +128,12 @@ void checkReadCompletion() { if (entries.size() < count && cursor.hasMoreEntries()) { // We still have more entries to read from the next ledger, schedule a new async operation if (nextReadPosition.getLedgerId() != readPosition.getLedgerId()) { - cursor.ledger.startReadOperationOnLedger(nextReadPosition); + cursor.ledger.startReadOperationOnLedger(nextReadPosition, callback); } // Schedule next read in a different thread cursor.ledger.getExecutor().execute(safeRun(() -> { - readPosition = cursor.ledger.startReadOperationOnLedger(nextReadPosition); + readPosition = cursor.ledger.startReadOperationOnLedger(nextReadPosition, callback); cursor.ledger.asyncReadEntries(OpReadEntry.this); })); } else { From 16ba1dc3597f669aa226d47b77cf078c1fe43814 Mon Sep 17 00:00:00 2001 From: "xiaolong.ran" Date: Mon, 9 Dec 2019 11:12:52 +0800 Subject: [PATCH 2/3] fix comments Signed-off-by: xiaolong.ran --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 b16569b389996..a20bc2364fad2 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 @@ -1793,7 +1793,8 @@ void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { PositionImpl startReadOperationOnLedger(PositionImpl position, ReadEntriesCallback callback) { Long ledgerId = ledgers.ceilingKey(position.getLedgerId()); if (null == ledgerId) { - callback.readEntriesFailed(new ManagedLedgerNotFoundException("read ledger failed, the ledgerID is null"), null); + callback.readEntriesFailed(new ManagedLedgerException("The ceilingKey(K key) method is used to return the " + + "least key greater than or equal to the given key, or null if there is no such key"), null); } if (ledgerId != position.getLedgerId()) { From 46a0d4a93e465cd647700092f1b7308b417a8b70 Mon Sep 17 00:00:00 2001 From: "xiaolong.ran" Date: Mon, 9 Dec 2019 15:25:07 +0800 Subject: [PATCH 3/3] replace ReadEntriesCallback with OpReadEntry Signed-off-by: xiaolong.ran --- .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 4 ++-- .../org/apache/bookkeeper/mledger/impl/OpReadEntry.java | 6 +++--- 2 files 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 a20bc2364fad2..015471df16e6f 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 @@ -1790,10 +1790,10 @@ void updateCursor(ManagedCursorImpl cursor, PositionImpl newPosition) { } } - PositionImpl startReadOperationOnLedger(PositionImpl position, ReadEntriesCallback callback) { + PositionImpl startReadOperationOnLedger(PositionImpl position, OpReadEntry opReadEntry) { Long ledgerId = ledgers.ceilingKey(position.getLedgerId()); if (null == ledgerId) { - callback.readEntriesFailed(new ManagedLedgerException("The ceilingKey(K key) method is used to return the " + + opReadEntry.readEntriesFailed(new ManagedLedgerException.NoMoreEntriesToReadException("The ceilingKey(K key) method is used to return the " + "least key greater than or equal to the given key, or null if there is no such key"), null); } 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 e3bc08025edfe..c881eb1de5e9b 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 @@ -49,7 +49,7 @@ class OpReadEntry implements ReadEntriesCallback { public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPositionRef, int count, ReadEntriesCallback callback, Object ctx) { OpReadEntry op = RECYCLER.get(); - op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, callback); + op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); op.cursor = cursor; op.count = count; op.callback = callback; @@ -128,12 +128,12 @@ void checkReadCompletion() { if (entries.size() < count && cursor.hasMoreEntries()) { // We still have more entries to read from the next ledger, schedule a new async operation if (nextReadPosition.getLedgerId() != readPosition.getLedgerId()) { - cursor.ledger.startReadOperationOnLedger(nextReadPosition, callback); + cursor.ledger.startReadOperationOnLedger(nextReadPosition, OpReadEntry.this); } // Schedule next read in a different thread cursor.ledger.getExecutor().execute(safeRun(() -> { - readPosition = cursor.ledger.startReadOperationOnLedger(nextReadPosition, callback); + readPosition = cursor.ledger.startReadOperationOnLedger(nextReadPosition, OpReadEntry.this); cursor.ledger.asyncReadEntries(OpReadEntry.this); })); } else {