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 82b4cfcc9b68b..02a6b03b943a9 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 @@ -622,9 +622,11 @@ 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); - ledger.asyncReadEntries(op); + if (op != null) { + PENDING_READ_OPS_UPDATER.incrementAndGet(this); + ledger.asyncReadEntries(op); + } } @Override @@ -762,6 +764,9 @@ public void asyncReadEntriesOrWait(int maxEntries, long maxSizeBytes, ReadEntrie } else { OpReadEntry op = OpReadEntry.create(this, readPosition, numberOfEntriesToRead, callback, ctx, maxPosition); + if (op == null) { + return; + } if (!WAITING_READ_OP_UPDATER.compareAndSet(this, null, op)) { callback.readEntriesFailed(new ManagedLedgerException("We can only have a single waiting callback"), 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 0823886e0bdbb..db494a03f1329 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 @@ -2141,6 +2141,7 @@ PositionImpl startReadOperationOnLedger(PositionImpl position, OpReadEntry opRea if (null == ledgerId) { 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); + return null; } if (ledgerId != position.getLedgerId()) { 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 91a6e26f567d0..78cf0d01390f9 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 @@ -48,7 +48,6 @@ class OpReadEntry implements ReadEntriesCallback { public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPositionRef, int count, ReadEntriesCallback callback, Object ctx, PositionImpl maxPosition) { OpReadEntry op = RECYCLER.get(); - op.readPosition = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); op.cursor = cursor; op.count = count; op.callback = callback; @@ -58,6 +57,11 @@ public static OpReadEntry create(ManagedCursorImpl cursor, PositionImpl readPosi } op.maxPosition = maxPosition; op.ctx = ctx; + PositionImpl position = cursor.ledger.startReadOperationOnLedger(readPositionRef, op); + if (position == null) { + return null; + } + op.readPosition = position; op.nextReadPosition = PositionImpl.get(op.readPosition); return op; }