Repository navigation
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
Activity
please help review, thanks! @danny0405 @cshuo
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:
- Make the RPC handler a fast in-memory lookup/registration path that returns
READY(instant)orPENDING. It must not wait on prior commits, table locks, or an unfinished future. - Reuse the existing
instantRequestExecutoras the serialized creation worker, and serve status independently through thread-safe operation state. This avoids adding a third executor. - 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.
- 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. - 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.timeoutis 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 withgetPendingInstantsBefore(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.
Reacted by fhan- Make the RPC handler a fast in-memory lookup/registration path that returns
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.
Reacted by fhan
JIRA Umbrella Issue
Description
Problem
In Flink streaming/append writes, each write task obtains its instant time by sending an
InstantTimeRequesttoStreamWriteOperatorCoordinator. The coordinator handles this request on a single-threadedinstantRequestExecutor, and the handler performs two potentially long-blocking operations inline:EventBuffers#awaitAllInstantsToCompleteIfNecessary()→CommitGuard#blockFor(...)— blocks until all prior instants are committed (blocking-instant-generation mode).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 —startCommitblocks 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
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):
ready(instant)ornotReady(phase, progress); the write task'sCorrespondentpolls with exponential backoff + jitter and fails fast only when a monotonic progress token stalls for a whole no-progress window.startInstant()(lock acquisition) runs on a dedicatedinstantCreationExecutor, decoupled from the RPC-serving thread; its result is published back on the request thread.CommitGuard#blockFor) becomes a non-blocking gate on the request thread (getPendingInstantsBefore(cid).isEmpty()); the client polls, and normal commit completion innotifyCheckpointCompletereleases the gate.EventBuffersguarantees exactly one instant percheckpointIdand rejects events carrying a mismatched instant.Non-goals / Invariants preserved
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 onclose/ global failover / restore.Acceptance Criteria
rpc.ask.timeout, the coordination RPC does not time out, no task restart occurs, and the instant is returned once the lock is released.Sub-tasks → PR breakdown
HUDI-XXXX-1 —
[HUDI-XXXX] Harden EventBuffers to enforce one instant per checkpointChange Logs
Make the checkpoint→instant mapping authoritative:
initNewEventBufferusescompute+checkState(buffer == null);EventBuffer#addEventrejects events whose instant differs from the checkpoint's assigned instant (force JVM exit for non-bootstrap/non-endInput, else throw);getOrCreateBootstrapBuffervalidates 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 protocolChange Logs
Extend
Correspondent.InstantTimeResponsewithready/phase/progress(addInstantWaitPhase{PENDING_PRIOR_COMMIT, CREATING}; keepgetInstance(instant)for compatibility). RewriteCorrespondent#requestInstantTimeas a backoff+jitter poll loop that tracks phase+progress and fails fast on a whole no-progress window; handleInterruptedException. Thread poll config fromconfintoCorrespondent/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
readyin 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 StreamWriteOperatorCoordinatorChange Logs
Add
instantCreationExecutor(runs blockingstartInstant()off the request thread), anInstantCreationstate machine driven solely oninstantRequestExecutor, and a monotonicinstantGenProgresstoken.handleInstantRequestbecomes O(1): returnsready(idempotent),notReady(PENDING_PRIOR_COMMIT)(non-blocking replacement forCommitGuard#blockFor, blocking-mode only),notReady(CREATING)(creation in flight), or starts a new creation. Publish instant on the request thread (initNewEventBuffer+ complete waiters);failInstantCreation+resetFailedInstantCreationon failover/restore. Lifecycle:isClosingguard,waitForTasksFinish(true), orderedclose(), reject requests when closing. RetireCommitGuardfrom 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 requestsChange Logs
Add
instantRequestTimeoutSchedulerand suspend requests as waiters on the active creation, returningnotReady(CREATING)afterwrite.instant.request.long-poll.timeout.mswithout 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)
write.instant.request.poll.interval.max.mswrite.instant.request.no-progress.timeout.mswrite.instant.request.long-poll.timeout.msrpc.ask.timeout