Repository navigation
[fix][broker] Check deliverAt before containsMessage in bucket addMessage - #26230
Conversation
|
@nodece Thanks for the fix — the reordering itself looks right (it matches the When Making the A complete fix would clear the index bit in Given Pulsar's at-least-once semantics this may well be acceptable as-is (an occasional redelivery beats the pre-fix indefinite skip), but it seemed worth recording — and possibly worth handling together with the dispatcher-side expiration check mentioned as follow-up work in the PR description. |
|
@lhotari Good catch. I only considered the case where the expired message could be skipped indefinitely, but I missed the potential duplicate-delivery window caused by leaving the stale entry in the priority queue. The current Bucket implementation has become quite complicated over time — different bucket types, index maintenance, storage lifecycle, and message counters are mixed together, which makes these edge cases harder to reason about. I'm currently refactoring this part and will try to address this issue as part of the redesign. The direction is to make the index state and queue state more consistent, and avoid relying on implicit cleanup during pop operations. |
…sage (apache#26230) (cherry picked from commit 094f270)
…sage (apache#26230) (cherry picked from commit 094f270)
…sage (apache#26230) (cherry picked from commit 094f270)
…sage (apache#26230) (cherry picked from commit 094f270)
Motivation
BucketDelayedDeliveryTrackeruses lazy loading for sealed bucket segments. During segment loading, expired messages are filtered based on theirdeliverAttimestamp and are not added intosharedBucketPriorityQueue.delayedIndexBitMapis persisted with the bucket metadata and restored from BookKeeper. It is used for position deduplication and does not indicate whether a message is currently loaded intosharedBucketPriorityQueue.Therefore, during lazy loading, a message can temporarily exist in
delayedIndexBitMapbut not insharedBucketPriorityQueue.The current
BucketDelayedDeliveryTracker.addMessageimplementation checkscontainsMessagebefore checking whether the message has already expired (deliverAt <= cutoffTime).This causes an issue when an expired message exists only in the delayed index:
containsMessagereturnstruebecause the position exists indelayedIndexBitMap.addMessagereturnstrue, and the dispatcher skips the message.sharedBucketPriorityQueue, so delayed delivery cannot return it either.The message is not lost from storage, but it can remain skip:
Modifications
Change the check order in
BucketDelayedDeliveryTracker.addMessage.deliverAt <= cutoffTimebeforecontainsMessage.falseso they can be delivered immediately by the dispatcher.Add
testExpiredIndexedMessageReturnsFalse.deliverAt.deliverAt.addMessagereturnsfalse.Add
testRecoverThenExpireAddMessage.delayedIndexBitMapbut has not been loaded intosharedBucketPriorityQueue.addMessagereturnsfalse.In addition, the dispatcher side will add an explicit expiration check to make the delivery path more robust. The tracker-side fix ensures that an expired message is not incorrectly blocked by the delayed index when it should bypass delayed delivery tracking.