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]