Repository navigation
Conversation
sunchao
left a comment
There was a problem hiding this comment.
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
FetchFailedand 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.
|
@sam-1112 could you fix conflicts? |
…ler-failures-5646 # Conflicts: # spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala
Head branch was pushed to by a user without write access
|
@andygrove I have fixed the conflict. |
…ler-failures-5646 # Conflicts: # native/core/src/execution/operators/iceberg_write.rs
|
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
left a comment
There was a problem hiding this comment.
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')") |
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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()) |
There was a problem hiding this comment.
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() { |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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 callsbatchWrite.abortafter 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
cleanupOnAborttotrue, substantiating the existing missing-committed-file concern. Comparison withbranch-1.1separates inherited writer changes from these new hooks. This PR does not widen write eligibility or change configuration defaults. - Key design decisions: Require both
FetchFailedand an increased writer stage attempt to establish stage re-execution. Keep strict cleanup assertions for speculation while allowing only attributable rejected-attempt orphans afterSIGKILL. - 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.
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?
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.
FetchFailedand an increased writer stage attempt.