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]

Reply via email to