Repository navigation
[fix][client] Sync ackSet in client with broker to stop acked messages reaching the DLQ - #26135
Conversation
… stop acked messages reaching the DLQ On every batch redelivery, ConsumerImpl created a fresh ackSetInMessageId bitset that never accounted for indices the broker already reports as acked. This kept the batch from ever being seen as fully acked client-side, so PersistentAcknowledgmentsGroupingTracker never cleaned up unAckedMessageTracker / possibleSendToDeadLetterTopicMessages for it -- letting a stale DLQ candidate list survive and later get published to the DLQ topic even though the application had explicitly acknowledged those exact messages. Fixes #29 See also upstream apache#26125 Assisted-by: Claude Code
Pulsar style keeps issue/PR references out of code comments; that context belongs in the PR description instead. Assisted-by: Claude Code
void-ptr974
left a comment
There was a problem hiding this comment.
The main fix in receiveIndividualMessagesFromBatch looks correct to me: syncing the broker-provided ack set into ackSetInMessageId fixes the normal batch path.
However, I think the coverage/fix is incomplete. When a consumer is configured with messagePayloadProcessor(...), ConsumerImpl.messageReceived returns through processPayloadByProcessor(...), and MessagePayloadContextImpl.get(...) still creates a fresh all-unacked ackSetInMessageId. It only uses the broker ack set for ackBitSet, so previously acked/skipped batch indexes are not reflected in the BatchMessageIdImpl.
That means consumers using messagePayloadProcessor(MessagePayloadProcessor.DEFAULT) or a custom processor can still hit the same fully-acked-batch detection failure, leaving unAckedMessageTracker / possibleSendToDeadLetterTopicMessages uncleared.
Could we apply the same ackSetInMessageId.and(BitSet.valueOf(ackSet)) logic to the payload-processor path as well, preferably via a shared helper, and add a regression variant that sets .messagePayloadProcessor(MessagePayloadProcessor.DEFAULT)?
Sounds reasonable. I'll fix it soon. |
…sagePayloadContextImpl MessagePayloadContextImpl#getMessageAt() builds the same kind of per-batch ackSetInMessageId as ConsumerImpl#receiveIndividualMessagesFromBatch, but for messages delivered through a custom MessagePayloadProcessor. It had the same gap fixed in bd6d636: the bitset started as "all unacked" and was never intersected with the broker-reported ackSet, so the MessageId handed back to callers could misreport already-acked indices as still outstanding. Add a regression test covering MessagePayloadContextImpl#getMessageAt directly, verifying the returned MessageId's ackSet is reconciled with the broker ack-set. Assisted-by: Claude Code
@void-ptr974 Fixed, feel free to take a look. |
|
The functional fix looks good now. One issue in the new test: |
…regression test getMessageAt() only releases the retained ByteBuf copy it creates internally (MessagePayloadUtils.convertToByteBuf retains a copy); it doesn't release the caller's MessagePayload. The test created payload0/payload1 via MessagePayloadImpl.create(...) but never released them, leaking ref-counted buffers. Assisted-by: Claude Code
|
That's true, thanks for point it out, fixed. |
|
Since other language clients (pulsar-client-cpp, pulsar-client-go, pulsar-dotnet) seem to follow Pulsar Java client's implementation details, a similar bug might exist there. @Shawyeok Would you like to check pulsar-client-cpp and pulsar-client-go? |
Surely, I'll take a look. |
…s reaching the DLQ (apache#26135) (cherry picked from commit e8df873)
…s reaching the DLQ (apache#26135) (cherry picked from commit e8df873)
Fixes #26125
Motivation
With a
Sharedsubscription, batch-index acknowledgment enabled, and aDeadLetterPolicyconfigured, a message that the application explicitly acknowledged on its final allowed
redelivery round (
redeliveryCount == maxRedeliverCount) could still be routed to the DLQ topic.On every batch redelivery,
ConsumerImpl.receiveIndividualMessagesFromBatchcreates a freshackSetInMessageIdbitset (BatchMessageIdImpl.newAckSet, all bits "unacked") that is neversynced against the broker's own
ackSetfor that redelivery. Indices the broker already knowsare acked are correctly skipped from delivery, but their bit in the fresh
ackSetInMessageIdisnever cleared. As a result,
MessageIdAdvUtils.acknowledge()never sees the batch as fully ackedonce some indices were acked in earlier rounds, so
PersistentAcknowledgmentsGroupingTracker#addIndividualAcknowledgmentnever takes the"fully acked" branch --
unAckedMessageTracker.remove()andpossibleSendToDeadLetterTopicMessages.remove()are never called, leaving a stale DLQ-candidatelist around to be published once the client's own ack-timeout bookkeeping fires again.
On older client versions where batch-index acknowledgment is not enabled by default, the same
root cause shows up as a different symptom: the broker never sends a per-index
ackSetin thismode, so
ackSetInMessageIdis the only thing that decides when the whole entry gets acked andthe cursor advances. Because it's rebuilt from scratch on every redelivery, earlier acks within
the batch are forgotten across rounds and the bitset can permanently fail to reach "fully acked" --
so even messages that were already published to the DLQ (
ConsumerImplacks the original messageright after a successful DLQ publish) may never actually get marked acked in the cursor on the
broker side.
Modifications
ConsumerImpl.receiveIndividualMessagesFromBatch: AND the freshly-createdackSetInMessageIdwith the broker-reported
ackSetbits before processing the batch, so indices already known tobe acked are correctly reflected from the start of this redelivery round.
ackSetfromList<Long>tolong[]end-to-end inmessageReceived/receiveIndividualMessagesFromBatch(and theZeroQueueConsumerImploverride) to avoidredundant List<->array conversions now that the same array is used for both
BitSetRecyclable.valueOfand the newand(...)call.DeadLetterTopicTest#testAckedBatchMessageNotSentToDeadLetterTopicOnFinalRedeliveryRound.Verifying this change
This change added tests and can be verified as follows:
DeadLetterTopicTest#testAckedBatchMessageNotSentToDeadLetterTopicOnFinalRedeliveryRound:sends a 5-message batch, deliberately lets batch indices 1 and 2 time out for
maxRedeliverCountrounds and acks them on the final allowed round, then asserts nothing isever delivered to the DLQ topic.
passes once the fix is applied.
issue, both on a non-partitioned topic and on a 3-partition topic.
Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes