Skip to content

test: cover speculation and executor loss in multi-task native Iceberg writes - #6601

Open
sam-1112 wants to merge 8 commits into
apache:mainfrom
sam-1112:test/iceberg-scheduler-failures-5646
Open

sam-1112 wants to merge 8 commits into
apache:mainfrom
sam-1112:test/iceberg-scheduler-failures-5646

Conversation

@sam-1112

@sam-1112 sam-1112 commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes part of #5646(speculative execution and executor-loss scenarios for native Iceberg writes).

Rationale for this change

Local write tests do not exercise Spark's selection of speculative task results or recovery after an executor is lost during a native write.

These tests assert both the visible table outcome and storage state, including files that remain unreferenced and are therefore invisible to readers.

What changes are included in this PR?

  • Add a manual multi-worker integration suite covering:
    • Original-attempt and speculative-attempt winner directions.
    • Executor-loss task replacement.
    • Shuffle-output-loss stage re-execution.
  • Add explicitly scoped JVM/native probes with bounded gates and atomic per-attempt markers. Executor termination occurs after native write progress and physical file creation, before writer close or payload handoff.
  • Record scheduler events, driver-accepted task-message file sets, and successful commit markers.
  • Assert the exact 48,000-row ID multiset, one new snapshot, manifest equality with accepted winner files, and exclusion of rejected-attempt files from manifest references.
  • Compare physical files with manifest references and export orphan/missing file audits. Executor-loss scenarios allow only orphans belonging to known rejected attempts; unknown orphans and missing referenced files fail. Speculation retains strict cleanup assertions.
  • Add a local probe smoke test, a Linux/SSH executor-termination helper, and documentation for the external-cluster requirements.
  • Exclude the manual integration suite from ordinary single-host CI discovery.

This PR adds test instrumentation and coverage; it does not implement an orphan-cleanup mechanism.

How are these changes tested?

Tests ran on Linux ARM64 through Docker Desktop on a Mac mini M4, using Spark 4.1.3, Java 17, Iceberg 1.11.0, two Spark Standalone workers, and a shared Docker volume.

  • Both speculative winner-direction tests passed (2 passed, 0 failed).
  • Both executor-loss tests passed (2 passed, 0 failed). Task replacement recovered on another executor; stage re-execution recorded FetchFailed and an increased writer stage attempt.
  • The mid-write retry, post-native cleanup, commit-failure, post-native retry-success, and scheduler-probe smoke tests each passed locally.
  • The executor-loss tests verified rows, snapshot count, accepted files against manifest references, and isolation of rejected-attempt files. Each storage audit found one known rejected-attempt orphan, zero unknown orphans, and zero missing referenced files. Storage evidence was captured before fixture teardown.

@github-actions github-actions Bot added enhancement New feature or request test Testing related area:writer Native Parquet writer area:Iceberg labels Oct 4, 2026
@sam-1112
sam-1112 marked this pull request as ready for review October 4, 2026 16:34

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Local Iceberg failure tests did not exercise speculative winner selection or recovery after executor and shuffle-output loss.
  • Design approach: Add a manual multi-host suite with scoped, serializable probes and per-attempt evidence. Check scheduler events, exact rows, snapshots, accepted files, manifests and physical storage together.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked relevant scheduler and write semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Existing cleanup ownership remains intact. The supplied discussion snapshot contains no reviews, comments or threads.
  • Key design decisions: Preserve strict cleanup assertions for speculation while permitting only attributable rejected-attempt orphans after executor termination. Require FetchFailed and an increased stage attempt to establish stage re-execution.
  • Implementation sketch: An optional RDD wrapper gates input, the native probe records write progress, and JVM callbacks record handoff, accepted messages and commit completion. Inactive probes add checks but no marker I/O or polling. Performance was not benchmarked.
  • Behavioral changes worth calling out: The multi-host suite is intentionally excluded from ordinary CI discovery. Compared touched production files with branch-1.1. This PR preserves write eligibility, defaults and commit/cleanup semantics. Timestamp-partition differences from that release already exist in the supplied base.
  • Suggested improvements: None supported by reproducible P1/P2 evidence.

Reviewed the entire 13-file diff from d98fd2494c9a6cc8552efe6aa0a73d5714486b23 through 45f2ad0ef7f073187ab69a0da3eb1601e8ac54e2, including all branch commits. The PR is not a draft. Routed skills: review-comet-pr, review-comet-iceberg-write-pr and review-comet-shuffle-pr.

Exact-head CI: Only the successful label check was visible. No build or test verdict was available, including the required Iceberg regression verdict.

Validation limits: Suite registration, protobuf generation, Python syntax and git diff --check passed. Native/JVM builds and runtime suites were not run. The required multi-host cluster is not configured here, so the author's scheduler results were not independently reproduced.

@andygrove
andygrove enabled auto-merge October 4, 2026 21:14
@andygrove

Copy link
Copy Markdown
Member

@sam-1112 could you fix conflicts?

…ler-failures-5646

# Conflicts:
#	spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala
auto-merge was automatically disabled October 5, 2026 02:46

Head branch was pushed to by a user without write access

@sam-1112

sam-1112 commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

@andygrove I have fixed the conflict.

…ler-failures-5646

# Conflicts:
#	native/core/src/execution/operators/iceberg_write.rs
@sam-1112

sam-1112 commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

I’ll be away for military service until Oct 30 and will only be available on Oct 17 and Oct 24. I’ll continue working on this PR, but responses and updates may be slower during this period.Thanks for your understanding.

…ler-failures-5646

# Conflicts:
#	native/core/src/execution/operators/iceberg_write.rs

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for taking on the scheduler part of #5646. The storage audit against the manifests is careful, and the post-handoff retry test is a good addition. I'm requesting changes because the two executor-loss scenarios can't pass on the current head, and one of the probe hooks can abort a write that already committed. Details are inline.

I checked this on the head merged onto main, in local mode, using the suite's table properties and row counts with your SharedIcebergSchedulerProbe. I also tried a probe whose committed() throws, on Spark 4.1 and on Spark 3.4.

Could you also add a row for CometIcebergSchedulerFailureSuite to the Testing table in iceberg-writes.md, pointing at the new page? That table is where people look for the write suites.

The branch conflicts with main again in IcebergCommitExec, since #6693 added a command argument to IcebergWriteSummaryShim.commit, so it will need another merge of main.

spark.sql(
s"CREATE TABLE $table (id INT) USING iceberg " +
"TBLPROPERTIES ('write.distribution-mode'='none', " +
"'write.target-file-size-bytes'='1')")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I reproduced the executor-loss setup in local mode on the current head merged onto main, and it can't get past the check before the kill. Since #6241 landed on October 5, the native writer holds back each partition's first data page, 20,000 rows by default, before it opens any file. This table has no write.parquet.page-row-limit, and each task writes 12,000 rows. So all four native gates fired at rows=2000 with an empty file list, and assert(target.rows >= 1000L && target.files.nonEmpty) fails. With 'write.parquet.page-row-limit'='1000', the same change #6241 made to the mid-write retry test, each gate saw two files and the first was already on disk. Could you add that, update the "Executor process loss" section of the doc page to match, and re-run both loss scenarios on your cluster? Your earlier runs predate #6241.

}
}

test("native acceleration: scheduler probe records handoff and accepted files") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This test runs in mode smoke, so beforeNative() returns None and CI never drives the native gate. That's why nothing caught the #6241 change. If the native gate stays, could we add a local test that uses an executor-loss probe with release-all touched up front, writes a few thousand rows at batch size 1000, and asserts that native-progress.json lists files with bytes on disk? I tried one locally and it runs in about a second.

