andygrove commented on code in PR #6496:
URL: https://github.com/apache/datafusion-comet/pull/6496#discussion_r4159547749


##########
native/core/src/parquet/parquet_exec.rs:
##########
@@ -231,13 +231,29 @@ pub(crate) fn init_datasource_exec(
         _ => Arc::new(parquet_source),
     };
 
-    let expr_adapter_factory: Arc<dyn PhysicalExprAdapterFactory> = Arc::new(
-        SparkPhysicalExprAdapterFactory::new(spark_parquet_options, 
default_values),
-    );
+    let spark_adapter_factory = Arc::new(SparkPhysicalExprAdapterFactory::new(
+        spark_parquet_options,
+        default_values,
+    ));
+    let expr_adapter_factory: Arc<dyn PhysicalExprAdapterFactory> =
+        Arc::<SparkPhysicalExprAdapterFactory>::clone(&spark_adapter_factory);
 
     let file_groups = file_groups
-        .iter()
-        .map(|files| FileGroup::new(files.clone()))
+        .into_iter()
+        .map(|files| {
+            FileGroup::new(
+                files
+                    .into_iter()
+                    .map(|mut file| {
+                        // Retain the concrete adapter for runtime-filter 
read-safety checks.
+                        // Extensions are keyed by type; existing access plans 
stay intact.
+                        file.extensions
+                            .insert_arc(Arc::clone(&spark_adapter_factory));

Review Comment:
   Putting the typed factory into every `PartitionedFile`'s extensions and 
checking pointer identity in `from_file_scan` works, but it adds an entry to 
every file of every Comet Parquet scan just to recover a type that 
`PhysicalExprAdapterFactory` erases. It is only needed because the trait has no 
`as_any`, and DataFusion main still doesn't have one.
   
   Would you be up for filing a DataFusion issue to add `as_any` (or a downcast 
helper like the one `PhysicalExpr` now has) to `PhysicalExprAdapterFactory`, 
and linking it in the comment here? Then we can drop the per-file entry once it 
lands.



##########
native/core/src/parquet/schema_adapter/runtime_filter.rs:
##########
@@ -0,0 +1,117 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Per-file read safety using the Spark adapter's resolved fields.
+
+use super::{
+    is_infallible_read_adaptation, SparkPhysicalExprAdapter, 
SparkPhysicalExprAdapterFactory,
+};
+use arrow::datatypes::{Schema, SchemaRef};
+use datafusion::common::Result;
+use datafusion::datasource::physical_plan::FileScanConfig;
+use datafusion::physical_expr::expressions::Column;
+use datafusion::physical_expr_adapter::{PhysicalExprAdapter, 
PhysicalExprAdapterFactory};
+use parquet::variant::VariantType;
+use std::collections::HashMap;
+use std::sync::Arc;
+
+impl SparkPhysicalExprAdapterFactory {
+    /// Recover a typed factory only when every file carries this scan's 
active factory.
+    /// File extensions retain the type erased by DataFusion's adapter 
interface. Checking
+    /// pointer identity keeps a replaced or custom adapter on the ordinary 
rewrite path.
+    pub(crate) fn from_file_scan(scan: &FileScanConfig) -> Option<Arc<Self>> {
+        let active = scan.expr_adapter_factory.as_ref()?;
+        let mut files = scan.file_groups.iter().flat_map(|group| 
group.files());
+        let factory = files.next()?.extensions.get_arc::<Self>()?;
+        let erased: Arc<dyn PhysicalExprAdapterFactory> = 
Arc::<Self>::clone(&factory);
+        if !Arc::ptr_eq(active, &erased) {
+            return None;
+        }
+        files
+            .all(|file| {
+                file.extensions
+                    .get_arc::<Self>()
+                    .is_some_and(|other| Arc::ptr_eq(&factory, &other))
+            })
+            .then_some(factory)
+    }
+
+    /// Create the normal per-file adapter and determine whether pruning can 
skip its reads.
+    pub(crate) fn create_with_read_safety(
+        &self,
+        logical_schema: SchemaRef,
+        physical_schema: SchemaRef,
+        read_columns: &[Column],
+    ) -> Result<(Arc<dyn PhysicalExprAdapter>, bool)> {
+        let adapter = self.create_adapter(logical_schema, 
Arc::clone(&physical_schema))?;
+        let allow_runtime_filter =
+            adapter.read_columns_are_infallible(read_columns, 
&physical_schema);
+        Ok((Arc::new(adapter), allow_runtime_filter))
+    }
+}
+
+impl SparkPhysicalExprAdapter {
+    fn read_columns_are_infallible(&self, columns: &[Column], physical_schema: 
&Schema) -> bool {
+        self.read_columns_are_direct(columns)
+            || columns.iter().all(|column| {
+                self.rewrite(Arc::new(column.clone()))
+                    .is_ok_and(|expr| is_infallible_read_adaptation(&expr, 
physical_schema))
+            })
+    }
+
+    fn read_columns_are_direct(&self, columns: &[Column]) -> bool {

Review Comment:
   This shortcut depends on `SparkPhysicalExprAdapter::rewrite` returning a 
bare `Column` for every matching non-Variant field. That holds today, but it is 
a second copy of `rewrite`'s behavior. If someone later adds another wrapper 
for matching fields, the way `wrap_direct_variant_column` was added, this would 
still report the column as infallible. Runtime pruning could then skip the row 
groups that would have raised the error #6067 was protecting.
   
   Could we add a `debug_assert!` that the full `rewrite` plus 
`is_infallible_read_adaptation` check agrees whenever this returns `true`? Then 
the existing native tests would catch any drift. A comment in `rewrite()` 
pointing back here would also help whoever touches it next.



-- 
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