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

   ## Which issue does this PR close?
   
   No issue filed; this fixes a flaky test seen in CI on #6664.
   
   ## Rationale for this change
   
   `CometIcebergWriteActionSuite` "write report records which writer ran each 
Iceberg write" fails intermittently. The failure on #6664 (Spark 4.1, JDK 17 
[scans]) was:
   
   ```
   ArraySeq("spark", "native", "jvm", "jvm", "spark") did not equal 
List("native", "jvm", "jvm", "spark")
   ```
   
   The extra leading record is `ReportedWrite(spark,AppendData,List(),false)`. 
#6693 added a setup write right before the listener window on Spark 3.5+, with 
the split operator disabled, so it plans Spark's `AppendData`. `reportedWrites` 
registers its `QueryExecutionListener` without draining the listener bus first. 
These listeners receive events asynchronously from the bus, and the listener 
manager hands each event to whichever listeners are registered when it is 
dispatched. So when the bus is behind, the setup write's event can arrive after 
registration and gets recorded. The helper already drains after the action, but 
not before registering.
   
   ## What changes are included in this PR?
   
   Call `CometListenerBusUtils.waitUntilEmpty` before registering the listener 
in every Iceberg test helper that captured events without draining first:
   
   - `reportedWrites` in `CometIcebergWriteActionSuite` (the failing test and 
the other write report tests).
   - `capturePlans` and `captureFailedPlans` in `CometIcebergTestBase`. 
`captureWrite` and many other tests go through these, and the setup writes 
before them could leak plans into the capture in the same way.
   - The task-end listener in "a failed write job deletes the data files of 
tasks that completed". A leftover task-end event from the seed insert could 
count down `JobAbortGate` before the write's own tasks finish.
   - `captureSqlPlans` in `CometIcebergRewriteActionSuite`, and `capturePlans` 
in `CometIcebergWriteBenchmark`.
   
   `captureWritePlan` in `CometIcebergWriteDetectionSuite` and 
`taskInputMetrics` in `CometIcebergNativeSuite` already drain before 
registering.
   
   ## How are these changes tested?
   
   Test-only change. On Spark 4.1, with the target test temporarily repeated 30 
times, all iterations pass. The race did not reproduce locally on its own, so 
to force it I temporarily added a `SparkListener` that sleeps 300 ms on each 
`SparkListenerSQLExecutionEnd`, which makes the shared queue lag behind:
   
   - Without the new drain in `reportedWrites`: fails with the same 
`ArraySeq("spark", "native", "jvm", "jvm", "spark")` mismatch seen in CI.
   - With it: passes.
   
   `CometIcebergWriteActionSuite` and `CometIcebergRewriteActionSuite` pass in 
full on Spark 4.1.
   
   This pull request and its description were written by Isaac.
   


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