Skip to content

[fix][client] Sync ackSet in client with broker to stop acked messages reaching the DLQ - #26135

Merged
lhotari merged 4 commits into
apache:masterfrom
Shawyeok:fix-dlq-misrouting-batch-ack
Jul 3, 2026
Merged

lhotari merged 4 commits into
apache:masterfrom
Shawyeok:fix-dlq-misrouting-batch-ack

Conversation

@Shawyeok

@Shawyeok Shawyeok commented Jul 2, 2026

Copy link
Copy Markdown
Contributor

Fixes #26125

Motivation

With a Shared subscription, batch-index acknowledgment enabled, and a DeadLetterPolicy
configured, 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.receiveIndividualMessagesFromBatch creates a fresh
ackSetInMessageId bitset (BatchMessageIdImpl.newAckSet, all bits "unacked") that is never
synced against the broker's own ackSet for that redelivery. Indices the broker already knows
are acked are correctly skipped from delivery, but their bit in the fresh ackSetInMessageId is
never cleared. As a result, MessageIdAdvUtils.acknowledge() never sees the batch as fully acked
once some indices were acked in earlier rounds, so
PersistentAcknowledgmentsGroupingTracker#addIndividualAcknowledgment never takes the
"fully acked" branch -- unAckedMessageTracker.remove() and
possibleSendToDeadLetterTopicMessages.remove() are never called, leaving a stale DLQ-candidate
list 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 ackSet in this
mode, so ackSetInMessageId is the only thing that decides when the whole entry gets acked and
the 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 (ConsumerImpl acks the original message
right after a successful DLQ publish) may never actually get marked acked in the cursor on the
broker side.

Modifications

  • ConsumerImpl.receiveIndividualMessagesFromBatch: AND the freshly-created ackSetInMessageId
    with the broker-reported ackSet bits before processing the batch, so indices already known to
    be acked are correctly reflected from the start of this redelivery round.
  • Minor: switch ackSet from List<Long> to long[] end-to-end in messageReceived /
    receiveIndividualMessagesFromBatch (and the ZeroQueueConsumerImpl override) to avoid
    redundant List<->array conversions now that the same array is used for both
    BitSetRecyclable.valueOf and the new and(...) call.
  • Added DeadLetterTopicTest#testAckedBatchMessageNotSentToDeadLetterTopicOnFinalRedeliveryRound.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • Added DeadLetterTopicTest#testAckedBatchMessageNotSentToDeadLetterTopicOnFinalRedeliveryRound:
    sends a 5-message batch, deliberately lets batch indices 1 and 2 time out for
    maxRedeliverCount rounds and acks them on the final allowed round, then asserts nothing is
    ever delivered to the DLQ topic.
  • Confirmed this test fails without the fix (receives an unexpected DLQ message) and
    passes once the fix is applied.
  • Manually verified against a live standalone broker with the original Java reproducer from the
    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

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Shawyeok added 2 commits July 1, 2026 18:41
… 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 void-ptr974 left a comment

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.

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)?

@Shawyeok

Shawyeok commented Jul 3, 2026

Copy link
Copy Markdown
Contributor Author

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
@Shawyeok

Shawyeok commented Jul 3, 2026

Copy link
Copy Markdown
Contributor Author

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.

@void-ptr974 Fixed, feel free to take a look.

@void-ptr974

Copy link
Copy Markdown
Contributor

The functional fix looks good now. One issue in the new test: payload0 and payload1 are created with MessagePayloadImpl.create(...) but are never released. getMessageAt(...) only releases the retained buffer it creates internally, not the original MessagePayloadImpl, so this test leaks ref-counted buffers.

…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
@Shawyeok

Shawyeok commented Jul 3, 2026

Copy link
Copy Markdown
Contributor Author

That's true, thanks for point it out, fixed.

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

LGTM, good work @Shawyeok. Just fix the minor issue that was pointed out earlier.

@lhotari

lhotari commented Jul 3, 2026

Copy link
Copy Markdown
Member

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?

@Shawyeok

Shawyeok commented Jul 3, 2026

Copy link
Copy Markdown
Contributor Author

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.

@lhotari
lhotari merged commit e8df873 into apache:master Jul 3, 2026
43 checks passed
lhotari pushed a commit that referenced this pull request Jul 22, 2026
lhotari pushed a commit that referenced this pull request Jul 22, 2026
sandeep-ctds pushed a commit to datastax/pulsar that referenced this pull request Jul 31, 2026
nodece pushed a commit to ascentstream/pulsar that referenced this pull request Aug 28, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[bug][client] Consumer sends already-acknowledged messages to DLQ on final redelivery round

3 participants