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]

Reply via email to