FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

[fix][broker] Fix unacked accounting races during consumer removal by Denovo1998 · Pull Request #26663 · apache/pulsar · GitHub

/ pulsar Public

[fix][broker] Fix unacked accounting races during consumer removal - #26663

Draft
Denovo1998 wants to merge 4 commits into
apache:masterfrom
Denovo1998:fix-unacked-message-settlement
Draft

Denovo1998 wants to merge 4 commits into
apache:masterfrom
Denovo1998:fix-unacked-message-settlement

Conversation

Copy link
Copy Markdown
Contributor

Motivation

Consumer removal can race with ACK completion, redelivery, or mark-delete cleanup, causing subscription and broker unacked counters to become inconsistent. Broker throttling also has lock-order and registration races that can block consumer removal or leave a dispatcher blocked after its unacked count has fallen below the low watermark.

These races can disrupt flow control and prevent subscriptions from resuming delivery.

Modifications

  • Serialize each consumer's unacked accounting with removal, settle its remaining balance once, and reject subsequent accounting updates after removal.
  • Apply the settlement lifecycle to both modern and classic shared dispatchers.
  • Resolve subscriptions before acquiring the broker unacked lock and schedule resumed reads after releasing it.
  • Serialize broker blocking transitions with dispatcher registration, and recheck the low watermark after registration to cancel a late block.
  • Coalesce eligible whole-entry ACK completions by consecutive consumer owner, using the actual removed message count. Preserve the original path for single-entry, batch-index, and transaction-enabled processing, and flush previously returned deltas if a later removal fails.
  • Extend the existing accounting test class with deterministic removal, blocking, batched-entry ACK, and Key_Shared ownership-transfer coverage.
  • Add ConsumerUnackedMessagesBenchmark under microbench/ to measure accounting contention and grouped ACK completion.

Verifying this change

  • Make sure that the change passes the CI checks.

(Please pick either of the following options)

This change is a trivial rework / code cleanup without any test coverage.

(or)

This change is already covered by existing tests, such as (please describe tests).

(or)

This change added tests and can be verified as follows:

(example:)

  • Added integration tests for end-to-end deployment with large payloads (10MB)
  • Extended integration test for recovery after broker failure

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

Denovo1998 marked this pull request as draft September 20, 2026 13:25
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
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.

1 participant


Back | FazBrowse Home | New Git URL