Conversation
…on native I/O executePlan's ScanExec path re-polled the plan's stream in a tight loop whenever it returned Pending, relying on pull_input_batches to block on the JVM iterators in between. Once every JVM-fed scan holds a batch or has reached EOF that pull is a no-op, so a stream pending on native I/O (a Parquet or Iceberg scan reading from S3 or HDFS) spun the executor thread at 100% CPU for the whole read. A broadcast hash join over a native scan hits this on every probe-side read. pull_input_batches and the two scan operators now report whether a buffer was refilled. When the stream is Pending and nothing was pulled, the loop parks the block_on task until a waker registered by that poll fires, bounded by a short safety timeout, then re-enters the loop so JVM-fed scans still get refilled. Awaiting the stream directly is not safe: ScanExec returns Pending without a waker when an operator drains and re-polls it within one poll. The metrics interval is checked every iteration now that iterations are no longer spins. Closes apache#6091 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
mbutrovich
left a comment
There was a problem hiding this comment.
Thanks @mixermt. The diagnosis in #6091 closes the gap #3553 left open. #3553 parked the executor thread only for plans without JVM-fed inputs, so plans with a ScanExec or ShuffleScanExec still spin while native I/O is pending. You're right that the loop can't .await the stream, because the JVM refill has to run between polls. Parking for one wake-up keeps that refill reachable.
The threading section of development.md (lines 46-48) describes this loop. Can you update it in this PR to say the loop parks until a waker fires when nothing was pulled?
andygrove
left a comment
There was a problem hiding this comment.
Thanks for the writeup on this one, the diagnosis in #6091 made it easy to follow. I traced through the alternative you rejected and I agree it isn't safe. If the loop awaits the stream, the wake can land inside the await, an operator drains a JVM buffer and re-polls that scan within the same poll, and the no-waker Pending that comes back has no loop left to rescue it. Parking for exactly one wake-up is the right shape. I also checked that build_runtime already calls .enable_all(), so the timer driver is there for the park, and that spark.comet.metrics.updateInterval defaults to 3000 ms, comfortably above the 100 ms park, so metrics stay timely.
Two things on top of Matt's comments.
ScanStream and ShuffleScanStream now register the poll's waker when their buffer is empty and get_next_batch wakes it after refilling, so every Pending from the plan carries a waker. The loop in executePlan moves into next_batch(stream, on_pending): poll, refill and check metrics on Pending, then park until a waker fires. The 100 ms timeout, the bool from get_next_batch and the tokio time feature are gone. EOF stays buffered so a re-poll of an exhausted scan returns Ready(None) again instead of another JNI round trip. The tracing memory sample sits behind the metrics interval, and development.md describes the loop. Tests drive next_batch with a stream pending on a sleep (with the park removed it pulls 135,311 times in 50 ms) and with a ScanExec refilled by the pull closure under a timeout, plus a ShuffleScanStream waker test. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
|
Thanks both. Everything below is in the latest push, and the two threads changed the design rather than patching it.
Verified with the rebuilt library: clippy with This rewrites the native execution loop, so the Spark SQL suites should report here rather than in the merge queue. Could a committer add |
The JVM data source path polls operators on the Spark executor thread inside block_on; only tasks they spawn run on tokio workers. The heading now names ShuffleScanExec as well, since pull_input_batches feeds both streams and both register a waker. The native-wait test gets the same ten second bound as the refill test, so a lost wake fails instead of hanging the suite. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
andygrove
left a comment
There was a problem hiding this comment.
The waker registration and the sticky EOF resolve both of my earlier threads, thanks. I found one more case, and I think it means the pulled flag from the first version needs to come back alongside the wakers.
A refill can run another Comet native plan on the same thread. CometNativeWriteExec and CometIcebergWriteExec feed child.executeColumnar() through CometArrowStream.inputObjects, so when the child is native, ScanExec::get_next_batch ends up in the child's executePlan, which calls block_on again inside the block_in_place. Every block_on on a thread shares tokio's one thread-local parker, and an unpark leaves a single token. If the outer plan is waiting on I/O that completes on another thread while the nested plan is parked, the nested park consumes that token, and when the pull returns, park_until_woken parks with nothing left to wake it. The refill's own wake doesn't cover this when the refilled scan wasn't polled empty in that round. That's the shape of the writer loop in ParquetWriterExec, which takes a batch and then awaits the write, so the scan is refilled with no waker registered. Without the timeout this becomes a hang rather than a slow wait.
I checked the mechanism outside Comet with a copy of next_batch whose pull closure runs a nested Handle::block_on(sleep(100ms)) while the stream waits on a 20 ms sleep. With the loop always parking, it only returned when an unrelated 5 s timer fired. Skipping the park after a pull brought it back to about 100 ms.
Could get_next_batch go back to reporting whether it made a JNI call, with next_batch parking only when nothing was pulled? A real pull is the only way the JVM runs native code on this thread, so that closes the gap and still fixes #6091, where the spin comes from the no-op pull. The wakers are still worth keeping, so both streams follow the Stream contract and the park needs no timeout. A test that runs a nested block_on inside the pull closure and asserts on elapsed time would guard it, since within_ten_seconds lets a single lost wake pass slowly.
|
@mixermt bringing back the |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Plans combining JVM-fed inputs with asynchronous native I/O could busy-poll after input refills became no-ops, consuming an executor core while waiting.
- Design approach: Extract
next_batchto poll, refill and park. Share anAtomicWakeracross scan clones, wake after refilling, and retain EOF in both scan streams. - Correctness / compatibility analysis: The previously reported nested-
block_onhang remains atnative/core/src/execution/jni_api.rs:1006. A refill can execute a child native plan that consumes the outer thread's wake token. Unconditional parking then stalls despite completed I/O. I reproduced this independently and traced the native writer input paths. Spark'sFileFormatDataWriterandInterruptibleIteratorsources were checked across supported versions 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No additional introduced P1/P2 issues found within this review. - Key design decisions: Keeping JVM pulls outside stream polling preserves the execution boundary. Waker registration under the batch lock and waking after unlocking avoid a registration/refill race. Buffer ownership and the one-batch buffering model remain intact.
- Implementation sketch:
executePlantemporarily takes its stream, callsnext_batch, restores the stream before propagating errors, and updates metrics throughupdate_metrics_on_interval. - Behavioral changes worth calling out: EOF no longer causes repeated upstream reads. Native I/O waits park, and in-loop memory tracing follows the metrics interval. The threading documentation correctly distinguishes executor-thread polling from spawned worker tasks.
- Suggested improvements: Resolve the existing lost-wakeup review by reporting whether a refill pulled input and skipping parking after a real pull. Add the already-requested nested-execution test with an elapsed-time assertion. This is an existing blocker, not a new duplicate finding.
Reviewed all four changed files and all three commits from base 09b44ad6fa17f58f4bbccf958c5cf02d790f7334 through head 013d566fa7b8cb45ea36681e9fbb800fe0d8c9e6. Read the discussion and resolved threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-shuffle-pr, and review-comet-memory-pr.
Exact-head CI: Comet CI passed Linux builds, Comet suites, TPC-H/TPC-DS checks and Rust tests. The Rust log confirms all three added tests passed and reports 1,597 tests passed, five skipped. Spark SQL, Iceberg and macOS jobs were skipped. The subsequent Spark-label run reports action_required with no jobs, so it provides no Spark SQL verdict.
Validation: A standalone harness using the exact helper bodies and locked Tokio 1.53.1 reproduced the hang with a 20 ms outer wait and 100 ms nested wait. The current helper exceeded an external two-second timeout. The busy-loop control and proposed pulled-input guard both completed in about 102 ms. Local validation used substitute batch/error types, not a complete Comet/JVM build. No end-to-end writer reproduction or production CPU benchmark was run. The checkout remains unchanged.
Regenerate the 1.1.0 changelog with the generator from #6358, for the same range as #6282 (1.0.0..021c378), and format it with prettier 3.9.9. The generator used for #6282 credited each PR only to the author of its merge commit, so it left out Michael Taranov, whose commits from #6092 are in #6219, and credited the backports only to the person who opened them. The PR lines for #5310, #6219, #6223 and the four backports now name their co-authors, and the credits count each co-authored PR, which adds Michael Taranov and brings the contributor count to 41.
generate-changelog.py credited each PR only to the author of the commit that merged it, both on the PR's line and in the credits, so anyone else whose commits were squashed into the PR went uncredited. apache#6219 carries the commits from apache#6092, and the 1.1.0 changelog left their author out. For a merge commit with a Co-authored-by trailer, the script now fetches the PR's commits and credits the GitHub account of each commit's author. It skips merges of the base branch, commits whose email is not linked to an account, and accounts that have never opened an issue or PR in the repository. GitHub links a commit to whichever account has its email address, so a placeholder identity or an AI agent's commits can otherwise be credited to an unrelated account. The script prints each account it skips for that reason. In the credits, a co-author counts once for each such PR, under the name git shortlog already lists them under, and counts toward the number of contributors in the header. As in git shortlog, a name counts once per PR, so someone who committed from two accounts is not counted twice.
Which issue does this PR close?
Closes #6091.
Rationale for this change
Java_org_apache_comet_Native_executePlanruns a plan that has any JVM-fed input (a broadcast build side,CometSparkRowToColumnar, a shuffle read) through a loop that polls the stream and, onPending, pulls the next batches from the JVM. That pull blocks inside JNI while the JVM produces data, so the loop never spun as long as the JVM was the only thing worth waiting for. With native scans reading from S3 or HDFS the stream is also pending on asynchronous I/O, and once every JVM-fed scan holds a batch or has reached EOF the pull is a no-op: the loop re-polls at full speed for the duration of every read, pinning one core per task.On a production workload (Iceberg on HDFS joined with a broadcast relation) the scan stages used 87 core-hours against 6.6 for Spark alone, ran 3x to 6x longer, and the saturated cores caused HDFS ack timeouts and retries. Details in #6091.
What changes are included in this PR?
ScanStreamandShuffleScanStreamhonor theStreamcontract.poll_nextregisterscx.waker()when the buffer is empty, andget_next_batchwakes it after refilling, through anAtomicWakershared by the exec's clones. EOF stays buffered, so a re-poll returnsReady(None)again andget_next_batchis a no-op once the reader is drained.executePlanmoves its loop intonext_batch(stream, on_pending): poll the stream; onPending, refill the JVM-fed scans and check the metrics interval, then park theblock_ontask until a waker fires (park_until_woken, apoll_fnthat yields once). A refill wakes the task before it parks, so it resumes at once; a stream waiting on native I/O sleeps until that I/O wakes it. There is no timeout, because everyPendingnow carries a waker.update_metrics_on_intervalreplaces the per-100-polls gate. It runs on every pending poll and once per returned batch, and the tracinglog_memory_usagesample sits behind the same interval, so trace density does not depend on how often the loop turns.development.mddescribes the pull-then-park loop.How are these changes tested?
next_batch_parks_while_the_stream_waits_on_native_io: a stream pending ontokio::time::sleepfor 50 ms; the pull closure runs once. With the park removed it ran 135,311 times.next_batch_resumes_on_a_refill_and_stops_pulling_after_eof: a realScanExecin test mode. Each park ends only on the refill's wake, under a timeout that turns a lost wake into a failure, and a re-poll after EOF pulls nothing. With EOF cleared on poll it fails.refill_wakes_the_pending_poll_and_eof_stays_bufferedinshuffle_scan.rs: the empty-buffer poll registers the waker, the refill wakes it, and EOF stays buffered.cargo clippy --all-targets -p datafusion-comet -- -D warningsis clean; all 433 core unit tests pass.CometTaskMetricsSuite,CometJoinSuite,CometNativeShuffleInputRDDSuite,CometNativeShuffleSuiteandCometIcebergNativeSuite: 235 tests pass and none hang; the one canceled test is the pre-existing Spark 4.1assumegate (SPARK-55626).This touches the native execution loop, so the Spark SQL suites (
run-spark-4.1-tests) should run before merge.AI Disclosure
Drafted, implemented and tested with AI assistance (Claude Code); reviewed before submission.
🤖 Generated with Claude Code