Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions .ai/skills/review-comet-iceberg-write-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
32 changes: 19 additions & 13 deletions docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
21 changes: 11 additions & 10 deletions docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
103 changes: 92 additions & 11 deletions native/core/src/execution/operators/iceberg_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<Item = DFResult<RecordBatch>> + 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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -1770,13 +1789,15 @@ mod tests {
CompressionCodec as ProtoCodec, IcebergParquetWriteSettings, IcebergWriteCommon,
IcebergWriterMode as ProtoIcebergWriterMode,
};
use futures::StreamExt;
use iceberg::spec::{
Manifest, NestedField, PartitionSpec, PrimitiveType, Schema, Transform, Type,
};
use parquet::file::properties::WriterProperties;
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use tempfile::TempDir;

fn user_schema() -> SchemaRef {
Expand Down Expand Up @@ -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<PathBuf> = data_files
.iter()
Expand All @@ -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<PathBuf> = 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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]()
Expand Down
Loading