Skip to content

Flink streaming write: instant-time coordination RPC can time out and force task restarts when the instant lock is held (e.g. by cleaning) #19902

Description

@fhan688

JIRA Umbrella Issue

  • Summary: Flink streaming write: instant-time coordination RPC can time out and force task restarts when the instant lock is held (e.g. by cleaning)
  • Issue Type: Bug (Umbrella)
  • Priority: Major
  • Component/s: flink
  • Affects Version/s: 1.0.0, 1.1.0, master
  • Labels: flink, streaming-write, coordinator

Description

Problem

In Flink streaming/append writes, each write task obtains its instant time by sending an InstantTimeRequest to StreamWriteOperatorCoordinator. The coordinator handles this request on a single-threaded instantRequestExecutor, and the handler performs two potentially long-blocking operations inline:

  1. EventBuffers#awaitAllInstantsToCompleteIfNecessary() → CommitGuard#blockFor(...) — blocks until all prior instants are committed (blocking-instant-generation mode).
  2. startInstant() → HoodieFlinkWriteClient#startCommit(...) — acquires the table lock.

The write task obtains the result via gateway.sendRequestToCoordinator(...).get(), which is backed by a Flink coordination RPC with a finite ask timeout. When the lock is held by another operation — most commonly an async/inline clean or another table-service txn — startCommit blocks the request thread beyond the RPC timeout. The task's .get() then fails with a timeout and the pipeline restarts, even though nothing is actually wrong.

Root Cause

Instant-time request handling is not O(1): it blocks on (a) prior-commit completion and (b) table-lock acquisition on the same thread that must answer the coordination RPC within the ask timeout.

Impact

  • Spurious full-pipeline restarts under normal lock contention (cleaning, clustering/compaction scheduling, multi-writer OCC).
  • Restarts amplify under high write parallelism (more concurrent instant requests waiting on the one blocked thread).
  • Worse with tighter rpc.ask.timeout / lock providers with longer hold times.

Proposed Fix

