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]
