From 0ec7b8e812eb9a44e49c3b0f98addce758b19d02 Mon Sep 17 00:00:00 2001 From: dezhiliu Date: Sun, 23 Feb 2020 22:23:12 +0800 Subject: [PATCH] fix when send a delayed message ,there is a case when a consumer restarts and pull duplicate messages. --- .../bookkeeper/mledger/impl/ManagedCursorImpl.java | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) 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 9498145de2823..05f6f5c762133 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 @@ -1116,7 +1116,14 @@ public synchronized void readEntryFailed(ManagedLedgerException mle, Object ctx) }; positions.stream().filter(position -> !alreadyAcknowledgedPositions.contains(position)) - .forEach(p -> ledger.asyncReadEntry((PositionImpl) p, cb, ctx)); + .forEach(p ->{ + if (((PositionImpl) p).compareTo(this.readPosition) == 0) { + this.setReadPosition(this.readPosition.getNext()); + log.warn("[{}][{}] replayPosition{} equals readPosition{}," + " need set next readPositio", + ledger.getName(), name, (PositionImpl) p, this.readPosition); + } + ledger.asyncReadEntry((PositionImpl) p, cb, ctx); + }); return alreadyAcknowledgedPositions; }