try {
messages.foreach(batchWrite.onDataWriterCommit)
IcebergWriteSummaryShim.commit(batchWrite, messages, child)
schedulerProbe.foreach(_.committed())

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

committed() runs inside the try whose catch aborts the write. If the probe throws after a successful commit, we call batchWrite.abort on messages that are already in a snapshot. I tried a probe whose committed() throws. On Spark 4.1 the INSERT failed even though the snapshot committed. On Spark 3.4 with Iceberg 1.5.2, where SparkWrite.cleanupOnAbort starts out true, the abort also deleted the committed data file, so the table pointed at a missing file. Your probe's Files.createFile throws on a duplicate marker by design, so the suite can reach this. Could the call move below the try/catch, once the commit has succeeded?

writer.write(decorated, &properties).await?;
timer.done();
reservation.try_resize(open_files.bytes() + writer.pending_bytes())?;
if let Some(probe) = common.test_probe.as_ref() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I'm wary of adding a test-only message to the wire format, and a gate that polls the filesystem inside the native write loop, for a suite that CI doesn't run. The page-row-limit problem shows how quietly it can drift. The native writer pulls its JVM input one batch at a time through the Arrow C stream from CometArrowStream.inputObjects, so holding that iterator after N rows parks the writer at the same point as this gate. Could the mid-write pause live in IcebergSchedulerProbeRDD instead? The attempt's files can be found by the attempt-unique prefix from file_name_prefix, which iceberg-writes.md already documents. That would remove the proto field and the native changes entirely.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Reviewed the entire 13-file diff from 8c783aa88104616dcf0f7876b4f7a9ba71e2bf31 to de7eb9c2d98c8a1f2468f8fac92ae3f9f968c4be, including all branch commits. The PR is not a draft. Routed skills: review-comet-pr, review-comet-iceberg-write-pr, and review-comet-shuffle-pr.

No additional introduced P1/P2 issues found within this review. Two substantiated existing correctness blockers remain unresolved. These are acknowledged below without duplicating their inline findings.

Summary

  • Prior state and problem: Local Iceberg write tests did not establish speculative winner selection or recovery after executor and shuffle-output loss. Successful reads alone also could not detect orphan files.
  • Design approach: Add a manual multi-worker suite with serializable probes, bounded gates, atomic attempt markers, scheduler events, and storage audits. Correlate accepted task messages with manifest references and physical files.
  • Correctness: The exact ID multiset, single-snapshot, winner-file equality, and rejected-file exclusion assertions provide useful independent checks. However, the existing executor-loss setup concern remains: the gate pauses after 2,000 rows while the dictionary feed holds up to 20,000 rows before opening files. The existing post-commit abort concern also remains: committed() can throw inside the catch region that calls batchWrite.abort after commit succeeds.
  • Compatibility analysis: Compared scheduler result acceptance, speculative cancellation, fetch-failure recovery, and V2 commit handling against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 sources. Checked Iceberg 1.5.2, 1.8.1, 1.10.0, and 1.11.0 abort behavior. Iceberg 1.5.2 initializes cleanupOnAbort to true, substantiating the existing missing-committed-file concern. Comparison with branch-1.1 separates inherited writer changes from these new hooks. This PR does not widen write eligibility or change configuration defaults.
  • Key design decisions: Require both FetchFailed and an increased writer stage attempt to establish stage re-execution. Keep strict cleanup assertions for speculation while allowing only attributable rejected-attempt orphans after SIGKILL.
  • Implementation sketch: An optional RDD wrapper gates parent iterator construction. The protobuf probe gates native write progress. JVM callbacks record handoff, accepted file sets, and commit completion. The local retry test checks successful recovery and removal of failed-attempt files.
  • Performance: Inactive probes add checks without marker I/O or filesystem polling. Active probes deliberately introduce waits and diagnostics. No evidence-backed P1/P2 performance regression was identified. Throughput and overhead were not benchmarked.
  • Design: Combining scheduler evidence with logical and physical storage checks makes the intended guarantees clear. Disabling Comet shuffle and preserving ordinary Spark shuffle lineage isolates the stage-loss scenario. The current buffering interaction prevents the executor-loss tests from reaching their intended injection point.
  • Abstraction & complexity: The small serializable interface keeps executor state independent of the driver singleton. The main maintenance cost is the test-only protocol spanning Scala, protobuf, and Rust. The existing discussion's input-RDD gate proposal addresses that coupling and does not warrant another duplicate comment.
  • Behavioral changes worth calling out: The multi-host suite is intentionally excluded from ordinary CI discovery. The regular smoke test exercises JVM probe plumbing but never activates the native progress gate. This PR adds attribution and auditing, not orphan cleanup after executor termination.
  • Suggested improvements: Resolve the existing gate-configuration and post-commit callback threads, add the already-proposed local native-progress regression, and rerun both loss scenarios on the current head. Obtain the required Iceberg regression verdict. No additional recommendations met the reproducible P1/P2 reporting bar.

Exact-head CI: Label pull requests passed. Comet CI and CodeQL report action_required. Comet CI has zero jobs, so no build, test, or Iceberg regression verdict is available.

Validation: Python syntax, suite registration, protobuf generation, generated Java compilation, and git diff --check passed. The new probe/support Scala sources compiled against Spark 4.1.3. A bounded four-task local Spark check passed probe serialization, gate configuration, attempt markers, scheduler events, accepted markers, duplicate-commit rejection, and scope reset.

Validation limits: That focused check does not execute the native writer. Full native/JVM builds, native write suites, upstream Iceberg regression tests, and multi-host speculation/executor-loss scenarios were not run. The required external cluster is not configured here, and the author's earlier cluster results were not independently reproduced on this head.

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

Labels

area:Iceberg area:writer Native Parquet writer enhancement New feature or request test Testing related

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants