andygrove commented on code in PR #6695:
URL: https://github.com/apache/datafusion-comet/pull/6695#discussion_r4210931454
##########
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:
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.
--
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]