Skip to content

[fix][client] Fix ZeroQueueConsumer stuck on pending receiveAsync after reconnection - #26814

Open
subodhwandile wants to merge 2 commits into
apache:masterfrom
subodhwandile:fix/zero-queue-consumer-receive-async-reconnect
Open

subodhwandile wants to merge 2 commits into
apache:masterfrom
subodhwandile:fix/zero-queue-consumer-receive-async-reconnect

Conversation

@subodhwandile

Copy link
Copy Markdown

Fixes #15679

Motivation

A consumer with receiverQueueSize(0) that is waiting in receiveAsync() never receives another message after its topic is unloaded or the broker restarts.

ZeroQueueConsumerImpl#internalReceiveAsync() sends one permit for the pending future but does not set waitingOnReceiveForZeroQueueSize (only blocking receive() sets it). On reconnection, consumerIsReconnectedToBroker() resets the available permits to 0 and re-sends a permit only for a blocking receive(), a non-empty queue, or an idle listener. Pending async receives are not counted, so the new broker-side consumer has 0 permits and the futures never complete.

This also explains why the fix attempt in #15852 stalled: it used a single volatile boolean flag set by each internalReceiveAsync() call, which breaks under concurrent calls — only the last caller's flag survives.

Modifications

  • In consumerIsReconnectedToBroker(), count the pending receiveAsync() futures that are not done, and send max(pendingAsyncReceives, existingCondition ? 1 : 0) permits.
  • Add ZeroQueueSizeTest#testReceiveAsyncNotStuckAfterTopicUnload and #testMultiplePendingReceiveAsyncNotStuckAfterTopicUnload. Both fail with TimeoutException on master and pass with this change.

Verifying this change

  • Covered by the new tests. The full ZeroQueueSizeTest passes (19 tests).

Documentation

  • doc-not-needed

@subodhwandile

Copy link
Copy Markdown
Author

@lhotari could you apply the ready-to-test label so CI can run? Thank you!

@hohoho1886

Copy link
Copy Markdown

You can run CI on your fork like this

@subodhwandile

subodhwandile commented Oct 3, 2026 •

Copy link
Copy Markdown
Author

Personal CI passed on my fork: subodhwandile#1
@lhotari could you apply the ready-to-test label?

Comment on lines +146 to +150
for (CompletableFuture<Message<T>> pendingReceive : pendingReceives) {
if (!pendingReceive.isDone()) {
pendingAsyncReceives++;
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think there's still a race here. ConsumerImpl.internalReceiveAsync() adds the future to pendingReceives asynchronously on internalPinnedExecutor, but ZeroQueueConsumerImpl.internalReceiveAsync() sends the permit right away.

If the executor task hasn't run when reconnection happens, this loop won't see the pending future. The original permit may have gone to the old connection, so the new broker gets no permit and the future can still hang.

The opposite timing could also result in duplicate permits.

Could we coordinate the pending-receive registration and permit sending with the reconnect path? A test that deliberately delays the executor task until after reconnection would help catch this.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. The permit is now sent inside the internalPinnedExecutor task, under synchronized(ZeroQueueConsumerImpl.this) , the same lock held during consumerIsReconnectedToBroker(). Either reconnect counts the future and sends its permit, or the task does after reconnection. Never lost, never duplicated. Same fix applied to fetchSingleMessageFromBroker(): waitingOnReceiveForZeroQueueSize = true moved inside the synchronized block.

Comment on lines +151 to 157
boolean shouldSendPermit = waitingOnReceiveForZeroQueueSize
|| currentQueueSize > 0
|| (listener != null && !waitingOnListenerForZeroQueueSize)) {
increaseAvailablePermits(cnx);
|| (listener != null && !waitingOnListenerForZeroQueueSize);
int permits = Math.max(pendingAsyncReceives, shouldSendPermit ? 1 : 0);
if (permits > 0) {
increaseAvailablePermits(cnx, permits);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could this undercount permits when a blocking receive() and some receiveAsync() calls are waiting at the same time?

For example, if there's one blocking receive and three pending async receives, Math.max(3, 1) only restores three permits, even though there are four waiting requests.

Since incoming messages are delivered to pending async receives first, the blocking receive could still get stuck after reconnection.

I think we need to account for the blocking receive separately here. Could we also add a regression test that mixes receive() and receiveAsync() across a topic unload?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. Changed to pendingAsyncReceives + (shouldSendPermit ? 1 : 0). With 1 blocking receive() and 3 receiveAsync(), this correctly sends 4 permits after reconnection.

…onnection

Register pending receiveAsync and send its permit atomically under the consumer
lock, move waitingOnReceiveForZeroQueueSize inside the synchronized block, and
send one permit per pending receive plus one for a waiting receive on reconnect.

Fixes apache#15679
@subodhwandile
subodhwandile force-pushed the fix/zero-queue-consumer-receive-async-reconnect branch from be2664f to 0f50540 Compare October 11, 2026 08:21
@subodhwandile

Copy link
Copy Markdown
Author

Both fixed in 0f50540. Added testReceiveAsyncRegisteredAfterReconnectionNotStuck (blocks executor until after reconnect, asserts exactly 1 broker permit) and testMixedReceiveAndReceiveAsyncNotStuckAfterTopicUnload (1 receive() + 3 receiveAsync(), asserts 4 permits before and after unload).

@subodhwandile

Copy link
Copy Markdown
Author

@lhotari fork CI passed on the updated commit: https://github.com/subodhwandile/pulsar/actions/runs/38125151511 and @Denovo1998 has approved. Could you apply the ready-to-test label?

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Topic rebalances cause consumers with receiverQueueSize=0 to get stuck when calling receiveAsync()

3 participants