Skip to content

[core][spark] Deduplicate a replayed Structured Streaming micro-batch - #9667

Open
zhuxiangyi wants to merge 1 commit into
apache:masterfrom
zhuxiangyi:spark-streaming-idempotent-commit
Open

[core][spark] Deduplicate a replayed Structured Streaming micro-batch#9667
zhuxiangyi wants to merge 1 commit into
apache:masterfrom
zhuxiangyi:spark-streaming-idempotent-commit

Conversation

@zhuxiangyi

Copy link
Copy Markdown
Contributor

Purpose

Closes #9666.

Structured Streaming delivers exactly-once only if the sink is idempotent for a repeated batch id.
When a query fails after the sink returns from addBatch but before Spark records the batch as
completed, the restarted query replays that micro-batch with its original batch id.

PaimonSink received the batch id but used it only to pace full compaction, and committed through
table.newBatchWriteBuilder(), whose commit user is a fresh random UUID per builder and whose
commit identifier is always BatchWriteBuilder.COMMIT_IDENTIFIER = Long.MAX_VALUE. Neither of the
dimensions Paimon deduplicates on could identify a replay, so the batch was committed a second
time: every row duplicated in an append-only table, and silently wrong values in an aggregation
merge-engine table (writing (1, 10) and replaying that batch yields v = 20).

The machinery already exists in core and is what the Flink sink uses; the Spark sink simply took
the batch write path.

Tests

PaimonSinkIdempotencyTest (new, 7 cases), each asserting the correct behaviour so that it fails
without the fix:

  • a replayed micro-batch of an append-only table, driven end to end by deleting the commit log
    entry of the batch, which is exactly the checkpoint state a driver failure leaves behind;
  • the same when only the query id is available, i.e. the checkpoint location never reaches the sink
    options because it comes from spark.sql.streaming.checkpointLocation;
  • a replay of a batch that is not the first one of the query;
  • a replay in complete output mode;
  • write.stream.commit-user as an option of the writer and as a spark.paimon. session conf;
  • addBatch called twice with the same batch id through the API directly.

Two of them assert the prefix of the commit user recorded in the snapshot, so that the case which
is meant to exercise the query id derivation cannot pass through the checkpoint derivation.

The full set of Spark streaming suites was run on the spark3 and spark4 profiles (Spark 3.2, 3.4,
3.5, 4.1): 34 suites, 349 tests.

API and Format

BatchWriteBuilderImpl gains withCommitUser, and its newCommit() return type is narrowed from
BatchTableCommit to InnerTableCommit. The narrowing keeps the caller in the connector free of a
downcast that could only fail at runtime; it is source compatible, and BatchWriteBuilderImpl has
no subclasses.

InnerTableCommit gains checkFilesExistence(boolean). filterAndCommit verified that every file
it is about to commit still exists, which guards a committable restored from an engine's state that
may reference files deleted long ago. A caller filtering a committable it has just produced knows
those files exist, so the Spark sink turns the check off; otherwise every micro-batch would pay a
file listing proportional to the number of files it wrote. The default is unchanged, so Flink keeps
the check.

Documentation

docs/docs/spark/structured-streaming.md gains an "Exactly-once" section covering how the commit
user is derived, the new write.stream.commit-user option, and the limits: starting from a new
checkpoint location gives a query a new commit user, the data files of a skipped replay are left to
orphan file cleaning, and a postpone bucket table committing through the staged committer cannot
deduplicate and logs a warning per micro-batch. The generated option reference is regenerated.

Structured Streaming guarantees exactly-once only if the sink is idempotent
for a repeated batch id: when a query fails between the sink returning from
addBatch and Spark recording the batch as completed, the restarted query
replays that micro-batch with its original batch id.

PaimonSink received the batch id but only used it to pace full compaction,
and committed through newBatchWriteBuilder(), whose commit user is a fresh
random UUID and whose commit identifier is always Long.MAX_VALUE. Neither
can identify a replay, so the whole batch was committed a second time,
duplicating its rows in an append table.

Commit every micro-batch under a commit user that survives a restart, and
use filterAndCommit with the batch id as the commit identifier, so a replay
that Paimon already committed is skipped. The commit user is derived from
the checkpoint location, falling back to the query id Spark persists in the
checkpoint metadata when the location does not reach the sink options, and
write.stream.commit-user overrides both, as an option of the writer or as a
spark.paimon.write.stream.commit-user session conf, like the read side takes
its read.stream.* options.

filterAndCommit verified that every file it is about to commit still exists,
a check meant for a committable restored from an engine's state that may
reference files deleted long ago. A caller filtering a committable it has
just produced knows those files exist, so InnerTableCommit can now turn the
check off, and the Spark sink does. Otherwise every micro-batch would pay a
file listing proportional to the number of files it wrote.

The data files of a skipped replay stay uncommitted and are reclaimed by
orphan file cleaning. A postpone bucket table committing through the staged
committer cannot deduplicate and logs a warning per micro-batch.

@JingsongLi JingsongLi left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Two issues found in the checkpoint identity and commit maintenance lifecycle. The existing seven Spark 3 tests pass; additional checkpoint edge-case tests and a focused maintenance probe reproduce the issues below.

Comment on lines +58 to +60
checkpointLocation
.map(derivedCommitUser("checkpoint", _))
.orElse(queryId.map(derivedCommitUser("query", _)))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Prefer the persisted query ID over the checkpoint path

If a checkpoint is deleted and a new query starts at the same path, Spark assigns a new query ID and restarts batch IDs at 0. This code nevertheless reuses the previous commit user, so filterCommitted silently skips fresh batches whose IDs are at or below the previous query's last committed ID. I reproduced this with different input in the new query: the query completed successfully, but the table contained only the old row instead of both rows.

The reverse also fails: restarting the same checkpoint with only a trailing / added to its path changes the commit user despite an unchanged query ID. Replaying batch 0 then produced two rows instead of one.

Please keep the explicit override, but prefer the persisted Spark query ID for the default identity and use the checkpoint path only as a fallback. Add coverage for both a fresh query reusing a checkpoint path and an existing query using an equivalent path spelling.

Comment on lines +517 to +520
tableCommit
.checkFilesExistence(false)
.filterAndCommit(
Collections.singletonMap(Long.box(identifier), commitMessages.toList.asJava))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Preserve the batch maintenance lifecycle when filtering commits

Unlike commit(List), filterAndCommit does not set TableCommitImpl.batchCommitted. Maintenance therefore runs through the streaming executor wrapper, but this writer still closes and discards the committer immediately after each batch. With snapshot.expire.execution-mode=async, close() calls shutdownNow() and can interrupt snapshot expiration before it finishes. Even with the default synchronous mode, the wrapper catches maintenance exceptions and stores them for the next commit; because this instance is discarded, those failures are never propagated to the caller.

A focused probe against the built classes reproduced both the asynchronous interruption and the loss of synchronous error propagation. The existing testBatchWriteAsyncExpireFallbackToSync also establishes that a batch committer must finish maintenance before closing.

Please preserve the one-shot batch maintenance semantics in the filtered commit path, or explicitly wait for maintenance and propagate its failure before closing the committer.

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug] Spark Structured Streaming write commits a replayed micro-batch twice

2 participants