kazuyukitanimura commented on code in PR #6699:
URL: https://github.com/apache/datafusion-comet/pull/6699#discussion_r4190395977


##########
spark/src/test/scala/org/apache/comet/CometIcebergWriteActionSuite.scala:
##########
@@ -2835,65 +2835,61 @@ class CometIcebergWriteActionSuite
         coalesceInsert(table, Seq((0, "seed", 0.0)))
         val committed = parquetFiles(dataDir(table))
         val before = countSnapshots(table)
-        val session = spark
-        import session.implicits._
-        withTempPath { dir =>
-          (1 to 30)
-            .map(i => (i, s"r$i", i.toDouble))
-            .toDF("id", "region", "amount")
-            .repartition(3)
-            .write
-            .parquet(dir.getAbsolutePath)
-          
spark.read.parquet(dir.getAbsolutePath).createOrReplaceTempView("job_abort_src")
-          JobAbortGate.reset(othersToFinish = 2)
-          spark.udf.register(
-            "boom_after_others",
-            (id: Int) => {
-              if (id == 25) {
-                JobAbortGate.awaitOthers()
-                throw new RuntimeException("boom")
-              }
-              id
-            })
-          val listener = new SparkListener {
-            override def onTaskEnd(taskEnd: SparkListenerTaskEnd): Unit =
-              if (taskEnd.reason == Success) JobAbortGate.taskFinished()
+        JobAbortGate.reset(othersToFinish = 2)

Review Comment:
   Thanks for running this one down. The COMET_WORKER_THREADS=1 reproduction 
and the check that the test still fails without deleteCompletedTaskFiles make 
this easy to trust. Moving the wait into the mapPartitions closure keeps it on 
the task thread, because CometExecRDD.compute resolves its input iterators 
eagerly. If it ever moved back, the gate
     would fail loudly. The test-order dependence you describe looks real. 
init_runtime does nothing once a runtime exists, so the worker count comes from 
whichever session builds it first in the JVM, and an earlier local[1] suite can 
leave it with one worker. Could you file an issue for that and link it from the 
comment here? Then the next test that
     blocks inside a native plan has somewhere to point.



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