viirya opened a new pull request, #6695:
URL: https://github.com/apache/datafusion-comet/pull/6695

   ## 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`, `executePlan` and `releasePlan`. Their logic 
was spread across the
   JNI 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 `Env` or JNI
   local reference in its signature:
   
   - `createPlan`: `prepare_plan_settings` parses the Spark configs and 
initializes the Tokio runtime.
     `create_execution_context(PlanInputs, memory_pool_for)` deserializes the 
plan and sets up its
     memory pool, DataFusion session and tracing registration.
   - `setShufflePartitionPusher`: `register_shuffle_partition_pusher(ctx, 
make_pusher)` checks the plan
     state, 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 next
     batch, or `None` at the end. The export still exports the batch through 
`prepare_output` inside
     the 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 an `Env` are passed 
in as callbacks: the
   memory pool factory (the pool acquires memory from Spark's task memory 
manager), the pusher
   factory, and `PublishMetrics`, which publishes to the `CometMetricNode`.
   `ExecutionContext::metrics` becomes an `Option`. `createPlan` always sets 
it, and `None` is only
   used without a JVM.
   
   Exception classes and messages are unchanged. In `createPlan` the JNI 
conversions now all happen
   before 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?
   
   - New Rust unit tests without a JVM:
     - A `RangeScan` plan 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_plan` reports a failed metrics publish.
     - A malformed plan is rejected.
     - A shuffle pusher can only be registered once, before execution, and a 
rejected pusher is not
       built.
   - Ran `CometExecSuite`, `CometNativeShuffleSuite`, 
`CelebornShufflePartitionPusherSuite`,
     `CometTaskMetricsSuite` and `ParquetEncryptionITCase` locally.
   
   This pull request and its description were written by Isaac.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to