Skip to content

Commit

Permalink
[fix][broker] Fix deadlock in PendingAckHandleImpl (#18989)
Browse files Browse the repository at this point in the history
(cherry picked from commit 22866bd)
  • Loading branch information
nicoloboschi committed Dec 20, 2022
1 parent 7c63e92 commit 523a0da
Showing 1 changed file with 10 additions and 12 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -146,20 +146,20 @@ public PendingAckHandleImpl(PersistentSubscription persistentSubscription) {
.getBrokerService().getPulsar().getTransactionPendingAckStoreProvider();

pendingAckStoreProvider.checkInitializedBefore(persistentSubscription)
.thenAccept(init -> {
.thenAcceptAsync(init -> {
if (init) {
initPendingAckStore();
} else {
completeHandleFuture();
}
})
.exceptionally(e -> {
}, internalPinnedExecutor)
.exceptionallyAsync(e -> {
Throwable t = FutureUtil.unwrapCompletionException(e);
changeToErrorState();
exceptionHandleFuture(t);
this.pendingAckStoreFuture.completeExceptionally(t);
return null;
});
}, internalPinnedExecutor);
}

private void initPendingAckStore() {
Expand Down Expand Up @@ -925,18 +925,16 @@ public TransactionPendingAckStats getStats(boolean lowWaterMarks) {
return transactionPendingAckStats;
}

public synchronized void completeHandleFuture() {
if (!this.pendingAckHandleCompletableFuture.isDone()) {
this.pendingAckHandleCompletableFuture.complete(PendingAckHandleImpl.this);
}
if (recoverTime.getRecoverStartTime() != 0L) {
public void completeHandleFuture() {
this.pendingAckHandleCompletableFuture.complete(PendingAckHandleImpl.this);
if (recoverTime.getRecoverStartTime() != 0L && recoverTime.getRecoverEndTime() == 0L) {
recoverTime.setRecoverEndTime(System.currentTimeMillis());
}
}

public synchronized void exceptionHandleFuture(Throwable t) {
if (!this.pendingAckHandleCompletableFuture.isDone()) {
this.pendingAckHandleCompletableFuture.completeExceptionally(t);
public void exceptionHandleFuture(Throwable t) {
final boolean completedNow = this.pendingAckHandleCompletableFuture.completeExceptionally(t);
if (completedNow) {
recoverTime.setRecoverEndTime(System.currentTimeMillis());
}
}
Expand Down

0 comments on commit 523a0da

Please sign in to comment.