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:
- Set
hoodie.prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES.
- Restart the job mid-write so an inflight commit is left behind.
- 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.
- 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
Bug Description
What happened:
#19902 established that
StreamWriteOperatorCoordinatorshould 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 requestsasynchronously and having the task poll up to
write.commit.ack.timeout, so a slow request no longerbreaks the coordination RPC.
Working in that same direction, we hit two remaining gaps when
hoodie.prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES(orFAILED_WRITES_CLEANER_POLICY=EAGER) isset.
Gap 1 — a third blocking operation on the path, and it is work rather than a wait.
With pre-write rollback enabled,
startCommitperforms a full rollback of failed writes before itreturns:
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:
hoodie.prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES.startCommit. JM will have as many as rollback.parallelism number of threads working on it.storage, the request can exceed
write.commit.ack.timeoutand fail the checkpoint, leaving anotherinflight commit behind.
Acceptance criteria
prewrite.cleaner.policy=ROLLBACK_FAILED_WRITES, instant-time request latency is independentof the number of files in the instant being rolled back.
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