andygrove opened a new pull request, #6153:
URL: https://github.com/apache/datafusion-comet/pull/6153
## Which issue does this PR close?
Closes #6143.
## Rationale for this change
#6143 reported that `IcebergCommitExec` rethrows the raw cause of a failed
write where Spark 3.x wraps it in a "Writing job aborted" `SparkException`.
Checking Spark's source shows that premise is only true for Spark 3.3 and
earlier, which Comet does not support.
At every Spark version Comet builds against (3.4, 3.5, 4.0, 4.1, 4.2),
`V2TableWriteExec.writeWithV2` catches any `Throwable` from the write job or
from `batchWrite.commit` and aborts. Then:
- if `batchWrite.abort` succeeds, it rethrows the original cause unchanged
- if `batchWrite.abort` itself throws, it attaches the abort failure to the
cause as suppressed and throws
`QueryExecutionErrors.writingJobFailedError(cause)`, a `SparkException`
("Writing job failed.") with the original failure as its cause
| Spark tag | `WriteToDataSourceV2Exec.scala`, wrap / rethrow |
`QueryExecutionErrors.writingJobFailedError` |
|---|---|---|
| v3.4.3 | 434 / 437 | 912 |
| v3.5.8 | 416 / 419 | 899 |
| v4.0.1 | 452 / 455 | 889 |
| v4.1.1 | 483 / 486 | 944 |
| v4.2.0 | 664 / 667 | 973 |
`IcebergCommitExec` already matched Spark when the abort succeeds. It
diverged only when the abort failed: it attached the abort failure as
suppressed but rethrew the raw cause instead of wrapping it. So the fix is
narrower than the issue suggests, and needs no version shim.
## What changes are included in this PR?
- `IcebergCommitExec`: on both the job-failure and the commit-failure path,
when `batchWrite.abort` throws, throw `writingJobFailedError(cause)` instead of
`cause`. A small helper runs the abort and reports whether it failed.
Everything else is unchanged: the abort itself, the deletion of completed
tasks' data files on a job failure, and attaching abort and cleanup failures to
the cause as suppressed exceptions. Cleanup failures alone do not cause
wrapping, because that cleanup is Comet's own addition and has no counterpart
in Spark's path.
- `CometIcebergWriteActionSuite`, four tests:
- a commit-time validation failure (a serializable overwrite validated
against an older snapshot, which fails deterministically without concurrency),
run with the split operator off and on, asserting the thrown exception and its
cause have the same classes, and that the split run went through
`IcebergCommitExec`
- a failed write job (a UDF that fails a task over a Parquet source), with
the same off/on comparison
- an abort that fails after a commit failure, and after a job failure,
driving `IcebergCommitExec` with a stub `BatchWrite`, asserting the
`writingJobFailedError` wrapping, the cause, and the suppressed abort failure
One thing noticed while writing the job-failure test: the existing "failed
write job aborts and leaves the table unchanged" test uses a local relation as
its source, so on Spark 4.1 the optimizer evaluates its UDF during planning and
the query fails before a write job starts. The new job-failure test reads from
Parquet so that a task actually fails. The existing test is unchanged here.
## How are these changes tested?
- Default profile (Spark 4.1): `./mvnw test -Dtest=none
-Dsuites="org.apache.comet.CometIcebergWriteActionSuite"`, 73 succeeded, 0
failed.
- Spark 3.5: the four new tests with `-Pspark-3.5`, 4 succeeded, 0 failed.
- The two abort-failure tests fail against `main`'s `IcebergCommitExec`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]