viirya commented on code in PR #6695:
URL: https://github.com/apache/datafusion-comet/pull/6695#discussion_r4211562197
##########
native/core/src/execution/jni_api.rs:
##########
@@ -3381,4 +3522,175 @@ mod tests {
"{error}"
);
}
+
+ /// A serialized plan whose only operator is a range of `num_elements`
longs from 0, which
+ /// needs no JVM input.
+ fn range_plan(num_elements: i64) -> Vec<u8> {
+ use datafusion_comet_proto::spark_operator::RangeScan;
+ use prost::Message;
+
+ Operator {
+ plan_id: 1,
+ op_struct: Some(OpStruct::RangeScan(RangeScan {
+ start: 0,
+ step: 1,
+ num_elements,
+ num_slices: 1,
+ })),
+ ..Default::default()
+ }
+ .encode_to_vec()
+ }
+
+ /// Creates a plan without a JVM: no metrics node or JVM input, and a
memory pool backed by a
+ /// fake Spark.
+ fn create_test_plan(serialized_plan: Vec<u8>) ->
CometResult<Box<ExecutionContext>> {
+ use crate::execution::memory_pools::create_memory_pool_with_fake_spark;
+ use datafusion_comet_proto::spark_config::ConfigMap;
+ use prost::Message;
+
+ let configs = ConfigMap {
+ entries: HashMap::from([(SPARK_EXECUTOR_CORES.to_string(),
"2".to_string())]),
+ };
+ let inputs = PlanInputs {
+ id: -6201,
+ settings: prepare_plan_settings(&configs.encode_to_vec())?,
+ serialized_plan,
+ partition_count: 1,
+ batch_size: 4,
+ off_heap_mode: true,
+ memory_pool_type: "greedy_unified".to_string(),
+ memory_limit: 1 << 30,
+ task_attempt_id: -6201,
+ task_cpus: 1,
+ local_dirs:
vec![std::env::temp_dir().to_str().unwrap().to_string()],
+ metrics_update_interval: 0,
+ start: Instant::now(),
+ jvm: PlanJvmRefs {
+ metrics: None,
+ input_sources: vec![],
+ key_unwrapper: None,
+ task_context: None,
+ class_loader: None,
+ },
+ };
+ create_execution_context(inputs, |config| {
+ create_memory_pool_with_fake_spark(config, -6201, 1 << 30)
+ })
+ }
+
+ #[test]
+ fn plan_lifecycle_without_a_jvm() {
+ let _guard = serial();
+ let mut exec_context = create_test_plan(range_plan(10)).unwrap();
+ assert!(
+ exec_context.root_op.is_none(),
+ "operators are created on first execution"
+ );
+
+ // Without an update interval, metrics are published only when the
plan is released.
+ let mut publishes = 0;
+ let mut values = Vec::new();
+ while let Some(batch) = execute_plan(&mut exec_context, 0, 0, &mut |_|
{
+ publishes += 1;
+ Ok(())
+ })
+ .unwrap()
+ {
+ assert!(
+ batch.num_rows() <= 4,
+ "batches follow the configured batch size"
+ );
+ let column =
+
arrow::array::AsArray::as_primitive::<arrow::datatypes::Int64Type>(batch.column(0));
+ values.extend(column.values().iter().copied());
+ }
+ assert_eq!(values, (0..10).collect::<Vec<i64>>());
+ assert_eq!(publishes, 0);
+ assert!(
+ exec_context.batch_producer.is_some(),
+ "a plan without JVM input runs on a Tokio task"
+ );
+
+ let plan_memory = Arc::clone(&exec_context.plan_memory);
+ release_plan(exec_context, &mut |ctx| {
+ assert!(
+ ctx.root_op.is_some(),
+ "the final metrics are of the executed plan"
+ );
+ publishes += 1;
+ Ok(())
+ })
+ .unwrap();
+ assert_eq!(publishes, 1);
+ assert_eq!(plan_memory.reserved(), 0);
Review 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_memory` asserts `plan_memory.reserved() >
0` before `release_plan` and `0` after, and that the final metrics are
published once.
- `release_plan_reports_a_failed_metrics_publish` uses 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.
--
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]