Skip to content

[fix][client] Apply layout changes that arrive while a V5 producer or consumer is created - #26854

Open
merlimat wants to merge 2 commits into
apache:masterfrom
merlimat:mmerli/checkpoint-consumer-missed-layout-update
Open

merlimat wants to merge 2 commits into
apache:masterfrom
merlimat:mmerli/checkpoint-consumer-missed-layout-update

Conversation

@merlimat

Copy link
Copy Markdown
Contributor

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 DagWatchClient stores every layout the broker pushes, but
delivered it only to a listener that was already set, and setListener did not hand over what had arrived
before. The broker pushes each layout once.

  • An ungrouped CheckpointConsumer and a QueueConsumer (including the per-topic consumers of a namespace
    QueueConsumer) apply the initial layout first and set the listener only once a reader or consumer exists for
    every 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.
  • A Producer set its listener before applying its initial layout, so a layout delivered in between was
    overwritten 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.setListener
already replays the current assignment.

Example.

CompletableFuture<CheckpointConsumer<String>> creating = client.newCheckpointConsumer(Schema.string())
        .topic("topic://tenant/ns/orders")
        .startPosition(Checkpoint.earliest())
        .createAsync();
// While the consumer is still creating the reader of segment 0:
admin.scalableTopics().splitSegment("topic://tenant/ns/orders", 0);
CheckpointConsumer<String> consumer = creating.get();
// Messages sent after the split never reach the consumer.

ScalableCheckpointConsumer.createUnmanagedAsync (ScalableQueueConsumer.createAsyncImpl does the same):

return consumer.applyAssignment(allSegmentsOf(initialLayout))   // the split arrives here...
        .thenApply(__ -> {
            dagWatch.setListener((newLayout, oldLayout) -> ...);   // ...but the listener is set only now

DagWatchClient:

// onUpdate
ClientSegmentLayout oldLayout = currentLayout.getAndSet(newLayout);   // the split is stored
LayoutChangeListener l = listener;
if (l != null) {                                                      // null while the consumer is created
    l.onLayoutChange(newLayout, oldLayout);
}

void setListener(LayoutChangeListener listener) {
    this.listener = listener;                                         // nothing hands the split over later
}

Modifications

Change. DagWatchClient.setListener now hands the listener the current layout when it is not the initial one
that start() completed with, which the caller applies itself. The listener calls made from setListener's
thread 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.

void setListener(LayoutChangeListener listener) {
    this.listener = listener;
    notifyListener();
}

private void notifyListener() {
    if (pendingNotifications.getAndIncrement() != 0) {
        return;
    }
    do {
        LayoutChangeListener l = listener;
        ClientSegmentLayout layout = currentLayout.get();
        ClientSegmentLayout previous = notifiedLayout;   // starts as the initial layout
        if (l != null && layout != previous) {
            notifiedLayout = layout;
            l.onLayoutChange(layout, previous);          // exceptions are caught and logged
        }
    } while (pendingNotifications.decrementAndGet() != 0);
}

ScalableTopicProducer now applies its initial layout before setting its listener. The consumers are unchanged.

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • V5LayoutChangeDuringCreateTest holds back the per-segment attaches with a v4 client subclass, splits the
    topic, waits until the DagWatchClient has the split, then lets the creation finish. It checks which segments
    are 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.
  • DagWatchClientTest covers the replay of a layout received before setListener, no replay when nothing
    arrived, and a layout arriving while the listener runs on another thread (delivered after it, never
    concurrently, without blocking the I/O thread).
  • The pulsar-client-v5 unit tests and all the V5 tests in pulsar-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

  • 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

… 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.
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.

1 participant