Convert the single synchronous instant-time RPC into a poll-based protocol with an asynchronous coordinator state machine, so every coordination RPC returns in O(1):

  • The coordinator returns either ready(instant) or notReady(phase, progress); the write task's Correspondent polls with exponential backoff + jitter and fails fast only when a monotonic progress token stalls for a whole no-progress window.
  • The blocking startInstant() (lock acquisition) runs on a dedicated instantCreationExecutor, decoupled from the RPC-serving thread; its result is published back on the request thread.
  • The blocking "wait for prior commits" (CommitGuard#blockFor) becomes a non-blocking gate on the request thread (getPendingInstantsBefore(cid).isEmpty()); the client polls, and normal commit completion in notifyCheckpointComplete releases the gate.
  • Companion hardening in EventBuffers guarantees exactly one instant per checkpointId and rejects events carrying a mismatched instant.

Non-goals / Invariants preserved

  • Idempotency per checkpointId; ordering (new instant only after prior commits) in blocking mode; unchanged behavior in non-blocking mode; unchanged checkpoint state format; no in-flight creation state persisted; correct behavior on close / global failover / restore.

Acceptance Criteria

  • With the instant lock artificially held for longer than rpc.ask.timeout, the coordination RPC does not time out, no task restart occurs, and the instant is returned once the lock is released.
  • No regression in idempotency, ordering, or exactly-once commit semantics.

Sub-tasks → PR breakdown

HUDI-XXXX-1 — [HUDI-XXXX] Harden EventBuffers to enforce one instant per checkpoint

Change Logs
Make the checkpoint→instant mapping authoritative: initNewEventBuffer uses compute + checkState(buffer == null); EventBuffer#addEvent rejects events whose instant differs from the checkpoint's assigned instant (force JVM exit for non-bootstrap/non-endInput, else throw); getOrCreateBootstrapBuffer validates instant consistency.

Impact
Defensive only; prevents stale-instant events from entering the wrong buffer. No API/state change.

Risk level: low

Documentation Update: none

Independent of the async work; can merge first.

HUDI-XXXX-2 — [HUDI-XXXX] Poll-based instant-time coordination protocol

Change Logs
Extend Correspondent.InstantTimeResponse with ready/phase/progress (add InstantWaitPhase{PENDING_PRIOR_COMMIT, CREATING}; keep getInstance(instant) for compatibility). Rewrite Correspondent#requestInstantTime as a backoff+jitter poll loop that tracks phase+progress and fails fast on a whole no-progress window; handle InterruptedException. Thread poll config from conf into Correspondent / MockCorrespondent. New internal options: write.instant.request.poll.interval.max.ms, write.instant.request.no-progress.timeout.ms.

Impact
Wire-format of the coordination response changes; coordinator and tasks are same-version so no cross-version concern. Coordinator side still returns ready in this PR (behavior-preserving) — the poll loop simply short-circuits.

Risk level: low

Documentation Update: advanced config docs for the new internal options.

HUDI-XXXX-3 — [HUDI-XXXX] Make instant creation non-blocking in StreamWriteOperatorCoordinator

Change Logs
Add instantCreationExecutor (runs blocking startInstant() off the request thread), an InstantCreation state machine driven solely on instantRequestExecutor, and a monotonic instantGenProgress token. handleInstantRequest becomes O(1): returns ready (idempotent), notReady(PENDING_PRIOR_COMMIT) (non-blocking replacement for CommitGuard#blockFor, blocking-mode only), notReady(CREATING) (creation in flight), or starts a new creation. Publish instant on the request thread (initNewEventBuffer + complete waiters); failInstantCreation + resetFailedInstantCreation on failover/restore. Lifecycle: isClosing guard, waitForTasksFinish(true), ordered close(), reject requests when closing. Retire CommitGuard from the instant-request path.

Impact
Core fix. Removes lock-induced RPC timeouts. Preserves idempotency/ordering; non-blocking mode unchanged; checkpoint state format unchanged; in-flight creation not persisted.

Risk level: medium

Documentation Update: none (internal behavior)

Tests: fault-injecting test lock provider that holds the lock past rpc.ask.timeout, asserting no RPC timeout / no restart / eventual success; coordinator concurrency, idempotency, failure-retry, close-cancellation tests.

HUDI-XXXX-4 (optional) — [HUDI-XXXX] Server-side long-poll for instant-time requests

Change Logs
Add instantRequestTimeoutScheduler and suspend requests as waiters on the active creation, returning notReady(CREATING) after write.instant.request.long-poll.timeout.ms without cancelling creation. Reduces tail latency and steady-state request rate.

Impact
Latency optimization only. The long-poll window MUST be strictly less than the coordination RPC timeout (rpc.ask.timeout, default 10s) — default conservatively (e.g. 0 = disabled, or 8000ms) and document the constraint.

Risk level: low


Dependency order: 1 (independent) → 2 → 3 → 4 (optional). 1 can proceed in parallel with 2/3.

New configuration options (internal use)

Option Suggested default Purpose
write.instant.request.poll.interval.max.ms 1000 Client polling backoff cap
write.instant.request.no-progress.timeout.ms checkpoint-timeout order Fail fast only when the progress token stalls a whole window
write.instant.request.long-poll.timeout.ms 0 (disabled) or 8000 Server-side long-poll window; MUST be < rpc.ask.timeout

Activity

  1. fhan688 commented on Sep 11, 2026

    @fhan688
    ContributorAuthor

    please help review, thanks! @danny0405 @cshuo

  2. danny0405 commented on Sep 11, 2026

    @danny0405
    Contributor

    The root-cause analysis makes sense. I checked the coordinator and EventBuffers on current master (a46909b09bbb). I would suggest a smaller first fix based on the standard asynchronous request–reply pattern: start once, return PENDING promptly, and poll the same operation until it is READY.

    The key distinction is that instant creation is already offloaded to an executor, but the coordination RPC still waits for it to finish. Moving creation to another executor alone does not fix that coupling.

    A minimal design could be:

    1. Make the RPC handler a fast in-memory lookup/registration path that returns READY(instant) or PENDING. It must not wait on prior commits, table locks, or an unfinished future.
    2. Reuse the existing instantRequestExecutor as the serialized creation worker, and serve status independently through thread-safe operation state. This avoids adding a third executor.
    3. Atomically register one operation per checkpoint within the current coordinator generation before submitting creation. Concurrent writers and repeated polls share that operation; polling must never enqueue duplicate creation.
    4. Keep the prior-commit wait and startInstant() on the worker. Install the event buffer before publishing READY, and preserve the checkpoint-to-instant mapping so a lost response does not cause another instant to be created.
    5. Let writers poll with capped backoff and jitter, with one RPC outstanding per writer and a bounded overall wait budget. Transient transport retries must reuse the same checkpoint identity; creation failures should remain terminal and propagate through the normal failure path.

    I would leave phase/progress tokens and server-side long polling out of the initial fix. A healthy operation waiting for a cleaner's lock may show no progress at all, while unrelated activity could advance a global progress token without helping this request. A monotonic overall deadline is simpler and easier to reason about. The existing write.commit.ack.timeout is worth evaluating as the wait budget, checking its suitability for both blocking and non-blocking modes; PENDING responses must not reset that deadline.

    A few correctness details still need explicit coverage:

    • Ordering: If retaining CommitGuard on the worker, wait in a predicate loop with the remaining deadline. The current blockFor(String) only waits for a signal; a wakeup does not establish that every required prior commit completed. The supplier-based overload is a useful starting point. Also, replacing the existing gate with getPendingInstantsBefore(cid).isEmpty() changes how empty buffers are treated, so that should be validated rather than assumed behavior-preserving.
    • Recovery/close: Fence results by coordinator generation and coordinate recovery with the running worker. Cancelling a future or ignoring its result alone does not stop startInstant() from mutating shared clients or the timeline. Prevent old work from racing with restored state or closed clients.
    • Failure/idempotency: Do not automatically retry creation after a potentially partial failure. Retire operation records with checkpoint lifecycle handling, and reject stale requests rather than interpreting them as new work.
    • Validation: Reject mismatched checkpoint/instant events through normal coordinator/job failure handling; introducing JVM exit seems unnecessary for this fix.

    The main integration test should hold the lock longer than the configured RPC timeout but shorter than the operation/checkpoint budgets, and verify prompt polling responses, no restart, eventual success, and exactly one instant across concurrent writers. Lost responses, creation failure, and recovery during creation are the other essential cases.

    This addresses lock-induced RPC timeouts while keeping the existing serialized creation model. It cannot guarantee success when contention exceeds the lock-acquisition or checkpoint deadline; those remain legitimate failure boundaries.

    This is a source-based design review, not a tested patch.

  3. cshuo commented on Sep 11, 2026

    @cshuo
    Collaborator

    This is a valid issue: the RPC can time out while instant creation waits for the lock, causing the writer task to fail and potentially triggering a job restart. Flink’s RPC ask timeout defaults to 10 seconds (pekko.ask.timeout; see the configuration definition), so lock contention lasting longer than that can trigger this failure.

    Before introducing a more complex state machine, could we first try bounded client-side retries with backoff and jitter? Since requests are processed serially and reuse the checkpoint-to-instant mapping, this could address transient contention with a smaller change.

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

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions