Repository navigation
Conversation
Separate the native logic of createPlan, setShufflePartitionPusher, executePlan and releasePlan from their JNI argument handling: - create_execution_context builds the execution context from plain values (PlanInputs) and a memory pool factory. prepare_plan_settings parses the Spark configs and initializes the Tokio runtime. - register_shuffle_partition_pusher checks the plan state before building the pusher. - execute_plan returns the plan's next batch; the export still exports it. - release_plan stops and drops the plan. These plans call back into the JVM, so the JVM objects stay JNI global references in the execution context (PlanJvmRefs), and the steps that need an Env are passed in as callbacks: the memory pool factory, the pusher factory, and PublishMetrics. ExecutionContext::metrics becomes an Option, which is always set by createPlan, so a plan can be created and run without a JVM in unit tests. Exception classes and messages are unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Plan lifecycle logic lived inside JNI exports, making it difficult to test without a JVM.
- Design approach: Extract creation, execution, shuffle callback registration and release into internal Rust functions, with callbacks for JVM-dependent operations.
- Correctness / compatibility analysis: JNI signatures, execution paths, Arrow export, metrics cadence and cleanup ordering are preserved. Checked Spark task-context and completion-listener semantics across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: Global references retain ownership across calls. Rejected shuffle registrations still avoid constructing a callback. The small input structs and callbacks provide useful test boundaries without adding batch copies or per-row JNI calls.
- Implementation sketch:
PlanSettings,PlanInputsandPlanJvmRefsfeedcreate_execution_context.execute_planreturns batches for the existing JNI exporter, andrelease_planpreserves cleanup even when metrics publication fails. - Behavioral changes worth calling out: JNI conversions now precede plan deserialization, so simultaneous creation failures can report a different first error. Compared with
branch-1.1, cancellation detection, metrics throttling and reservation-release waiting are inherited changes, not introduced here. - Suggested improvements: None meet the requested P1/P2 reporting threshold.
Reviewed full SHA eddc691be21972957c1d215b370479df7eb37ef0 against base 1aa166fde858b2bd20ac1dfa83eb2c16e2f6208c. The complete PR diff changes jni_api.rs. Verified that the additional cache/Iceberg differences in a direct tree comparison originate exclusively from two newer base commits. Existing reviews, comments and threads were empty.
Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-memory-pr and review-comet-shuffle-pr. Also consulted review-comet-iceberg-write-pr while checking the base-only differences.
Exact-head CI: 38 successful checks, 31 skipped and one cancelled. Rust CI recorded 2,089 passing tests, including all four new tests. Linux Comet suites, TPC checks and all nine Spark 4.1 SQL shards passed. The label-run aggregate passed. The main workflow reports failure with its Required Checks job cancelled, so its aggregate verdict is not green.
Validation: All 35 focused execution::jni_api::tests passed locally with --no-default-features. The default-feature build failed because this environment lacks jni.h for HDFS. No local JVM suite or performance benchmark was run. macOS, other Spark runtime profiles and Iceberg suites were skipped in exact-head CI.
andygrove
left a comment
There was a problem hiding this comment.
I built the merge preview (main plus this PR) and the jni_api unit tests pass, including repeated runs under load. One question on the lifecycle test inline.
| }) | ||
| .unwrap(); | ||
| assert_eq!(publishes, 1); | ||
| assert_eq!(plan_memory.reserved(), 0); |
There was a problem hiding this comment.
The last assertion in this test, assert_eq!(plan_memory.reserved(), 0), can't fail for this plan. RangeScan runs on LazyMemoryExec, which never reserves memory, so reserved() is 0 before the plan starts, while it runs and after it is released, whatever release_plan does. To check, I changed release_plan to skip BatchProducer::stop and mem::forget the context. All four new tests still passed.
Release is the step where a mistake costs the most, because Spark can hand a task's memory to another task as soon as the task ends. Could the lifecycle test run a plan that holds a reservation, for example a Sort over the RangeScan, and assert plan_memory.reserved() > 0 before release_plan and 0 after? With that, the leaking release_plan fails the test. The sort holds nothing after its first batch at the 10 rows used here (batch size 4) and held 704 bytes with 100 rows, so the input needs to be larger. The > 0 check before the release would also tell a later reader if the plan stops reserving.
The same plan could also let release_plan_reports_a_failed_metrics_publish check that a failed publish still frees the plan.
There was a problem hiding this comment.
You're right, thanks for checking it with the leaking release_plan. The release tests now run a descending Sort over a RangeScan of 10,000 rows and release it after its first batch, while the sort still holds its buffered input:
release_plan_returns_the_plans_memoryassertsplan_memory.reserved() > 0beforerelease_planand0after, and that the final metrics are published once.release_plan_reports_a_failed_metrics_publishuses the same plan and also checks that the plan's memory is returned when the publish fails.
Draining a plan to its end moved to a separate execute_plan_drains_a_plan test on the plain RangeScan. With your variant of release_plan (no BatchProducer::stop, mem::forget of the context), both release tests now fail with the sort's reservation still held.
…tion The lifecycle test ran a RangeScan, which never reserves memory, so its check that nothing stays reserved after release_plan could not fail. Run a descending sort over the range instead, release it after its first batch while it still holds its buffered input, and check that the plan holds a reservation before the release and none after. Do the same when the final metrics publish fails. Draining a plan to the end stays covered by a separate test.
Which issue does this PR close?
Part of #6434.
Rationale for this change
#6435 and #6534 separated the native logic of the JNI entry points that do not call back into the
JVM from their JNI argument handling. This does the same for the plan entry points:
createPlan,setShufflePartitionPusher,executePlanandreleasePlan. Their logic was spread across theJNI exports, so a plan could not be created, run and released in a Rust unit test.
What changes are included in this PR?
Each export keeps converting its JNI arguments and calls a core function with no
Envor JNIlocal reference in its signature:
createPlan:prepare_plan_settingsparses the Spark configs and initializes the Tokio runtime.create_execution_context(PlanInputs, memory_pool_for)deserializes the plan and sets up itsmemory pool, DataFusion session and tracing registration.
setShufflePartitionPusher:register_shuffle_partition_pusher(ctx, make_pusher)checks the planstate, then builds the pusher, so a rejected call still creates no global reference.
executePlan:execute_plan(ctx, stage_id, partition, publish_metrics)returns the plan's nextbatch, or
Noneat the end. The export still exports the batch throughprepare_outputinsidethe same trace span.
releasePlan:release_plan(ctx, publish_metrics)stops the plan, publishes its final metrics,drops it and waits for its reservations.
These plans call back into the JVM, so the JVM objects remain JNI global references held by the
execution context (
PlanJvmRefs). The steps that need anEnvare passed in as callbacks: thememory pool factory (the pool acquires memory from Spark's task memory manager), the pusher
factory, and
PublishMetrics, which publishes to theCometMetricNode.ExecutionContext::metricsbecomes anOption.createPlanalways sets it, andNoneis onlyused without a JVM.
Exception classes and messages are unchanged. In
createPlanthe JNI conversions now all happenbefore the plan is deserialized, so when two steps fail at once, a different one of the two errors
can be reported. Successful plan creation is unchanged.
How are these changes tested?
RangeScanplan is created, executed to the end and released. The test checks the values,the batch size, that the plan runs on a Tokio task, that metrics are published once on
release, and that no memory stays reserved.
release_planreports a failed metrics publish.built.
CometExecSuite,CometNativeShuffleSuite,CelebornShufflePartitionPusherSuite,CometTaskMetricsSuiteandParquetEncryptionITCaselocally.This pull request and its description were written by Isaac.