andygrove commented on code in PR #4746:
URL: https://github.com/apache/datafusion-comet/pull/4746#discussion_r4124374901
##########
spark/src/main/scala/org/apache/comet/serde/operator/NativeWriteUtils.scala:
##########
@@ -141,12 +141,12 @@ object NativeWriteUtils {
* byte for byte (see [[hdfsPathDivergence]] for why it may not):
*
* - the destination directory, and
- * - `fileNamePrefix`, the basename every file name is built from. On
Spark 4.0+ that is
+ * - `fileNamePrefix`, the basename every file name is built from. That is
* `mapreduce.output.basename`, which
`HadoopMapReduceCommitProtocol.getFilename`
- * interpolates into `<basename>-<split>-<jobId>`; on 3.x Comet names
the files itself and
- * the basename is always the literal `part`. A basename holding `?` or
`#` is the dangerous
- * one: the native URL parser truncates there, so *every* task writes a
file with the same
- * truncated name and they overwrite each other during commit.
+ * interpolates into `<basename>-<split>-<jobId>`; both native writers
take their file names
Review Comment:
I asked for this comment to change, but the new wording overshoots on 3.x.
Spark 3.4 and 3.5 hardcode the prefix in
`HadoopMapReduceCommitProtocol.getFilename` as `part-$split%05d-$jobId`, and
reading `mapreduce.output.basename` there only arrived in 4.0. So with
`option("mapreduce.output.basename", "out")` on 3.5 this writer still produces
`part-...` files, which is also what Spark does. That means the new
`hadoopConf.get(BASE_OUTPUT_NAME, "part")` in `CometDataWritingCommand` can
decline an HDFS write over a basename that never reaches a file name. Could the
3.x serde go back to `DEFAULT_BASE_OUTPUT_NAME` with a comment pointing at
3.x's `getFilename`, and could this doc say the basename only matters on 4.0+?
`checkNativeWriteDestination` still catches a custom committer that does use it.
##########
spark/src/test/scala/org/apache/comet/parquet/CometParquetWriterSuite.scala:
##########
@@ -1080,16 +1080,13 @@ class CometParquetWriterSuite extends
CometParquetWriterTestBase {
}
//
---------------------------------------------------------------------------------------------
- // 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
Review 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.
##########
spark/src/main/scala/org/apache/spark/sql/comet/CometNativeWriteExec.scala:
##########
@@ -104,219 +78,124 @@ case class CometNativeWriteExec(
"rows_written" -> SQLMetrics.createMetric(sparkContext, "number of written
rows"))
override def doExecute(): RDD[InternalRow] = {
- // Setup job if committer is present
- committer.foreach { c =>
- val jobContext = createJobContext()
- c.setupJob(jobContext)
- }
-
- // Execute the native write with commit protocol
- val resultRDD = doExecuteColumnar()
-
- // Force execution by consuming all batches
- resultRDD
- .mapPartitions { iter =>
- iter.foreach(_.close())
- Iterator.empty
- }
- .count()
-
- // Extract write statistics from metrics
- val filesWritten = metrics("files_written").value
- val bytesWritten = metrics("bytes_written").value
- val rowsWritten = metrics("rows_written").value
-
- // Collect TaskCommitMessages from accumulator
- val commitMessages = taskCommitMessagesAccum.value.asScala.toSeq
-
- // Commit job with collected TaskCommitMessages
- committer.foreach { c =>
- val jobContext = createJobContext()
- try {
- c.commitJob(jobContext, commitMessages)
- logInfo(
- s"Successfully committed write job to $outputPath: " +
- s"$filesWritten files, $bytesWritten bytes, $rowsWritten rows")
- } catch {
- case e: Exception =>
- logError("Failed to commit job, aborting", e)
- c.abortJob(jobContext)
- throw e
- }
- }
-
- // Return empty RDD as write operations don't return data
+ executeWriteAndCommit()
sparkContext.emptyRDD[InternalRow]
}
override def doExecuteColumnar(): RDD[ColumnarBatch] = {
- // Comet replaces DataWritingCommandExec entirely, so Spark's
- // InsertIntoHadoopFsRelationCommand.run() never runs. That method is
where Spark handles
- // SaveMode semantics (path-exists check, delete-before-Overwrite, Ignore
short-circuit) -
- // port the non-partitioned, non-catalog branch of that logic here. See
Spark 3.5's
- // InsertIntoHadoopFsRelationCommand.run doInsertion match. This runs on
the driver before
- // any executor tasks fire, mirroring where Spark does the delete.
+ executeWriteAndCommit()
+ sparkContext.emptyRDD[ColumnarBatch]
+ }
+
+ private def executeWriteAndCommit(): Unit = {
if (!prepareOutputPathForMode()) {
logInfo(s"Skipping insertion into $outputPath - already exists
(SaveMode.$mode)")
- return sparkContext.emptyRDD[ColumnarBatch]
+ return
}
- // Get the input data from the child operator
+ val job = Job.getInstance(new Configuration(serializableHadoopConf.value))
+ // Like FileFormatWriter, only abort after setupJob has succeeded.
+ committer.setupJob(job)
+ Utils.tryWithSafeFinallyAndFailureCallbacks(block = {
+ // Include configuration changes made by setupJob in the task contexts.
+ val commitMessages = runNativeWriteJob(new
SerializableConfiguration(job.getConfiguration))
+ committer.commitJob(job, commitMessages.toSeq)
+ logInfo(
+ s"Successfully committed native write job to $outputPath: " +
+ s"${metrics("files_written").value} files, " +
+ s"${metrics("bytes_written").value} bytes,
${metrics("rows_written").value} rows")
+ })(catchBlock = committer.abortJob(job))
+ }
+
+ private def runNativeWriteJob(
+ hadoopConf: SerializableConfiguration): Array[TaskCommitMessage] = {
val childRDD = if (child.supportsColumnar) {
child.executeColumnar()
} else {
- // If child doesn't support columnar, convert to columnar
child.execute().mapPartitionsInternal { _ =>
- // TODO this could delegate to CometRowToColumnar, but maybe Comet
- // does not need to support this case?
throw new UnsupportedOperationException(
"Row-based child operators not yet supported for native write")
}
}
- // Capture metadata before the transformation
- val numPartitions = childRDD.getNumPartitions
- val numOutputCols = child.output.length
+ // SPARK-23271: a zero-partition input would spawn no task and therefore
write no file at all,
Review Comment:
Could this swap get a 3.x test? `CometEmptyRelationParquetWriterSuite` is
under `spark-4.x`, and the SPARK-23271 test covers a single partition with no
batches, so nothing on 3.x reaches this branch. A write from an empty directory
does, on every version: `spark.read.schema("id INT, name
STRING").parquet(emptyDir)` plans a zero-partition native scan. On `main` that
write leaves no output directory at all and the read-back fails with
`PATH_NOT_FOUND`, while this branch writes `_SUCCESS` and a schema-only file.
That's #5303, which this PR fixes on 3.x.
--
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]