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

[fix][client]in chunk message consumer, pendingChunkedMessageCount is not synced with count of entries present in chunkedMessagesMap by programmerahul · Pull Request #26689 · apache/pulsar · GitHub

/ pulsar Public

[fix][client]in chunk message consumer, pendingChunkedMessageCount is not synced with count of entries present in chunkedMessagesMap - #26689

Open
programmerahul wants to merge 8 commits into
apache:masterfrom
programmerahul:fix/chunked-message-pending-count-drift
Open

programmerahul wants to merge 8 commits into
apache:masterfrom
programmerahul:fix/chunked-message-pending-count-drift

Conversation

programmerahul commented Sep 22, 2026 •
edited
Loading

Copy link
Copy Markdown
Contributor

Fixes #26688

Motivation

The count pendingChunkedMessageCount is not decrease while the actual count of incomplete chunked messages present in chunkedMessagesMap is decreased when out-of-order or duplicate chunk is received. So this cause unnecessary eviction, since eviction policy is based on this pendingChunkedMessageCount.

Modifications

Decrease the count when ctx uuid entries is removed from chunkedMessagesMap

Verifying this change

  • Make sure that the change passes the CI checks.

(Please pick either of the following options)

This change added tests and can be verified as follows:

  • Added a test-case for this

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

Rahul Prasad and others added 5 commits September 22, 2026 22:38
…/discard

On a duplicate/resent first chunk (or producer restart reusing the sequenceId),
the consumer replaced the existing ChunkedMessageCtx but still incremented
pendingChunkedMessageCount and re-added the uuid to pendingChunkedMessageUuidQueue,
so both drifted above the real chunkedMessagesMap size on every duplicate. The count
drift eventually crosses maxPendingChunkedMessage and evicts a good in-flight message;
the queue growth is unbounded and a stale head can block
removeExpireIncompleteChunkedMessages().

Handle replacement as reuse of the same tracking slot (no count++/queue.add), and only
count+enqueue for a genuinely new uuid. Also decrement the count and remove the uuid
from the queue in the out-of-order discard path when a tracked context is dropped.

Adds a test asserting chunkedMessagesMap.size == pendingChunkedMessageCount ==
pendingChunkedMessageUuidQueue.size after repeated duplicate first chunks.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Thanks for tracking down the pendingChunkedMessageCount drift and for adding regression tests; splitting the replace case from the new-uuid case in ConsumerImpl.java:1639-1681 is a clear improvement, and re-queueing the uuid on replacement addresses the ordering concern @void-ptr974 raised. I confirmed that both new tests fail with the corresponding change reverted and pass with it.

The main thing left is sequencing with #26658 (same author, same bookkeeping). It also edits the completion path and adds the same getPendingChunkedMessageUuidQueueSizeForTest() accessor at the same spot in ConsumerImpl.java, so whichever lands second will conflict, and the test's claim that map, count and queue stay in sync only holds fully once #26658's removal of completed uuids is in. I would suggest landing #26658 first (or combining them) and rebasing this one. #26661 (merged) only touches the discard-branch permit return and does not conflict; chunk reassembly lives only in ConsumerImpl, so the v5 client is not affected.

A few smaller points are inline. For context, not for this PR: processMessageChunk runs on the Netty thread (ClientCnx.handleMessage calls messageReceived directly) while expiry runs on internalPinnedExecutor, so the plain int counter and the recycle/replace of a context are touched from two threads. That race pre-exists, but it is a reason to derive occupancy from the map rather than a separate counter.

} else {
// Genuinely new uuid: count it and enqueue it exactly once. Eviction is only checked
// here because only a new uuid grows the number of in-flight chunked messages.
pendingChunkedMessageCount++;

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

[SUGGESTION] consider deriving the pending count from chunkedMessagesMap instead of a separate counter

The bug here is that a hand-maintained int has to be adjusted on every path that touches the map (completion, replace, out-of-order discard, eviction, expiry), and this PR adds two more adjustments. Since chunkedMessagesMap is the source of truth, the eviction check could be chunkedMessagesMap.size() >= maxPendingChunkedMessage before inserting a new uuid, and pendingChunkedMessageCount (and its ForTest getter) could go away, which removes this whole class of drift. It would also avoid the plain int being modified from both the Netty thread and the expiry task on internalPinnedExecutor.

If you keep the counter, that is fine too; it just means every future change to these paths has to keep it in step by hand.

// (When chunkedMsgCtx is null nothing was ever tracked for this uuid, so there is
// nothing to decrement or dequeue.)
pendingChunkedMessageCount--;
pendingChunkedMessageUuidQueue.remove(msgMetadata.getUuid());

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

[NIT] queue.remove(uuid) is a linear scan on a locked array queue

GrowableArrayBlockingQueue.remove(Object) scans and shifts the backing array while holding its locks, and this runs on the Netty thread. With the default maxPendingChunkedMessage that is small, but with 0 (unbounded) a stream of out-of-order chunks against a large queue gets more expensive. The replace path now does the same remove-then-add. Not a blocker; an insertion-ordered map keyed by uuid would make these O(1).

}

Awaitility.await().atMost(10, TimeUnit.SECONDS)
.untilAsserted(() -> assertEquals(consumerImpl.chunkedMessagesMap.size(), 1));

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

[SUGGESTION] this test can assert before the duplicate chunks are processed, and the out-of-order path is not covered

producer.send() returning does not mean the consumer has processed the chunk, and chunkedMessagesMap.size() == 1 is already true after the first of the ten duplicates, so the assertions below can run before the rest have arrived (and between a replacement's remove and computeIfAbsent the map is briefly empty). It failed reliably for me with the fix reverted, so it is not vacuous in practice, but it is timing-dependent. A deterministic way is to send a complete second chunked message afterwards, receive it (messages are processed in order, so all duplicates are done), and then assert the three sizes.

It would also be good to cover the out-of-order discard (count-- plus queue.remove), and eviction order after a replacement.


// Keep refreshing A's first chunk continuously so A's receivedTime is never older than the
// 2s expiry window -- A must never be eligible for expiry while we observe. We check B's
// removal WHILE still refreshing A: where the

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

[NIT] tidy the new expiry test's comments

This comment has a fragment (WHILE still refreshing A: where the), and the Javadoc refers to "the reviewer's" scenario. Please reword both so the test reads standalone. The loop also sleeps up to ~12s; it passes quickly in the normal case, but it is worth checking it is not close to the 10s-style Awaitility limits used elsewhere in this class if CI is slow.

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug] in chunk message consumer, pendingChunkedMessageCount is not synced with count of entries present in chunkedMessagesMap

3 participants


Back | FazBrowse Home | New Git URL