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]