Skip to content

fix: park the native scan loop instead of busy-polling while waiting on native I/O - #6092

Closed
mixermt wants to merge 3 commits into
apache:mainfrom
mixermt:fix/park-native-scan-loop
Closed

mixermt wants to merge 3 commits into
apache:mainfrom
mixermt:fix/park-native-scan-loop

Conversation

@mixermt

@mixermt mixermt commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6091.

Rationale for this change

Java_org_apache_comet_Native_executePlan runs a plan that has any JVM-fed input (a broadcast build side, CometSparkRowToColumnar, a shuffle read) through a loop that polls the stream and, on Pending, 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?

  • ScanStream and ShuffleScanStream honor the Stream contract. poll_next registers cx.waker() when the buffer is empty, and get_next_batch wakes it after refilling, through an AtomicWaker shared by the exec's clones. EOF stays buffered, so a re-poll returns Ready(None) again and get_next_batch is a no-op once the reader is drained.
  • The ScanExec path of executePlan moves its loop into next_batch(stream, on_pending): poll the stream; on Pending, refill the JVM-fed scans and check the metrics interval, then park the block_on task until a waker fires (park_until_woken, a poll_fn that 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 every Pending now carries a waker.
  • Awaiting the stream directly is still not an option: the JVM refill has to run between polls, and the loop is what keeps that step reachable.
  • update_metrics_on_interval replaces the per-100-polls gate. It runs on every pending poll and once per returned batch, and the tracing log_memory_usage sample sits behind the same interval, so trace density does not depend on how often the loop turns.
  • development.md describes the pull-then-park loop.

How are these changes tested?

  • next_batch_parks_while_the_stream_waits_on_native_io: a stream pending on tokio::time::sleep for 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 real ScanExec in 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_buffered in shuffle_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 warnings is clean; all 433 core unit tests pass.
  • JVM, with the rebuilt library: CometTaskMetricsSuite, CometJoinSuite, CometNativeShuffleInputRDDSuite, CometNativeShuffleSuite and CometIcebergNativeSuite: 235 tests pass and none hang; the one canceled test is the pre-existing Spark 4.1 assume gate (SPARK-55626).
  • Not measured here: the CPU reduction on the production workload itself, which needs a run with this build.

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

…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>
@github-actions github-actions Bot added bug Something isn't working area:scan Parquet scan / data reading labels Sep 21, 2026
@mbutrovich
mbutrovich self-requested a review September 21, 2026 21:34

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Comment thread native/core/src/execution/jni_api.rs Outdated
Comment thread native/core/src/execution/jni_api.rs Outdated
Comment thread native/core/src/execution/jni_api.rs Outdated

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread native/core/src/execution/jni_api.rs Outdated
Comment thread native/core/src/execution/operators/scan.rs Outdated
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>
@mixermt

mixermt commented Sep 22, 2026

Copy link
Copy Markdown
Contributor Author

Thanks both. Everything below is in the latest push, and the two threads changed the design rather than patching it.

  • Wakers instead of a timeout (@mbutrovich). ScanStream and ShuffleScanStream now register cx.waker() when the buffer is empty, and get_next_batch wakes it after refilling, through an AtomicWaker shared by the exec's clones. Every Pending carries a waker, so the 100 ms timeout, the bool return from get_next_batch and the tokio time feature are gone; park_until_woken is a poll_fn that yields once. I dropped the timeout rather than logging its firings: with wakers in place a firing could not be told apart from a legitimately slow read, so a counter would not have found a lost wake.
  • EOF is sticky (@andygrove). poll_next leaves InputBatch::EOF in the buffer, so a re-poll returns Ready(None) again and get_next_batch stays a no-op instead of making another JNI round trip into a drained reader. The resting-state invariant in Native executePlan busy-polls at 100% CPU while a plan with a JVM-fed input waits on native I/O #6091 is now what the code does.
  • A loop test that fails without the fix (@mbutrovich). The poll, pull and park steps moved into next_batch(stream, on_pending); update_metrics and prepare_output stay in executePlan. next_batch_parks_while_the_stream_waits_on_native_io drives it with a stream pending on tokio::time::sleep and asserts the pull closure ran a few times; with the park deleted it ran 135,311 times in one 50 ms wait. next_batch_resumes_on_a_refill_and_stops_pulling_after_eof runs a real ScanExec in 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 again it fails. refill_wakes_the_pending_poll_and_eof_stays_buffered covers ShuffleScanStream.
  • Tracing density (@andygrove). log_memory_usage sits behind the same interval check as update_metrics, in update_metrics_on_interval, which runs on every pending poll and once per returned batch, so in-loop trace density no longer depends on how often the loop turns.
  • Docs (@mbutrovich). The JVM data source paragraph in development.md describes the pull-then-park loop. poll_fn is imported next to task::Poll and the Park struct is gone.

Verified with the rebuilt library: clippy with -D warnings is clean, the 433 core unit tests pass, and CometTaskMetricsSuite, CometJoinSuite, CometNativeShuffleInputRDDSuite, CometNativeShuffleSuite and CometIcebergNativeSuite pass (235 tests, no hangs).

This rewrites the native execution loop, so the Spark SQL suites should report here rather than in the merge queue. Could a committer add run-spark-4.1-tests?

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 andygrove added this to the 1.1.0 milestone Sep 22, 2026

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@andygrove andygrove added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 24, 2026
@andygrove

Copy link
Copy Markdown
Member

@mixermt bringing back the pulled flag and the nested block_on test from my last review is the last change I'm looking for here, and I'd like this to make 1.1.0 before we cut the branch. I've approved the CI workflows and added run-spark-4.1-tests, so your next push will run the Spark SQL suites too. If you won't get to it by Friday, would you mind if I push that change to your branch? It would stay your PR.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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_batch to poll, refill and park. Share an AtomicWaker across scan clones, wake after refilling, and retain EOF in both scan streams.
  • Correctness / compatibility analysis: The previously reported nested-block_on hang remains at native/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's FileFormatDataWriter and InterruptibleIterator sources 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: executePlan temporarily takes its stream, calls next_batch, restores the stream before propagating errors, and updates metrics through update_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.

@andygrove

Copy link
Copy Markdown
Member

Thanks for this @mixermt. I'd like to get this into 1.1.0, so I branched from your branch and fixed the one remaining issue in #6219. I hope you don't mind.

@andygrove andygrove closed this Sep 25, 2026
andygrove added a commit that referenced this pull request Sep 28, 2026
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.
LinSimon-901101 pushed a commit to LinSimon-901101/datafusion-comet that referenced this pull request Sep 29, 2026
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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:scan Parquet scan / data reading bug Something isn't working run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native executePlan busy-polls at 100% CPU while a plan with a JVM-fed input waits on native I/O

4 participants