[fix][ml] Decode the shared cached message metadata before publishing it to other threads - #26824
Merged
Merged
Conversation
Assisted-by: Claude Code (claude-opus-5-5)
This was referenced Oct 4, 2026
dao-jun
approved these changes
Oct 4, 2026
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
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Motivation
When several Key_Shared subscriptions read the same entries from the broker's entry cache, a dispatcher can stop dispatching for good. The Pulsar Performance Testing Framework's catch-up scenario, in which 3 subscriptions catch up on a backlog together, hit this in 3 of 14 runs on master. One subscription stopped with 1.8 million messages in its backlog after this error:
The
Stringthat the partition key was decoded into had no contents. Two things combine to cause this:MessageMetadataonce, inEntryImpl.initializeMessageMetadataIfNeeded. Every copy of the entry shares that instance, including the copies that different subscriptions' dispatchers handle on different threads.getPartitionKey()decodes the field from the parsed buffer on its first call and caches it in a plain field.LightProtoCodec.readStringallocates theStringwithUnsafe.allocateInstanceand stores its byte array with a plain write. Without a constructor, the string doesn't get the safe publication that a final field gives.Stringbefore its contents. It can also decode the field concurrently with the other thread.The exception escaped
readEntriesComplete, so the dispatcher didn't read again, and the subscription stayed stalled until the run timed out.Modifications
EntryImpl.initializeMessageMetadataIfNeededcallsMessageMetadata.materialize()after parsing and before the volatile write that publishes the metadata.materialize()decodes every lazily decoded field, including the properties, the repeated fields and the bytes fields, and drops the reference to the parsed buffer.EntryImplTest.testInitializedMessageMetadataDoesNotDecodeFromTheBufferLateroverwrites the entry's buffer after the metadata is initialized, and checks that every lazily decoded field still has its value. It also compares the serialized metadata with the original, which covers every field, and checks that a copy of the entry shares the instance. Without the change, the test fails.Not in this change:
readEntriesCompletestops reading. That deserves its own fix.Verifying this change
This change added tests and can be verified as follows:
EntryImplTest.testInitializedMessageMetadataDoesNotDecodeFromTheBufferLater, which fails without the change.RangeEntryCacheImplTest.testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntryrelied on the source entry's metadata decoding lazily from its buffer; it now checks that this metadata doesn't depend on the buffer either.The existing
EntryImplTest, the managed ledger's cache tests andManagedLedgerTestpass.The catch-up scenario with the change ([improve][test] Add a catch-up scenario: consumers joining a live stream later #26826): 10 runs, with no stalled subscription and no
readEntriesComplete failedin the broker's log. Without the change, 3 of 14 runs stalled.No measurable cost: in an A/B of the catch-up scenario, the same commit with and without the change, 4 interleaved runs per side, nothing else on the host:
A concurrency test of the race itself would depend on the JIT's reordering, and a passing run would be weak evidence.
Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
This change was prepared with the assistance of Claude Code (claude-opus-5-5); I have reviewed and verified it.