andygrove opened a new pull request, #6699:
URL: https://github.com/apache/datafusion-comet/pull/6699

   ## Which issue does this PR close?
   
   Closes #6643.
   
   ## Rationale for this change
   
   The `boom_after_others` UDF waited up to two minutes for the write's other 
two tasks to finish. Comet evaluates that UDF inside the native plan. The 
plan's only input was a native Parquet scan, so it ran on a worker of Comet's 
process-wide Tokio runtime, and the waiting task held that worker. The other 
two tasks need workers for their plans as well, so on a one-worker runtime they 
never run and the gate times out.
   
   Setting `COMET_WORKER_THREADS=1` on main reproduces the trace in the issue. 
One of the three tasks finished, and `JobAbortGate.awaitOthers` threw after 120 
s. The trace has no executor frames below `CometUdfBridge.evaluate`, so the 
wait ran on a Tokio worker rather than a Spark task thread. The suite has run 
on `local[5,2]` since #6111, and in that run the retry hid the stall as a 
two-minute pass.
   
   The runtime is process-wide and takes its size from whichever session 
creates it, so its worker count depends on which suites ran earlier in the JVM. 
I haven't confirmed how many workers it had in the failing run.
   
   ## What changes are included in this PR?
   
   The test now reads a three-slice RDD instead of three Parquet files. The 
task whose slice holds id 25 waits in the `mapPartitions` function that builds 
its slice. That function runs in JVM code on the task's own thread, before any 
of the task's native plans run. The UDF, renamed `boom_at_25`, now only throws. 
`CometTestBase` turns on `spark.comet.convert.rdd.enabled`, so the project with 
the UDF and the native Iceberg writer still run natively.
   
   ## How are these changes tested?
   
   This PR changes only the test.
   
   - With `COMET_WORKER_THREADS=1`, both variants pass in seconds on Spark 3.4 
and 4.1, with no gate timeouts. A temporary print showed the wait running on an 
`Executor task launch worker` thread.
   - With the `deleteCompletedTaskFiles` call in `IcebergCommitExec` removed, 
both variants fail on Spark 4.1 and leave the two completed tasks' files 
behind. On Spark 3.4 they still pass, because Iceberg 1.5.2's own 
`SparkWrite.abort` deletes those files.
   - The full `CometIcebergWriteActionSuite` passes on Spark 3.4 and 4.1.
   


-- 
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