Repository navigation
fix: honor Spark’s commit protocol in Spark 3.x native writes - #4746
Conversation
|
Thanks @peterxcli please check on your fork https://github.com/apache/datafusion-comet/actions/workflows/spark_sql_writer_tests.yml |
|
@comphead would you like to take a look? Thanks! |
|
Thanks @peterxcli there is another angle #5293 |
915b174 to
278f28e
Compare
andygrove
left a comment
There was a problem hiding this comment.
#5763 landed on the 19th as 5442c937d, so this branch still carries its pre-merge commits and conflicts with main in six files. Replaying only 278f28e0a and e9f5e5e7f onto main conflicts just in the CometDataWritingCommand imports, where #5821 added the hasEmptyRelationInput guard, and in the two suite lists. Could you rebase down to those two commits? Two comments describe the old 3.x writer and stop being true with this change. NativeWriteUtils.escapedHdfsDestination says that on 3.x Comet names the files itself and the basename is always part, and the ParquetWriterExec field docs say work_dir is the Spark 3.x path. After this nothing sets work_dir, job_id or task_attempt_id, so could the Some(work_dir) arm in ParquetWriterExec::execute go as well, with the three proto fields reserved?
The task closure calls createTaskContext, which is an instance method because it reads jobTrackerID. That pulls the whole CometNativeWriteExec, child subtree included, into every task despite the captured* locals, which is the thing CometWriteFilesExec goes out of its way to avoid. Could createTaskContext move to the companion object and take jobTrackerID as a parameter?
runNativeWriteJob now runs the same task loop as CometWriteFilesExec, but it leaves out the two empty-input cases that one copies from FileFormatWriter. A zero-partition child runs no task at all, so the output gets a _SUCCESS and no schema-bearing file, and spark.read.parquet on it fails to infer a schema (SPARK-23271, #5303). And every empty partition still writes its own file, because create_arrow_writer runs before the first batch. Could this take the same parallelize(Seq.empty[ColumnarBatch], 1) swap and the sparkPartitionId == 0 || batches.hasNext check, and drop the assume(isSpark40Plus) from the SPARK-23271 and empty-partition tests in CometParquetWriterSuite?
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The legacy writer bypassed Spark’s configured commit protocol, discarded its chosen filename, and handled job completion inconsistently between execution entry points.
- Design approach: Spark 3.x now runs the configured job/task lifecycle. The included Spark 4 prerequisite replaces
WriteFilesExec, leaving job commit, SaveMode handling and catalog refresh with Spark. - Correctness / compatibility analysis: Compared the write contracts against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Filename ownership, task identifiers, commit-message delivery and exception preservation follow Spark’s contracts. The existing P2 concerning percent-bearing HDFS basenames remains in this checkout. A local probe using the current guard methods and real Hadoop
Pathconfirmed thatpart%fooandpart%25pass planning but fail the task guard. Ordinarypartsucceeds. This duplicates an existing finding, so no new inline finding is returned. - Key design decisions: Separate version paths are justified by Spark’s concrete-class versus base-trait write discovery. Explicit task state keeps the Spark 4 serialization boundary clear. The existing Spark 3 closure-capture concern remains visible. No measured performance regression was established, and the per-row statistics callback cost remains unmeasured.
- Implementation sketch: Prepare the Hadoop job, allocate the committer’s exact filename, execute and close the native iterator, then commit the task and deliver its message through
runJob. Failures preserve the original exception while invoking abort callbacks. - Behavioral changes worth calling out: Both Spark 3 execution entry points now complete job commit, and filenames follow Spark’s convention. Spark 4 retains Spark’s surrounding write framework. Existing Spark 3 empty-input concerns are already covered by the prior review.
- Suggested improvements: Resolve the existing HDFS basename P2 by applying the Java URI-escaping check to the basename during planning, allowing these writes to fall back before execution.
Reviewed the entire 21-file diff from 481aefea9c60592612650de73bd2e1c7aa173979 to e9f5e5e7fe22baa4d5ff74499a13c7c72aa3beb6, including prerequisites. Confirmed non-draft status and read existing discussion. Routed skill: review-comet-pr. No sibling skill applies to this operator review.
Exact-head CI: 54 successful checks and 10 skipped, with no failures. The inspected CI merge e40331c50bd2f0435798837b835249ec019596ac has the same full tree as the reviewed head. Logs confirm all eight Spark 3.4 commit-lifecycle tests passed, plus 39 writer tests with nine cancellations. Spark 4.0 passed all 48 writer tests.
Local validation: the six-case path probe, suite inventory check and git diff --check passed. No full native/JVM build was run locally. System Maven rejected the repository’s maven.config. No local HDFS integration, automatic retry, speculation or throughput benchmark was run.
No additional introduced P1/P2 issues found within this review. Existing blockers remain.
|
Thanks, all addressed:
With the swap in place, the |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The Spark 3.x writer hardcoded the committer, discarded its allocated filename and handled job completion inconsistently between execution entry points.
- Design approach: Use Spark’s configured
FileCommitProtocol, prepared Hadoop configuration and exact task filename. Return task commit messages throughrunJob. - Correctness / compatibility analysis: Compared lifecycle and filename contracts against Spark 3.4.3 and 3.5.9, and checked the shared native path against supported Spark 4.x sources. Task identifiers, message delivery and cleanup ordering follow those contracts. The earlier HDFS basename, closure-capture and empty-input concerns are addressed in the current tree.
- Key design decisions: Moving task-context creation to the companion object avoids capturing the execution plan. Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. No introduced performance regression was established. Throughput was not benchmarked.
- Implementation sketch: Prepare the job, allocate the committer’s filename, execute and close the native iterator, then commit the task and job. Failure paths abort while preserving the original exception.
- Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition input gets a schema-bearing file, and empty nonzero partitions skip file creation. Removed protobuf field numbers and names are reserved.
- Suggested improvements: None at P1/P2 severity.
Reviewed the entire 10-file diff from a86c9672a63c909f0fd7b752c86b5a679df4e84b to 086bc1eefc49dfa694e210c077984533e8de8450. Confirmed non-draft status and read the existing reviews and discussion. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI: Checks attached to this SHA currently show 20 successful, 13 skipped and two running, with no failures. The running checks are Spark 4.1 execution and expression suites. Inspected logs confirm 49 Parquet writer tests, the empty-relation writer test and four native writer tests passed. All eight Spark 3.x lifecycle tests were canceled by their version gate. CI actually checked out merge 110c2dc6711a86ac089d710eddf9eb23041c19a5, which includes four additional changed files from main.
Validation: An eight-case probe compiled from the current path-guard methods passed using real Hadoop Path, including both previously failing percent-bearing basenames. Suite registration and git diff --check passed. No full local native/JVM build was run. Earlier Spark 3.x runtime evidence predates the latest edits. Spark SQL CI was skipped, and real HDFS, retries, speculation and performance were not validated locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
andygrove
left a comment
There was a problem hiding this comment.
I ran this on Spark 3.5 merged onto current main, and CometNativeWriteSuite and the un-gated writer tests pass. I also added run-all-spark-profiles, since the PR build only runs Spark 4.1 where every new 3.x test cancels itself. The [scans] bucket with both writer suites is green on 3.4, 3.5, 4.0 and 4.2.
This also fixes silent data loss on main. The old 3.x writer named files part-<partition>-<taskAttemptId>.parquet and wrote them straight into the destination, and task attempt ids restart at zero in every application. So two applications appending to the same path wrote the same names and the second overwrote the first. On main, two runs that each append 10 rows read back 10. With this PR they read back 20.
On your question about hasEmptyRelationInput in CometDataWritingCommand, I'd drop it in this PR. CometEmptyRelationExec only exists on Spark 4.0+ (the 3.x ShimCometEmptyRelation.emptyRelationClass is None), and since #5763 that serde only runs on 3.x, so the guard can't return true anymore. Its comment about the writer only mapping existing partitions also stops being true with your swap. The same sentence in CometEmptyRelationParquetWriterSuite at line 67 needs updating too.
The commit list still has the seven #5763 commits in it, and the squash message is built from the commit list, so it would repeat all of #5763's messages. Could you rebase onto main so only your own commits remain? git push --force-with-lease is fine here.
| // Spark 4.0+ only. These cover behavior that comes from leaving Spark's write framework in | ||
| // place, which is only possible where `V1WritesUtils.getWriteFilesOpt` matches the | ||
| // `WriteFilesExecBase` trait. See CometWriteFilesExec. | ||
| // Commit-protocol checks run on both writers. Tests requiring the surrounding Spark write |
There was a problem hiding this comment.
dynamic partition overwrite falls back to Spark passes on 3.5 with its assume removed. The 3.x path now hardcodes dynamicPartitionOverwrite = false and relies on the partitioned-write decline, so the reasoning in that test's comment applies on 3.x as well. Could it lose the gate? The maxRecordsPerFile test next to it can't yet, because the 3.x serde never declines maxRecordsPerFile: on 3.5 the write stays native and produces one file instead of ten. That predates this PR, so it doesn't need to hold this one up.
086bc1e to
f8bddaa
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Checked lifecycle, filenames, task identifiers and empty-input handling against Spark 3.4.3 and 3.5.9, plus relevant shared contracts in 4.0.4, 4.1.3 and 4.2.0. Spark 3.x correctly retains its fixed
partprefix. Earlier review concerns are addressed in the current code and tests. - Key design decisions: Removing the accumulator, iterator wrapper and native filename generation simplifies ownership. Capturing task inputs avoids serializing the execution-plan subtree. No introduced performance regression was established. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit each task, deliver its message through
runJob, then commit the job. Failure callbacks preserve the original exception while attempting cleanup and abort. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, while empty nonzero partitions skip file creation. Removed protobuf identifiers are reserved.
- Suggested improvements: None supported at P1/P2 severity.
Reviewed all five commits and all 11 changed files from base 65a0cda1cf62877a37ec1a1f5f9ebc9dd1405ebd to head f8bddaac1aecbee9ab4f1568a5c39455509569c2. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI at review time: 15 successful checks, 13 skipped, two running and 22 queued, with no failures. Spark SQL suites were skipped. CI checked out merge 4b3743e4968f30a7eaaac198d1219b2dd15976fb, whose full tree matches the reviewed head. Runtime validation remains pending.
Local validation: A 14-case Scala probe using current path guards, verified Spark filename methods and real Hadoop path/context classes passed. Suite registration, protobuf reservations and git diff --check passed. The default Java installation lacked JNI headers. Retrying with a cached JDK passed that obstacle, but the 150-second build and 240-second focused native-test attempts expired during compilation before tests ran. No full JVM/Spark SQL suites, real HDFS integration, automatic retries, speculation or benchmarks were run locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
|
This is a light fully automated review since there are so many PRs open. With |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Compared filenames, task identifiers, commit callbacks and empty-input handling against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Spark 3.x correctly retains its fixed
partprefix. Earlier filename, closure-capture and empty-input concerns are addressed. - Key design decisions: Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. Explicit task inputs avoid capturing the execution-plan subtree. No concrete performance regression was identified. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit the task, deliver its message through
runJob, then commit the job. Failure callbacks preserve the original exception while attempting cleanup and abort. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, while empty nonzero partitions skip file creation. Removed protobuf identifiers are reserved.
- Suggested improvements: None supported at P1/P2 severity.
Reviewed all five commits and all 11 changed files from base 65a0cda1cf62877a37ec1a1f5f9ebc9dd1405ebd to head f8bddaac1aecbee9ab4f1568a5c39455509569c2. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI: 40 successful checks and 14 skipped, with no failures or pending checks. CI tested merge 4b3743e4968f30a7eaaac198d1219b2dd15976fb, whose full tree matches the reviewed head. Inspected logs confirm all nine legacy-writer tests passed on both Spark 3.4 and 3.5, writer and empty-input coverage passed across all five profiles, and the native test run passed 1,776 tests with five skipped.
Local validation: A recompiled 14-case Scala probe passed using current path guards, verified Spark filename methods and real Hadoop classes. Suite registration, protobuf reservations and git diff --check passed. No full local native/JVM build was run. Spark SQL and macOS CI were skipped. Real HDFS integration, automatic retries, speculation and performance were not validated locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
Mirror CometWriteFilesExec and FileFormatWriter: swap a zero-partition child for a single empty partition so the output still carries a schema (SPARK-23271), and only write a file from partition 0 or a partition that has rows. Move createTaskContext to the companion object so the task closure no longer captures the exec node and its child subtree. The SPARK-23271 and empty-partition tests now run on Spark 3.x too.
f8bddaa to
c5868dc
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing filename collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Compared job preparation, filenames, task identifiers, commit callbacks and empty-input handling against Spark 3.4.3 and 3.5.9, plus relevant shared contracts in 4.0.4, 4.1.3 and 4.2.0. Spark 3.x correctly retains its fixed
partprefix. Earlier functional concerns are addressed. - Key design decisions: Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. Captured task inputs avoid serializing the execution-plan subtree. No concrete performance regression was identified. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit each task, deliver its message through
runJob, then commit the job. Failure callbacks preserve the original exception while attempting cleanup and abort. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, while empty nonzero partitions skip file creation. Removed protobuf identifiers are reserved.
- Suggested improvements: None supported at P1/P2 severity.
Reviewed all five commits and all 11 changed files from 098f58db88a250ab79c48575983428ce0c704c9f to c5868dcdef57ccea4511738019bbf6874155bb6a. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI: 40 successful checks and 15 skipped, with no failures or running checks. CI tested merge 687c9c1ff50bafd00af28dc9b5d8b8da5dcb8256, whose full tree matches the reviewed head. Inspected logs confirm all nine legacy-writer tests passed on Spark 3.4 and 3.5, writer and empty-input coverage passed across all five profiles, and the main native test run passed 1,855 tests with five skipped.
Local validation: A freshly compiled 14-case differential path/filename probe passed using current guards, verified Spark filename methods and real Hadoop classes. Suite registration, protobuf reservations and git diff --check passed. No full local native/JVM build was run. Spark SQL and macOS CI were skipped. Real HDFS integration, automatic retries, speculation and performance were not validated locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
# Conflicts: # native/core/src/execution/planner.rs
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing filename collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Compared job preparation, filenames, task identifiers, commit callbacks and empty-input handling against Spark 3.4.3 and 3.5.9, plus relevant shared contracts in 4.0.4, 4.1.3 and 4.2.0. Spark 3.x correctly retains its fixed
partprefix. Earlier functional concerns are addressed. - Key design decisions: Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. Captured task inputs avoid serializing the execution-plan subtree. No concrete performance regression was identified. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit each task, deliver its message through
runJob, then commit the job. Failure callbacks preserve the original exception while attempting cleanup and abort. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, while empty nonzero partitions skip file creation. Removed protobuf identifiers are reserved.
- Suggested improvements: None supported at P1/P2 severity.
Reviewed the entire 11-file diff from 63fd1c9ebd66a4159feb1079ede794ed95f0e53a to a90b9d038fac6e591ec91a17b1f4249bdaac526b. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI: Checks attached to this SHA show 31 successful, 14 skipped and nine running, with no failures. CI tested merge 37b4600ce2feec0819418e7480ea86ac10e4c346, which differs from the reviewed head in 21 files. The writer implementations and writer suites are unchanged in that delta. Inspected logs confirm all nine legacy-writer tests passed on both Spark 3.4 and 3.5, writer and empty-input coverage passed across all five profiles, and the main native run passed 1,936 tests with five skipped.
Local validation: A freshly compiled 14-case differential path/filename probe passed using current guards, verified Spark filename methods and real Hadoop classes. Suite registration, protobuf reservations and git diff --check passed. No full local native/JVM build was run. Spark SQL and macOS CI were skipped. Real HDFS integration, automatic retries, speculation and performance were not validated locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
With the empty-relation guard gone from CometDataWritingCommand, the user guide and the CometWriteFilesExec comment still described Parquet writes over an empty relation falling back to Spark's writer. Both native write paths now swap a zero-partition input for one empty partition, so such a write stays native and still leaves a readable schema-only file.
|
I would like to do a deeper review of this PR before it is merged. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Compared job preparation, filenames, task identifiers, commit callbacks and empty-input handling against Spark 3.4.3 and 3.5.9, plus relevant shared contracts in 4.0.4, 4.1.3 and 4.2.0. Spark 3.x correctly retains its fixed
partprefix. Earlier functional review concerns are addressed. - Key design decisions: Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. Captured task inputs avoid serializing the execution-plan subtree. No concrete performance regression was identified. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit each task, deliver its message through
runJob, then commit the job. Failure callbacks preserve the original exception while attempting cleanup and abort. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, while empty nonzero partitions skip file creation. Removed protobuf identifiers are reserved.
- Suggested improvements: None supported at P1/P2 severity.
Reviewed the entire 11-file diff from 63fd1c9ebd66a4159feb1079ede794ed95f0e53a to a90b9d038fac6e591ec91a17b1f4249bdaac526b. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI: 41 successful checks and 15 skipped, with no failures or pending checks. CI tested merge 37b4600ce2feec0819418e7480ea86ac10e4c346, which differs from the reviewed head in 21 files. Writer implementations and writer suites are unchanged by that delta. Inspected logs confirm all nine legacy-writer tests passed on both Spark 3.4 and 3.5, writer and empty-input coverage passed across all five profiles, and the main native run passed 1,936 tests with five skipped.
Local validation: A freshly compiled 14-case path/filename probe passed using current guards, verified Spark filename methods and real Hadoop classes. Suite registration, CI configuration and git diff --check passed. No full local native/JVM build or Spark SQL run was performed. Spark SQL and macOS CI were skipped. Real HDFS integration, automatic retries, speculation and performance were not validated locally.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
andygrove
left a comment
There was a problem hiding this comment.
Thanks for working through this. I compared it with branch-1.1 and with Spark 3.4.3 and 3.5.9. For Spark 3.x users who enable native writes, this moves the writer onto Spark's commit protocol and fixes real problems in 1.1: a second application's append could overwrite the first one's files, and failed or speculative attempts could leave partial or duplicate files. Nothing changes on Spark 4.x. A few things I'd like to see before it merges:
1. docs/source/user-guide/latest/compatibility/operators.md:32
Following up on the 28 September comment, a few places still describe the guard this PR removes. docs/source/user-guide/latest/compatibility/operators.md:32-33 tells users that a Parquet write over an empty relation uses Spark's writer, and the EmptyRelationExec row at docs/source/user-guide/latest/operators.md:57 links there for "writer fallback". On Spark 4.x CometEmptyRelationParquetWriterSuite asserts the opposite, and on 3.x the guard could never fire. The comment at spark/src/main/scala/org/apache/spark/sql/comet/CometWriteFilesExec.scala:131-135 also still says CometDataWritingCommand declines empty-relation inputs because its writer "has nowhere to put this swap", and runNativeWriteJob now does exactly that swap. Could these be changed to say that a write over an empty input stays native and leaves a readable schema-only file, on both writers?
2. spark/src/main/scala/org/apache/spark/sql/comet/CometNativeWriteExec.scala:155-164
The file name's extension comes from writerFactory.getFileExtension(taskContext), which reads the codec that prepareWrite put in the job configuration. But the codec the native writer actually uses is still the planning-time value from parseCompressionCodec in convert. The two agree today. I checked both against ParquetOptions on 3.4.3 and 3.5.9. CometWriteFilesExec.executeTask reads the codec back from CodecConfig.from(taskAttemptContext) and sets it on the per-task ParquetWriter, so that the name and the footer agree by construction, and NativeWriteUtils.protoCompressionCodec documents that as the way to do it. This PR is what puts the codec into 3.x file names. Could the 3.x task do the same thing, next to setOutputPath?
3. Spark SQL writer tests
The last run of the manual Spark SQL Writer Tests workflow on your fork was on 13 July at 468cc03. That was before the zero-partition swap, the empty-partition rule, the abort changes and the admission change went in. The Spark SQL lanes are skipped in this PR's CI. Could you run it again for Spark 3.5 on the current head and link the run here? Spark's own writer suites are the broadest check we have that the 3.x lifecycle now matches FileFormatWriter.
4. PR description
For anyone who turned the 3.x writer on in 1.1, this fixes more than file naming. In 1.1, tasks write straight into the destination with names that repeat across applications, so a second application's Append overwrites the first one's files (the 10 rows instead of 20 from the 28 September comment). A failed attempt also leaves its partial file behind, and a retry or speculative attempt adds a second copy. Could the description say that plainly, so the changelog entry tells 1.1 users what they were exposed to? Could it also say which of #2827, #3015 and #5303 it closes? With this PR both writers use the commit protocol, stage their output and swap in an empty partition, so #5303 at least looks done. Separately, is a branch-1.1 backport planned? The JVM helpers this needs all exist there. The native part can be left out, because the 1.1 native writer already writes to output_path as is when work_dir is unset.
5. HDFS coverage
On 3.x the native HDFS writer now opens its file inside the committer's _temporary tree and relies on HDFS renames at commit, the same as the 4.x writer already does. No test runs either writer against HDFS, and WithHdfsCluster already starts a MiniDFSCluster for the read benchmark. Could you open a tracking issue for an end-to-end HDFS write test that covers both writers, and link it here?
# Conflicts: # docs/source/user-guide/latest/operators.md
The Spark 3.x task named its file after the codec in the job configuration, because getFileExtension reads it from there. The native writer used the codec that convert resolved at planning time. The two agree today only because both copy the precedence rule of ParquetOptions. Read the codec with CodecConfig.from(taskContext) next to setOutputPath, as CometWriteFilesExec.executeTask does, so the file name and the footer agree by construction. A new CometNativeWriteSuite test sets a job codec that differs from the planned one and checks the name and the footer. The mixed-case codec test now checks the file name on Spark 3.x too.
Since apache#6416, CometMetricNode reports a sum metric with add instead of set. The test metric threw only from set, so the two cleanup-failure tests in CometNativeWriteSuite stopped injecting their failure and failed on Spark 3.4 and 3.5.
The Parquet writer test step caps the sbt heap at 1 GB for the forked test JVM, but this workflow compiles Spark first, and sbt ran out of heap while it compiled sql/core. Compile the SQL test classes in an earlier step with a 3 GB heap. This branch had the same step in July, before a history rewrite dropped it.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spark 3.x native writes discarded the committer’s filename and handled job completion inconsistently, allowing collisions across append jobs.
- Design approach: Prepare the Hadoop job, use Spark’s configured
FileCommitProtocol, and pass its exact task filename to native code. - Correctness / compatibility analysis: Compared job preparation, task identifiers, commit callbacks, filenames, codecs and empty-input handling against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The inspected paths follow those contracts. Earlier functional review concerns are addressed.
- Key design decisions: Removing the accumulator, iterator wrapper and native filename-generation branch simplifies ownership. Captured task inputs avoid serializing the execution-plan subtree. No concrete performance regression was identified. Throughput was not benchmarked.
- Implementation sketch: Allocate the task filename, drain and close the native iterator, commit each task, deliver its message through
runJob, then commit the job. Cleanup and abort failures preserve the original exception. - Behavioral changes worth calling out: Both execution entry points complete job commit. Zero-partition inputs produce a schema-bearing file, and empty nonzero partitions skip file creation. Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, these are intended fixes. Spark 4.x retains its existing write lifecycle. - Suggested improvements: None supported at P1/P2 severity.
Reviewed the full 15-file diff from 33f21da181c14df7fa3e6a3104c3d7215810147e to b09a70b9fa791e75ce04aee484508ac29039b0ab. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-memory-pr.
Exact-head CI at 12:18 UTC: 17 successful checks, 14 skipped and 23 running, with no failures. These checks test merge e8b9c62210201754c28c79821584d03d04dda701, which includes 24 additional changed files from main, including JNI changes. Standard Spark SQL lanes are skipped. The manual Spark 3.5 writer run uses the exact reviewed head and remains in progress.
Local validation: A freshly compiled 14-case filename/path probe passed. A Spark 3.5.9 probe verified all six supported codec names and task-configuration overrides against prepareWrite, CodecConfig and file extensions. Suite registration, CI configuration and git diff --check passed. No full local native/JVM build, Spark SQL suite, real HDFS integration, automatic retry/speculation test or benchmark was run.
No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
# Conflicts: # spark/src/main/scala/org/apache/comet/serde/operator/CometDataWritingCommand.scala # spark/src/main/scala/org/apache/spark/sql/comet/CometNativeWriteExec.scala
The merge with main brings in two callers written against the old constructor. apache#5957's RevertNativeForTransitionHeavyStagesSuite relied on the removed committer and jobTrackerID defaults, and apache#6247's write_reserving_from helper still passed the removed work_dir, job_id and task_attempt_id arguments to ParquetWriterExec::try_new. CleanupFailingNativeWriteExec also needs apache#5957's originalPlan argument.
andygrove
left a comment
There was a problem hiding this comment.
I merged current main into the branch and pushed it, because #5957 landed last night and conflicted with this. The merge commit keeps #5957's originalPlan parameter and its output and sparkFallback overrides on CometNativeWriteExec, with your new constructor parameters after child. createExec passes op for it. A second commit fixes three test call sites that merged cleanly but no longer compiled. cometNativeWrite in RevertNativeForTransitionHeavyStagesSuite relied on the old constructor defaults. The write_reserving_from helper that #6247 added to the parquet_writer.rs tests still passed the removed arguments. And CleanupFailingNativeWriteExec needed the new originalPlan argument. With those changes, CometNativeWriteSuite, CometParquetWriterSuite and RevertNativeForTransitionHeavyStagesSuite pass locally on 3.4 and 3.5. The writer and revert suites pass on 4.1, and so do the native writer tests.
About the Spark SQL writer run I asked for on 2 October: thanks for running it, but it turns out that workflow never runs the native writer, which I should have checked before asking. ENABLE_COMET_WRITE only turns on spark.comet.write.parquet.enabled. The 3.x write also needs spark.comet.operator.DataWritingCommandExec.allowIncompatible=true, and nothing sets it. I filed #6749 for that. Instead I ran the same suites from Spark 3.5.9's test jar with both keys set. Only 10 tests actually commit a native write, because most of these suites write from local data. The only failure the native writer causes is "Write Spark version into Parquet metadata", which is #3427. The other twelve failures happen the same way with the writer off. The tests behind #3417 and #3428 now pass natively and fail on main, so this fixes both on 3.x as well.
I also compared the writer against Spark's on 3.4 and 3.5 directly. A task that fails its first commit and gets retried leaves one file per partition and no duplicate rows. The file suffix and footer codec match Spark for every codec we admit. Filters that empty some or all partitions leave the same files Spark does. Committer algorithm v2, parquet.summary.metadata.level=ALL and a custom spark.sql.parquet.output.committer.class all work.
I filed #6750 for the end-to-end HDFS write test and #6751 for maxRecordsPerFile on 3.x, so neither needs to come from you. Thanks for sticking with this one.
Which issue does this PR close?
Addresses the Spark 3.x write issues in #2827, #3015, and #5303.
Rationale for this change
Spark 3.x native writes choose their own filenames instead of using the commit protocol's paths. This can overwrite data across append jobs and prevents configured committers from tracking the files correctly.
What changes are included in this PR?
Spark 3.4/3.5 now prepare the Hadoop write job and use Spark's configured commit protocol and exact task filename. Both execution entry points collect task commit messages and complete the job lifecycle, preserving the original error when cleanup or abort also fails. The unused native filename-generation fields are removed and their protobuf numbers reserved.
Zero-partition inputs produce a schema-only file. HDFS admission uses Spark 3.x's fixed
partprefix, with the actual filename checked before writing. The obsolete empty-relation fallback is removed. Spark 4.0+ retains the write-command lifecycle introduced by the now-merged #5763.How are these changes tested?
The writer and commit-protocol suites pass on Spark 3.4 (53 tests) and 3.5 (55 tests). Tests assert native execution, exact filenames, abort handling, schema read-back from an empty directory, and dynamic partition overwrite fallback. Spark 4.0 writer and AQE empty-relation regression checks also pass (51 tests).
Native build, formatting, Scala style, suite registration, and CI configuration checks pass. Validation uses local storage; HDFS admission is tested without a cluster. The existing Spark 3.x
maxRecordsPerFilelimitation remains outside this change.