From 738a0319bda833ecfe6189afeba03a3d631abe56 Mon Sep 17 00:00:00 2001 From: Unik Dahal <61407386+unikdahal@users.noreply.github.com> Date: Sun, 4 Oct 2026 00:47:57 +0530 Subject: [PATCH] fix Iceberg cleanup handoff --- .../review-comet-iceberg-write-pr/SKILL.md | 8 +- .../contributor-guide/iceberg-writes.md | 32 +++--- .../user-guide/latest/iceberg-writes.md | 21 ++-- .../src/execution/operators/iceberg_write.rs | 103 ++++++++++++++++-- .../sql/comet/CometIcebergWriteExec.scala | 16 +++ .../comet/CometIcebergWriteActionSuite.scala | 73 +++++++++++++ 6 files changed, 216 insertions(+), 37 deletions(-) diff --git a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md index 9ad714e44fd..9752f865a5d 100644 --- a/.ai/skills/review-comet-iceberg-write-pr/SKILL.md +++ b/.ai/skills/review-comet-iceberg-write-pr/SKILL.md @@ -123,9 +123,11 @@ only the latest. ## 4. Failure Handling and Cleanup -Files must have exactly one owner at every moment. Read the ownership table in the contributor -guide before reviewing any change near `AbortOnDrop`, `TrackingLocationGenerator`, -`WrittenFileCleanup`, `drainNativePayload` or `IcebergCommitExec.collectAndCommit`. +Cleanup must never have an ownership gap. During the handoff, native remains armed until the JVM +has taken the locations and acknowledged that by polling EOF, so a brief overlap is intentional. +Read the ownership table in the contributor guide before reviewing any change near `AbortOnDrop`, +`TrackingLocationGenerator`, `WrittenFileCleanup`, `drainNativePayload` or +`IcebergCommitExec.collectAndCommit`. - [ ] A new failure point between writing a file and the JVM taking the locations is covered by the native guard, including the path where the plan is dropped mid-write rather than returning an diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index 3214bfed178..8f52afc251c 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -272,19 +272,25 @@ discards them. The price is one footer read per written file. ## Failure Handling and Cleanup Ownership A failed attempt must leave no data files behind, as iceberg-java's `DataWriter.abort()` does, and a -failed job must not commit anything. The files are always owned by exactly one side: - -| Phase | Owner | Mechanism | -| ------------------------------------------------------------------- | ---------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| Writing, closing writers, encoding the manifest, building the batch | Native | `TrackingLocationGenerator` records every location. `AbortOnDrop` deletes them on an error, and also when the plan is dropped mid-write, from inside or outside a Tokio runtime. | -| After the batch reaches the JVM, until the task succeeds | JVM task | `WrittenFileCleanup`, a task failure listener registered before the payload is read, takes the locations before the manifest is decoded. | -| Job failure after some tasks completed | Driver (`IcebergCommitExec`) | Calls `BatchWrite.abort` with the completed messages, then deletes their files through the table's `FileIO`. | -| Commit failure | Iceberg | `SparkWrite.abort` on the genuine `TaskCommit` messages, the same as the stock path. | - -The handoff between the first two rows is why the payload carries the locations separately from the -manifest: a failure decoding the manifest would otherwise lose the list of files to delete. All -deletion is best effort and logged. It must never replace the original exception, and anything it -misses is unreferenced and reclaimed by Iceberg's `remove_orphan_files`. +failed job must not commit anything. Cleanup must never have an ownership gap. During the handoff, +the native guard intentionally stays armed after the JVM takes the locations and is disarmed only +when the JVM polls the output stream to EOF: + +| Phase | Owner | Mechanism | +| ------------------------------------------------------------------- | ---------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------- | +| Writing, closing writers, encoding the manifest, building the batch | Native | `TrackingLocationGenerator` records every location. `AbortOnDrop` deletes on error or plan drop and stays armed after yielding the batch. | +| JVM has recorded locations, before its final EOF poll | Native + JVM task | `WrittenFileCleanup` owns the decoded locations while `AbortOnDrop` is still armed. A failure may make both sides attempt best-effort deletion. | +| After the JVM's EOF poll, until the task succeeds | JVM task | The EOF poll disarms `AbortOnDrop`; `WrittenFileCleanup` remains registered for the rest of the task. | +| Job failure after some tasks completed | Driver (`IcebergCommitExec`) | Calls `BatchWrite.abort` with the completed messages, then deletes their files through the table's `FileIO`. | +| Commit failure | Iceberg | `SparkWrite.abort` on the genuine `TaskCommit` messages, the same as the stock path. | + +The handoff between the first two rows intentionally overlaps: `WrittenFileCleanup` takes the +locations before the final EOF poll, and that poll is the acknowledgement that lets `AbortOnDrop` +disarm. Deleting the same file from both sides during a failure in that narrow window is harmless +because cleanup is best effort. The payload carries the locations separately from the manifest so a +failure decoding the manifest cannot lose the list of files to delete. Cleanup is logged and must +never replace the original exception; anything it misses is unreferenced and reclaimed by Iceberg's +`remove_orphan_files`. The same planning-time rule applies to failures: a failure that happens because the native side rejected something the gate admitted is a bug in the gate, even if cleanup works. diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index d1b14c1e718..af7a5b94501 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -246,16 +246,17 @@ Partial results are never committed. The commit set is exactly the commit messag successful tasks — a failed task contributes none — and if the job fails, the driver-side commit operator aborts without committing anything. A failed task attempt also deletes the data files it created, as iceberg-java's writer abort does. The native writer records every -location it hands to a file writer, and exactly one side owns deleting them at any moment: the -native writer owns them until its output batch reaches the JVM (so it cleans up a failed write, -a task torn down before the write completed — for example because the operator feeding it threw -— and a failure encoding the manifest or building that batch), and the JVM owns them from then -on through a task failure listener. The handoff does not depend on decoding the manifest: the -native side reports the locations in the output batch alongside it, and the listener is handed -them before the manifest is decoded, so a failure in that decode still cleans up. Both -deletions are best-effort and never mask the original failure; anything they miss is invisible -to every reader, since readers resolve files through committed manifests only, and is reclaimed -by Iceberg's normal `remove_orphan_files` maintenance. +location it hands to a file writer, and cleanup has no ownership gap. Native keeps its cleanup +guard armed after yielding the output batch. The JVM reads those locations into a task failure +listener first, then polls the native output to EOF; that EOF is the acknowledgement that lets +native disarm. A failure in this narrow handoff window can therefore trigger best-effort deletion +on both sides, which is harmless. The native guard still covers a failed write, a task torn down +before the write completed — for example because the operator feeding it threw — and a failure +encoding the manifest or building the output batch. The handoff does not depend on decoding the +manifest: the locations are read before the manifest is decoded, so a failure in that decode still +cleans up. Cleanup never masks the original failure; anything it misses is invisible to every +reader, since readers resolve files through committed manifests only, and is reclaimed by Iceberg's +normal `remove_orphan_files` maintenance. When one task fails, the tasks that had already completed leave committed-nothing data files too. The committer collects each task's commit message as that task finishes, so on a job diff --git a/native/core/src/execution/operators/iceberg_write.rs b/native/core/src/execution/operators/iceberg_write.rs index 52138702d6f..ac2357966ca 100644 --- a/native/core/src/execution/operators/iceberg_write.rs +++ b/native/core/src/execution/operators/iceberg_write.rs @@ -125,13 +125,13 @@ impl LocationGenerator for TrackingLocationGenerator { } } -/// Deletes the tracked files if the write task is dropped before it finished. +/// Deletes the tracked files if the write task is dropped before the JVM acknowledges its output. /// /// A task can end without its future ever observing an error: when the JVM-side input iterator /// throws, `executePlan` returns that error straight from the JNI batch pull and the JVM then /// releases the plan, dropping this future mid-flight. The guard turns that drop into the same -/// cleanup the explicit error path performs. It stays armed until the task's output batch has -/// been handed to the JVM, which is the point where the JVM takes over cleanup ownership. +/// cleanup the explicit error path performs. It stays armed until the JVM polls past the output +/// batch, after recording the locations in its task-failure listener. struct AbortOnDrop { file_io: FileIO, generator: TrackingLocationGenerator, @@ -198,6 +198,28 @@ impl Drop for AbortOnDrop { } } +/// The JVM reads the locations from the single output batch, registers its cleanup listener, +/// and then polls once more to verify that the stream ended. Keep native cleanup armed through +/// that last poll: a cancelled task or a failure decoding the locations before then must still +/// delete the files, even though the output batch was successfully constructed. +fn output_with_cleanup_ack( + batch: RecordBatch, + abort_guard: AbortOnDrop, +) -> impl futures::Stream> + Send { + futures::stream::unfold( + (Some(batch), abort_guard), + |(batch, mut abort_guard)| async move { + match batch { + Some(batch) => Some((Ok::<_, DataFusionError>(batch), (None, abort_guard))), + None => { + abort_guard.disarm(); + None + } + } + }, + ) +} + /// Best-effort deletion of every file a failed task attempt created, the native counterpart of /// iceberg-java's `SparkCleanupUtil.deleteTaskFiles`. Failures are logged rather than returned: /// the original task failure must stay the one Spark reports, and anything left behind is still @@ -401,12 +423,9 @@ impl ExecutionPlan for IcebergWriteExec { } .await; match packaged { - // The batch carries the locations, and the JVM takes cleanup ownership of them - // before it decodes the manifest, so the guard's job is done. - Ok(batch) => { - abort_guard.disarm(); - Ok::<_, DataFusionError>(futures::stream::iter(vec![Ok(batch)])) - } + // The JVM registers the locations before polling for EOF. Keep the guard armed + // until that poll, so dropping the stream during the handoff still cleans up. + Ok(batch) => Ok::<_, DataFusionError>(output_with_cleanup_ack(batch, abort_guard)), Err(e) => { abort_guard.abort().await; Err(e) @@ -454,7 +473,7 @@ impl DisplayAs for IcebergWriteExec { /// depending on `writer_mode`. /// /// On success the still-armed [`AbortOnDrop`] is returned along with the data files: the caller -/// owns cleanup until the output batch has been handed to the JVM. +/// owns cleanup until the JVM acknowledges the output after recording its locations. #[allow(clippy::too_many_arguments)] async fn run_write_task( mut input: SendableRecordBatchStream, @@ -1770,6 +1789,7 @@ mod tests { CompressionCodec as ProtoCodec, IcebergParquetWriteSettings, IcebergWriteCommon, IcebergWriterMode as ProtoIcebergWriterMode, }; + use futures::StreamExt; use iceberg::spec::{ Manifest, NestedField, PartitionSpec, PrimitiveType, Schema, Transform, Type, }; @@ -1777,6 +1797,7 @@ mod tests { use std::collections::HashMap; use std::path::PathBuf; use std::sync::Arc; + use std::time::Duration; use tempfile::TempDir; fn user_schema() -> SchemaRef { @@ -1974,7 +1995,7 @@ mod tests { assert!( abort_guard.armed, - "the caller owns cleanup until the JVM does" + "the caller owns cleanup until the JVM acknowledges the handoff" ); let written: Vec = data_files .iter() @@ -1990,6 +2011,66 @@ mod tests { assert!(!abort_guard.armed, "aborting also gives up ownership"); } + #[tokio::test] + async fn output_stream_waits_for_jvm_eof_poll_before_releasing_cleanup() { + for acknowledge in [false, true] { + let temp_dir = TempDir::new().unwrap(); + let data_location = format!("file://{}", temp_dir.path().display()); + let schema = iceberg_user_schema(); + let spec = PartitionSpec::builder(Arc::new(schema.clone())) + .build() + .unwrap(); + let common = common( + data_location, + serde_json::to_string(&spec).unwrap(), + serde_json::to_string(&schema).unwrap(), + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, + ); + let (data_files, guard) = run_write_task( + input_stream(vec![batch(&[1], &["us"])]), + common, + Arc::new(schema), + Arc::new(spec), + ProtoIcebergWriterMode::IcebergWriterUnpartitioned, + WriterProperties::builder().build(), + Some(0), + Some(0), + Time::default(), + ) + .await + .unwrap(); + + let written: Vec = data_files + .iter() + .map(|file| PathBuf::from(file.file_path().trim_start_matches("file:"))) + .collect(); + assert!(!written.is_empty()); + assert!(written.iter().all(|path| path.exists())); + + let output = + build_output_batch(vec![], &guard.locations(), &build_output_schema()).unwrap(); + let mut stream = Box::pin(output_with_cleanup_ack(output, guard)); + assert!(stream.next().await.unwrap().is_ok()); + + if acknowledge { + assert!(stream.next().await.is_none()); + } + drop(stream); + + if acknowledge { + assert!(written.iter().all(|path| path.exists())); + } else { + tokio::time::timeout(Duration::from_secs(5), async { + while written.iter().any(|path| path.exists()) { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("dropping an unacknowledged output must delete its files"); + } + } + } + #[tokio::test] async fn fanout_partitioned_write_produces_one_file_per_partition() { let temp_dir = TempDir::new().unwrap(); diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala index a8a33a7a206..877eb84c1f1 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometIcebergWriteExec.scala @@ -307,6 +307,7 @@ case class CometIcebergWriteExec( require( batch.numCols() == 2, s"iceberg_write expected 2 output columns per task, got ${batch.numCols()}") + CometIcebergWriteExec.beforeNativeHandoff() val locations = CometIcebergWriteExec.decodeLocations(batch.column(1).getBinary(0)) cleanup.own(locations) CometIcebergWriteExec.afterNativeHandoff(locations) @@ -321,6 +322,21 @@ case class CometIcebergWriteExec( object CometIcebergWriteExec { + // Local-executor test hook after the native output batch arrives but before the JVM decodes and + // takes cleanup ownership of its locations. The callback is absent outside a scoped test. + private val preHandoffFailpoint = new AtomicReference[() => Unit]() + + private[apache] def withPreNativeHandoffFailpoint[T](callback: () => Unit)(body: => T): T = { + val previous = preHandoffFailpoint.getAndSet(callback) + try body + finally preHandoffFailpoint.set(previous) + } + + private[comet] def beforeNativeHandoff(): Unit = { + val callback = preHandoffFailpoint.get() + if (callback != null) callback() + } + // Local-executor test hook for the boundary between owning the native payload's paths and // decoding its manifest. The callback is absent outside a scoped test invocation. private val handoffFailpoint = new AtomicReference[Seq[String] => Unit]() diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala index e8aa36e2c52..68718a9771a 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala @@ -2053,6 +2053,79 @@ class CometIcebergWriteActionSuite } } + test("native acceleration: a pre-native handoff failure cleans up task files") { + assumeNativeAcceleration() + withIcebergCatalog { warehouseDir => + createTable(warehouseDir, "pre_handoff_target", partitionSpec = "") + coalesceInsert("pre_handoff_target", Seq((0, "seed", 0.0))) + val before = countSnapshots("pre_handoff_target") + val root = dataDir("pre_handoff_target").toPath.toAbsolutePath + + def relativePath(location: String): String = { + val uri = new java.net.URI(location) + val file = if (uri.getScheme == null) new File(location) else new File(uri) + root.relativize(file.toPath.toAbsolutePath).toString + } + + def metadataFiles: Set[String] = spark + .sql(s"SELECT file_path FROM $catalog.$ns.pre_handoff_target.files") + .collect() + .map(row => relativePath(row.getString(0))) + .toSet + + val committed = metadataFiles + assert(committed.nonEmpty, "seed write did not create a data file") + assert(parquetFiles(root.toFile) == committed) + + val session = spark + import session.implicits._ + (1 to 1000) + .map(i => (i, s"r$i", i.toDouble)) + .toDF("id", "region", "amount") + .coalesce(1) + .createOrReplaceTempView("pre_handoff_src") + + val attempts = new AtomicReference[Vector[(Int, Int)]](Vector.empty) + val (failedPlans, error) = withNativeEnabled { + CometIcebergWriteExec.withPreNativeHandoffFailpoint { () => + val tc = TaskContext.get() + attempts.getAndUpdate(_ :+ ((tc.partitionId(), tc.attemptNumber()))) + throw new RuntimeException("pre-native handoff injected failure") + } { + captureFailedPlans(spark) { + spark.sql(s"INSERT INTO $catalog.$ns.pre_handoff_target " + + "SELECT id, region, amount FROM pre_handoff_src") + } + } + } + + assert( + error.toSeq + .flatMap(exceptionChain) + .exists(t => + Option(t.getMessage).exists(_.contains("pre-native handoff injected failure"))), + s"expected the handoff failure to reach Spark, got $error") + assert( + failedPlans.exists(p => + collectWithSubqueries(p) { case w: CometIcebergWriteExec => w }.nonEmpty), + s"failed write did not run natively:\n${failedPlans.mkString("\n--\n")}") + + val handoffs = attempts.get() + assert(handoffs.nonEmpty, "the pre-handoff failpoint was never reached") + assert(handoffs.forall(_._1 == 0), s"expected only partition 0: $handoffs") + assert( + handoffs.map(_._2) == handoffs.indices.toVector, + s"expected consecutive attempts: $handoffs") + + assert(countSnapshots("pre_handoff_target") == before, "failed write must not commit") + assertRows("pre_handoff_target", expectedIds = Seq(0)) + val physical = parquetFiles(root.toFile) + val referenced = metadataFiles + assert(physical == referenced, s"orphan files: ${physical -- referenced}") + assert(referenced == committed, s"failed write changed the table files: $referenced") + } + } + test("native acceleration: a post-native handoff failure cleans up task files") { assumeNativeAcceleration() withIcebergCatalog { warehouseDir =>