Skip to content

Flink streaming write: pre-write rollback adds blocking operation to the instant-time request path #20222

Description

@vamshipasunuru1

Bug Description

What happened:
#19902 established that StreamWriteOperatorCoordinator should serve instant-time requests in O(1),
and enumerated two blocking operations on that path: EventBuffers#awaitAllInstantsToCompleteIfNecessary()
(commit guard) and table-lock acquisition inside startCommit. #19960 followed by serving requests
asynchronously and having the task poll up to write.commit.ack.timeout, so a slow request no longer
breaks the coordination RPC.

Working in that same direction, we hit two remaining gaps when
hoodie.prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES (or FAILED_WRITES_CLEANER_POLICY=EAGER) is
set.

Gap 1 — a third blocking operation on the path, and it is work rather than a wait.

With pre-write rollback enabled, startCommit performs a full rollback of failed writes before it
returns:

StreamWriteFunction.snapshotState -> flushRemaining(false) [StreamWriteFunction.java:186]
-> instantToWrite(hasData()) [:515]
-> Correspondent.requestInstantTime
-> StreamWriteOperatorCoordinator.handleInstantRequest [StreamWriteOperatorCoordinator.java:421]
-> awaitAllInstantsToCompleteIfNecessary [:429] <- #19902, blocking #1
-> startInstant [:430, :538]
-> HoodieFlinkWriteClient.startCommit [BaseHoodieWriteClient.java:1150]
-> table lock <- #19902, blocking #2
-> rollbackFailedWrites <- gap 1

What you expected:

Pre-write rollback should be usable in the Flink writer without instant-time request latency scaling
with the number of files being rolled back.

Why this matters to us rather than simply leaving the config off: we currently rely on async cleans
running outside Flink to clean up failed writes, and enabling pre-write rollback is what surfaced the
above. The fallback works, but it overloads clean with a responsibility that is not its own. Clean and
rollback of failed writes are functionally different operations with different urgency. Clean reclaims
space and can tolerate running hours later. An uncleaned failed write actively blocks downstream
consumers: incremental readers such as DeltaStreamer block or fail on hollow commits, and a blocking
reader stalls silently, with every run reporting success and zero rows. Tying failed-write cleanup to
the clean schedule sets the latency of a correctness-affecting operation by a space-reclamation one.

Steps to reproduce:

  1. Set hoodie.prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES.
  2. Restart the job mid-write so an inflight commit is left behind.
  3. On restart, the first instant-time request triggers rollback of that inflight commit inside
    startCommit. JM will have as many as rollback.parallelism number of threads working on it.
  4. Observe checkpoint duration and JobManager thread count scaling with the marker count; on object
    storage, the request can exceed write.commit.ack.timeout and fail the checkpoint, leaving another
    inflight commit behind.

Acceptance criteria

  • With prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES, instant-time request latency is independent
    of the number of files in the instant being rolled back.
  • A restart leaving an inflight commit with ~100k markers is rolled back without checkpoint duration
    growing with marker count.

Environment

Hudi version: master
Query engine: (Spark/Flink/Trino etc) Flink
**Relevant configs:**prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES

Logs and Stack Trace

No response

Activity

  1. danny0405 commented on Oct 8, 2026

    @danny0405
    Contributor

    An uncleaned failed write actively blocks downstream
    consumers: incremental readers such as DeltaStreamer block or fail on hollow commits, and a blocking
    reader stalls silently, with every run reporting success and zero rows.

    we have switched to completion time based inc queries since 1.x, what release did you use then?

    instant-time request latency is independent
    of the number of files in the instant being rolled back.

    The rollback before instant creation is by-design for now, are you prososing to make it async and non-blocking? Making it async does not resolve the DeltaStreamer hollow commits block/fail issue right?

  2. vamshipasunuru1 commented on Oct 10, 2026

    @vamshipasunuru1
    ContributorAuthor

    we have switched to completion time based inc queries since 1.x, what release did you use then?
    consumer are on 0.14, we are doing a rolling upgrade.

    The rollback before instant creation is by-design for now, are you prososing to make it async and non-blocking?
    yes the proposal is to spin up rollback in async thread and do zero work within startInstant. Deletes can be throttled in rare cases so we would like the thread to be unbounded.

    Making it async does not resolve the DeltaStreamer hollow commits block/fail issue right?
    blocking downstream users is the desired behaviour, we don't want to change that.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    type:bugBug reports and fixes

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions