Skip to content

[fix][ml] Decode the shared cached message metadata before publishing it to other threads - #26824

Merged
lhotari merged 1 commit into
apache:masterfrom
lhotari:lh-fix-ml-cached-metadata-race
Oct 5, 2026
Merged

lhotari merged 1 commit into
apache:masterfrom
lhotari:lh-fix-ml-cached-metadata-race

Conversation

@lhotari

@lhotari lhotari commented Oct 4, 2026 •

Copy link
Copy Markdown
Member

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:

ERROR org.apache.bookkeeper.mledger.impl.OpReadEntry - readEntriesComplete failed {..., op=public/default/persistent/iot-telemetry-0 iot-application-1{ readPosition: 2:4623525, ... entries count: 31
java.lang.NullPointerException: Cannot invoke "[B.clone()" because "<parameter2>" is null
	at java.base/java.lang.String.encode8859_1(Unknown Source)
	at java.base/java.lang.String.getBytes(Unknown Source)
	at java.base/java.util.Base64$Decoder.decode(Unknown Source)
	at org.apache.pulsar.common.protocol.Commands.resolveStickyKey(Commands.java:2528)
	at org.apache.pulsar.broker.service.EntryAndMetadata.getStickyKey(EntryAndMetadata.java:66)
	...
	at org.apache.pulsar.broker.service.persistent.PersistentStickyKeyDispatcherMultipleConsumers.filterAndGroupEntriesForDispatching(PersistentStickyKeyDispatcherMultipleConsumers.java:492)
	...
	at org.apache.bookkeeper.mledger.impl.cache.RangeEntryCacheImpl.doAsyncReadEntriesByPosition(RangeEntryCacheImpl.java:464)

The String that the partition key was decoded into had no contents. Two things combine to cause this:

  • The cached metadata is shared across threads: since [improve][broker] Reduce unnecessary MessageMetadata parsing by caching the parsed instance in the broker cache #24682, the entry cache parses an entry's MessageMetadata once, in EntryImpl.initializeMessageMetadataIfNeeded. Every copy of the entry shares that instance, including the copies that different subscriptions' dispatchers handle on different threads.
  • LightProto decodes fields lazily without safe publication:
    • A generated getter such as getPartitionKey() decodes the field from the parsed buffer on its first call and caches it in a plain field.
    • For an ASCII string, LightProtoCodec.readString allocates the String with Unsafe.allocateInstance and stores its byte array with a plain write. Without a constructor, the string doesn't get the safe publication that a final field gives.
    • A dispatcher that reads the cached field while another thread decodes it can see the String before 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

  • Materialize before publishing: EntryImpl.initializeMessageMetadataIfNeeded calls MessageMetadata.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.
    • The getters then only read fields that the volatile write published, so the threads that share the instance don't decode anything.
    • All the paths that share the parsed metadata go through this method: a cache hit, entries read from storage and inserted into the cache, and the shared reads with the cache disabled.
  • Test: EntryImplTest.testInitializedMessageMetadataDoesNotDecodeFromTheBufferLater overwrites 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.
  • Cost: decoding the fields once per entry adds a little work when the entry is first read. The metadata also no longer references the entry's buffer.

Not in this change:

  • Dispatcher recovery: a dispatcher that gets an exception from readEntriesComplete stops reading. That deserves its own fix.
  • LightProto itself: its unsafe string creation could add a release fence, which would make its other lazy getters safe to share.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • EntryImplTest.testInitializedMessageMetadataDoesNotDecodeFromTheBufferLater, which fails without the change.

  • RangeEntryCacheImplTest.testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry relied 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 and ManagedLedgerTest pass.

  • 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 failed in 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:

    Catch-up scenario Master (3 valid runs; the fourth stalled) This change (4 runs)
    The first of 3 applications joining together caught up after 89.1 s (84.8–92.1) 90.1 s (87.0–92.6)
    The application joining alone at 60 s caught up after 102.8 s (98.5–107.2) 104.3 s (101.4–109.3)
    Broker CPU per million messages 154.9 s (154.7–155.0) 155.1 s (154.3–155.7)
  • 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

  • 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

This change was prepared with the assistance of Claude Code (claude-opus-5-5); I have reviewed and verified it.

@lhotari
lhotari merged commit 24f2fd5 into apache:master Oct 5, 2026
83 of 85 checks passed
lhotari added a commit that referenced this pull request Oct 5, 2026
… it to other threads (#26824)

(cherry picked from commit 24f2fd5)
lhotari added a commit that referenced this pull request Oct 5, 2026
… it to other threads (#26824)

(cherry picked from commit 24f2fd5)
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.

2 participants