Repository navigation
test: drain the listener bus before registering Iceberg test listeners - #6809
Merged
Merged
Conversation
Co-authored-by: Isaac <no-reply@databricks.com>
manuzhang
approved these changes
Oct 9, 2026
Member
There was a problem hiding this comment.
LGTM. Follow-ups that need not block this PR:
- Apply the same pre-registration drain to the three helpers in
CometMergeRowsNativeSuiteBase, which carry the identical race. - Decide on a timeout policy for the drain, either catching
TimeoutExceptionascaptureWritePlandoes or passing a longer deadline, and apply it consistently. - Consolidate the idiom into
withQueryExecutionListenerandwithSparkListenerhelpers onCometListenerBusUtilsso new call sites cannot forget the drain.
Member
Author
|
Thanks @manuzhang |
This was referenced Oct 9, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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: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'sAppendData.reportedWritesregisters itsQueryExecutionListenerwithout 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.waitUntilEmptybefore registering the listener in every Iceberg test helper that captured events without draining first:reportedWritesinCometIcebergWriteActionSuite(the failing test and the other write report tests).capturePlansandcaptureFailedPlansinCometIcebergTestBase.captureWriteand many other tests go through these, and the setup writes before them could leak plans into the capture in the same way.JobAbortGatebefore the write's own tasks finish.captureSqlPlansinCometIcebergRewriteActionSuite, andcapturePlansinCometIcebergWriteBenchmark.captureWritePlaninCometIcebergWriteDetectionSuiteandtaskInputMetricsinCometIcebergNativeSuitealready 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
SparkListenerthat sleeps 300 ms on eachSparkListenerSQLExecutionEnd, which makes the shared queue lag behind:reportedWrites: fails with the sameArraySeq("spark", "native", "jvm", "jvm", "spark")mismatch seen in CI.CometIcebergWriteActionSuiteandCometIcebergRewriteActionSuitepass in full on Spark 4.1.