fix: [branch-1.0] check each fair_unified reservation against its own share (#6205) - #6228
Conversation
…e#6205) * fix: check each fair_unified reservation against its own share Since the DataFusion 53 upgrade (apache#3629), CometFairMemoryPool::try_grow has compared the pool's total reserved bytes against pool_size / num_consumers. That caps the whole task at one consumer's share, and every new registration tightens the ceiling on the consumers already running. The comment that justified the change holds for shrink but not for try_grow. MemoryReservation::try_grow calls pool.try_grow before adding to its size, so reservation.size() is the size before the request. Check the reservation against its share again, and separately keep the pool's total within pool_size, since the shares alone do not bound it once a consumer registers late or one consumer holds sibling reservations. Update the tuning guide, the memory management guide and the memory PR review skill, which described the pool-wide comparison. Closes apache#5961. * fix: charge sibling fair_unified reservations to their consumer's share Checking each reservation against the share let an operator that holds several reservations under one consumer, such as a sort's streaming merge with one per batch, take other consumers' shares, bounded only by the pool total. Keep a running total for each consumer id instead, so all of a consumer's reservations count against one share. The pool no longer reads reservation.size(), so the check also stops depending on when DataFusion updates it. Note in the tuning guide that tasks can reserve more before they spill than they could in Comet 0.15.0 through 1.0.0. (cherry picked from commit e900844) Adapted for branch-1.0: - Dropped the changes to the memory management contributor guide and the review-comet-memory-pr skill. Neither file exists on branch-1.0; they come from apache#5933 and apache#6018. - tuning.md: the new paragraph on 0.15.0 through 1.0.0 no longer links to "Sizing the Overhead from the Memory Usage Log". That section and the native memory usage log it describes come from apache#6162, which is not on branch-1.0, so the link had no target.
|
This may be too big a change to backport to 1.0.x ... I need to run benchmarks first |
After testing this change in 1.1 branch, I believe that we should backport it to 1.0 to fix one of the major causes of OOM |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedcompared total pool reservations against one consumer’s share, causing premature allocation refusals and spilling. - Design approach: Track bytes by consumer ID, enforce each consumer’s share, and separately enforce the pool-total limit.
- Correctness / compatibility analysis: DataFusion 54.1’s reservation lifecycle supports this accounting. Both pool wrappers forward registration correctly. Spark’s acquire/release contracts were checked against sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Spark retains its task-level arbitration, and partial grants remain correctly returned.
- Key design decisions:
new_empty(),split()andtake()share one consumer quota. The total limit handles late registrations without reclaiming existing reservations. Infalliblegrowand overcommit repayment remain intact. - Implementation sketch: A consumer-ID map replaces the registration counter under the existing mutex. This adds map bookkeeping without another locking layer or additional JNI calls per accepted request. Four regression tests cover the accounting change.
- Behavioral changes worth calling out: Tasks with multiple consumers can reserve more before spilling. The tuning guide documents this change and the need to reassess executor headroom. This change does not bound actual process RSS.
- Suggested improvements: No introduced P1/P2 issues found within this review. No substantiated existing P1/P2 blocker remains.
Reviewed full SHA 47b98a8a9396d501e4a4dedffaff3467451f58c6 against base f037304537166c93319ad6fe3c935442478b9e63, covering the entire one-commit, two-file diff. The PR is not a draft. Both discussion comments were read. The snapshot contains no reviews, inline comments or review threads.
Routed skills: review-comet-pr and review-comet-memory-pr, loaded from the local 6268-c86c6a3a4676 Comet checkout because this branch omits .ai/skills.
Exact-head CI: 62 checks succeeded and nine were skipped, with no failed or pending check runs. The Rust job reports 877 tests passed and four skipped, including all five fair-pool tests. Successful checks also cover Linux Comet suites across Spark 3.4–4.2, macOS, and Spark SQL 3.5/4.1. Benchmark and Spark SQL 3.4/4.0 jobs were skipped.
Validation: A disposable harness using this head’s pool code, DataFusion 54.1, both wrappers and FakeSpark passed 13 tests, including concurrent sibling admissions and partial-grant/overcommit accounting. Running the current fair-pool tests against the base implementation reproduced three expected failures. The harness stubs JNI. No full local Comet/JVM build, end-to-end Spark workload, benchmark or RSS measurement was performed. Project source remains unchanged.
Backport of #6205 to
branch-1.0.Cherry-picked from
e90084420ab7acc4d1fd3a6db41f9ef9cb5a5087. The change tofair_pool.rsis byte-identical to upstream. Two docs files were dropped and one link was cut from the tuning guide, as described under "What changes are included" below.Which issue does this PR close?
Closes #5961 on
branch-1.0.Rationale for this change
The bug ships in 1.0.0 through the same code as on
mainbefore #6205.fair_pool.rsonbranch-1.0is identical tomain's copy just before #6205. Itstry_growcompares the pool's total reserved bytes againstpool_size / num_consumers, so all of a task's operators together are capped at one operator's share, and every new registration lowers that cap for the consumers already running. The comparison came in with the DataFusion 53 upgrade (#3629), so every release from 0.15.0 through 1.0.0 has it.fair_unifiedis the default off-heap pool onbranch-1.0.The cap costs the most on executors that run few tasks at once, where Spark's own limit on each task is loosest. See #6205 for the TPC-H measurements on
main, and the caveat about them under "How are these changes tested" below.What changes are included in this PR?
The fix is the original one. #6205's description covers only its first commit. The second commit charges all of a consumer's reservations against one share, and it is described in the commit message. The adaptations:
docs/source/contributor-guide/memory_management.mdand.ai/skills/review-comet-memory-pr/SKILL.mdare dropped, because neither file exists onbranch-1.0. They come from docs: add contributor guide page on memory management #5933 and docs: split the PR review skill by area and correct the shuffle contributor docs #6018.docs/source/user-guide/latest/tuning.md, the new paragraph about 0.15.0 through 1.0.0 ends at "check that executors still have enough headroom." Onmainit goes on to link to "Sizing the Overhead from the Memory Usage Log". That section, and the executor memory usage log it describes, come from feat: always count native allocations and log executor native memory usage #6162, which is not onbranch-1.0, so the link would have had no target. The rest of the tuning guide change is upstream's, and its wording holds here: onbranch-1.0the pool is also shared by all of a task's native plans, sonum_consumerscounts the consumers of all of them.How are these changes tested?
Same tests as the original PR, run locally on
branch-1.0:fair_pool.rspass, along with the rest of thedatafusion-cometlib tests: 169 passed and 4 ignored.branch-1.0, and the tests catch it. Withbranch-1.0's currentCometFairMemoryPooland the new tests kept, three of them fail:each_consumer_is_limited_to_its_own_share,a_consumer_registered_late_is_limited_by_the_pool_totalandsibling_reservations_draw_on_one_share. The fourth,split_and_take_keep_the_bytes_on_their_consumer, passes either way, because it guards the per-consumer ledger that fix: check each fair_unified reservation against its own share #6205 introduces rather than the bug.MemoryConsumer::id(), and the pool expects every reservation's consumer to have registered with it. Onbranch-1.0the fair pool sits under DataFusion 54.1'sTrackConsumersPool, and under Comet'sLoggingMemoryPoolwhenspark.comet.debug.memoryis on, and both forwardregisterandunregisterto it. DataFusion 54.1'sMemoryReservationalso behaves as fix: check each fair_unified reservation against its own share #6205 relies on:try_growcalls the pool before adding to the reservation's size,split,takeandnew_emptykeep the reservation's consumer, and a consumer unregisters only after its last reservation is dropped.cargo fmt --all -- --check,cargo clippy --all-targets --workspace -- -D warningsandprettier --checkontuning.mdpass.The JVM suites run this pool through JNI, since
CometTestBaseenables 2 GiB of off-heap memory. That path has no Rust unit tests, so it is left to CI. The TPC-H SF10 runs in #6205's description were made with its first commit, before the second commit moved the share check from each reservation to each consumer, and I did not repeat them onbranch-1.0.Are there any user-facing changes?
The same as #6205, with no config or API changes.
fair_unifiednow refuses atry_grow, without asking Spark, only when the requesting consumer would go over itspool_size / num_consumersshare or the pool's total would go overpool_size. The pool's own checks never refuse a request that 1.0.0's would accept, so tasks with several operators can reserve more memory before they spill, up to what Spark grants the task. That is worth a line in the 1.0.x release notes, because deployments may have sized executor memory against the accidental cap. The tuning guide change says the same.