Repository navigation
[fix][client] Fix ZeroQueueConsumer stuck on pending receiveAsync after reconnection - #26814
subodhwandile wants to merge 2 commits into
Conversation
…er reconnection Fixes apache#15679
|
@lhotari could you apply the ready-to-test label so CI can run? Thank you! |
|
You can run CI on your fork like this |
|
Personal CI passed on my fork: subodhwandile#1 |
| for (CompletableFuture<Message<T>> pendingReceive : pendingReceives) { | ||
| if (!pendingReceive.isDone()) { | ||
| pendingAsyncReceives++; | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| 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); | ||
| } |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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
be2664f to
0f50540
Compare
|
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). |
|
@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? |
Fixes #15679
Motivation
A consumer with
receiverQueueSize(0)that is waiting inreceiveAsync()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 setwaitingOnReceiveForZeroQueueSize(only blockingreceive()sets it). On reconnection,consumerIsReconnectedToBroker()resets the available permits to 0 and re-sends a permit only for a blockingreceive(), 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 booleanflag set by eachinternalReceiveAsync()call, which breaks under concurrent calls — only the last caller's flag survives.Modifications
consumerIsReconnectedToBroker(), count the pendingreceiveAsync()futures that are not done, and sendmax(pendingAsyncReceives, existingCondition ? 1 : 0)permits.ZeroQueueSizeTest#testReceiveAsyncNotStuckAfterTopicUnloadand#testMultiplePendingReceiveAsyncNotStuckAfterTopicUnload. Both fail withTimeoutExceptionon master and pass with this change.Verifying this change
ZeroQueueSizeTestpasses (19 tests).Documentation
doc-not-needed