test: [branch-1.0] wait for Parquet write plan callbacks (#6108) - #6338
Conversation
Backport of apache#6108 to branch-1.0. apache#6108 merged to main as an empty commit (df8e153) because apache#6156 had already made the same change to the helper, so this ports the captureWritePlan hunk of dd68a53. captureWritePlan registered its QueryExecutionListener while an earlier write's end event could still be queued on the asynchronous listener bus, and returned the first plan it saw. A test that seeds data with Comet disabled could then assert on the seed write's plan. The helper now drains the bus before registering and after the write. On branch-1.0 the helper lives in CometParquetWriterSuite, not in CometParquetWriterTestBase. The added and removed lines match dd68a53.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
captureWritePlancould mistake a delayed setup-write callback for the write under test. - Design approach: Drain Spark’s listener bus before registration and after the write, replacing the polling loop.
- Correctness / compatibility analysis: Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 confirm that draining waits for queued events and callbacks in progress.
AtomicReferencesafely publishes the captured plan across threads. - Key design decisions: Reuses
CometListenerBusUtilsand preserves listener removal infinally, without introducing another synchronization abstraction. - Implementation sketch: Drain, register, write, drain, read the captured plan, then unregister.
- Behavioral changes worth calling out: Waiting now uses Spark’s bounded 10-second drain instead of the 15-second polling loop. Synchronization overhead is confined to tests. Production execution and SQL semantics are unchanged.
- Suggested improvements: None at P1/P2 severity. No introduced P1/P2 issues found within this review.
Reviewed the entire diff from 5ffd67464f974b56a8e8b9229edebc444d1057c3 to 6ee7c6a4ddc2fea69adf3b1cd6b3636fc12a9b34. The PR is not a draft. Existing reviews, comments and threads were empty.
Routed skills: review-comet-pr, loaded from another local repository checkout because this release-branch head lacks .ai/skills. No subsystem sibling applies.
Validation: A bounded Spark 4.1.3/JDK 21 reproduction using the base and head helpers reproduced stale-plan capture at base and captured the intended write at head. git diff --check passed. The full Comet suite was not run locally because this checkout has no built native/JVM artifacts. Other Spark versions received source-level validation.
Exact-head CI: Preflight, change detection, title and labeling checks passed. Three lint jobs remained queued. Ten conditional checks were skipped. No failures were reported, but the build/test matrix had not completed.
Backport of #6108 to
branch-1.0.#6108 merged to
mainas an empty commit (df8e153). #6156 had already made the same change tocaptureWritePlanon 09-23, when this flake failed the Spark 4.2 nightly. So this PR ports thecaptureWritePlanhunk of #6156 (dd68a53). The rest of #6156 corrects a Spark 3.4 fallback that #6041 added, and #6041 is not onbranch-1.0.branch-1.1needs nothing. It was cut after #6156, and itsCometParquetWriterTestBase.scalais byte-identical to the one onmain.Which issue does this PR close?
None. #6075 is closed on
mainby #6108.Rationale for this change
branch-1.0still has the flaky helper.captureWritePlanregisters aQueryExecutionListenerwithout first draining the asynchronous listener bus, then returns the first matching plan it sees. If an earlier write's end event is still queued, the helper returns that write's plan, and the test reports that the write under test fell back to Spark. That is the failure described in #6075.On
branch-1.0, 21 of the suite's 30 tests write their seed or input data without the native writer just before the write under test:writeComplexTypeData.SaveModetests, throughmaterializeAsCometSource. Two of them also seed the target.createTestData.What changes are included in this PR?
captureWritePlaninCometParquetWriterSuitenow drains the listener bus withCometListenerBusUtils.waitUntilEmpty, once before it registers its listener and again after the write. It holds the plan in anAtomicReferenceand no longer polls for up to 15 s.The only adaptation is the file. On
mainthe helper lives inCometParquetWriterTestBase, which #5821 introduced and which is not onbranch-1.0. The added and removed lines are identical to #6156's hunk.CometListenerBusUtilsis already onbranch-1.0.How are these changes tested?
This changes a test helper only. I ran it locally on
branch-1.0with the default Spark 4.1 profile and JDK 17:CometParquetWriterSuitepass. Spotless, scalastyle and scalafix (-Psemanticdb -Pspark-3.5 -Pscala-2.12) are clean.branch-1.0with a throwaway test that isn't included here. It adds aSparkListenerthat sleeps 2 s on eachSparkListenerSQLExecutionStart, so the shared listener queue falls behind the writes. It then runs the body ofSaveMode.Append adds new files alongside existing data. With the old helper it fails withExpected exactly one CometNativeWriteExec in the plan, but found 0, because the captured plan is the Comet-disabledmaterializeAsCometSourcewrite tosource.parquet. With this change it passes.The PR build runs
CometParquetWriterSuiteon all five Spark profiles on Linux and on macOS.