Repository navigation
Conversation
… consumer is created DagWatchClient stored every layout the broker pushed but delivered it only to a listener that was already set, and setListener did not hand over what had arrived before it. An ungrouped checkpoint consumer and a queue consumer apply the initial layout first and set their listener only once a reader or consumer exists for every segment, so a split or merge pushed in between was dropped: the consumer kept reading the sealed parent, never attached to its children, and stopped receiving the keys they cover until the broker happened to push a layout again. The producer set its listener before applying its initial layout, so a layout delivered in between was overwritten by the older initial one. setListener now hands the listener the current layout when it is not the initial one that start() completed with. The listener calls made from setListener's thread and from the I/O thread are serialized without holding a lock, so the listener never runs twice at once and gets the newest layout last. The producer applies its initial layout before setting its listener.
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
Problem. A V5 producer or consumer on a scalable topic can miss a split or merge that the broker pushes while
the producer or consumer is being created. Its
DagWatchClientstores every layout the broker pushes, butdelivered it only to a listener that was already set, and
setListenerdid not hand over what had arrivedbefore. The broker pushes each layout once.
CheckpointConsumerand aQueueConsumer(including the per-topic consumers of a namespaceQueueConsumer) apply the initial layout first and set the listener only once a reader or consumer exists forevery segment. A split or merge pushed in between is dropped: the consumer keeps reading the sealed parent,
never attaches to its children, and stops receiving the keys they cover until the broker happens to push a
layout again. With ZooKeeper, that is the next split or merge. With Oxia or RocksDB, which notify every
ancestor of a created metadata node, the broker also pushes the unchanged layout again when the first load
record of a new segment is written (within ~10 s); that rescues the consumer only if it comes after the
listener is set.
Producerset its listener before applying its initial layout, so a layout delivered in between wasoverwritten by the older initial one, and a producer created off the I/O thread after a newer layout arrived
missed it. Its sends then go to the sealed parent, are retried for ~4 s, and fail.
Consumers in a consumer group, and stream consumers, are not affected:
ScalableConsumerClient.setListeneralready replays the current assignment.
Example.
ScalableCheckpointConsumer.createUnmanagedAsync(ScalableQueueConsumer.createAsyncImpldoes the same):DagWatchClient:Modifications
Change.
DagWatchClient.setListenernow hands the listener the current layout when it is not the initial onethat
start()completed with, which the caller applies itself. The listener calls made fromsetListener'sthread and from the I/O thread are serialized with a work-in-progress counter, without holding a lock while the
listener runs: the listener never runs twice at once and gets the newest layout last. Once the listener is set,
layouts are delivered as before, one call per update on the I/O thread.
ScalableTopicProducernow applies its initial layout before setting its listener. The consumers are unchanged.Verifying this change
This change added tests and can be verified as follows:
V5LayoutChangeDuringCreateTestholds back the per-segment attaches with a v4 client subclass, splits thetopic, waits until the
DagWatchClienthas the split, then lets the creation finish. It checks which segmentsare attached by the time the creation completes, then that the messages arrive: an ungrouped checkpoint
consumer and a queue consumer attach the children, and a producer created with the pre-split layout sends to a
child. All three fail without the fix (the consumers attach only the sealed parent, and the producer's first
send goes to it) and pass with it. They check the attached segments rather than only waiting for messages
because the broker pushes the unchanged layout again when a node is created under the topic's metadata (here,
the first load record of the new segments, a few seconds later), which let a message-only version pass
without the fix.
DagWatchClientTestcovers the replay of a layout received beforesetListener, no replay when nothingarrived, and a layout arriving while the listener runs on another thread (delivered after it, never
concurrently, without blocking the I/O thread).
pulsar-client-v5unit tests and all the V5 tests inpulsar-broker(org.apache.pulsar.client.api.v5,org.apache.pulsar.client.impl.v5) pass on this branch merged with current master.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes