From 8eaf2bbecf00093eb6d28184d817488a55716fe5 Mon Sep 17 00:00:00 2001 From: peterxcli Date: Fri, 11 Sep 2026 06:18:51 +0800 Subject: [PATCH] refactor: destructure dynamic filter and scalar subquery proto hooks --- datafusion/execution/src/async_stream.rs | 3 +- .../src/expressions/dynamic_filters/mod.rs | 56 +++++++++++++------ .../physical-expr/src/scalar_subquery.rs | 38 ++++++++----- 3 files changed, 64 insertions(+), 33 deletions(-) diff --git a/datafusion/execution/src/async_stream.rs b/datafusion/execution/src/async_stream.rs index 0462c53a0ffd3..a5e1c49ae983b 100644 --- a/datafusion/execution/src/async_stream.rs +++ b/datafusion/execution/src/async_stream.rs @@ -302,7 +302,6 @@ mod test { use crate::{async_stream, async_try_stream}; use futures::stream::FusedStream; use futures::{Stream, StreamExt, pin_mut}; - use std::assert_matches; use std::pin::Pin; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -464,7 +463,7 @@ mod test { pin_mut!(s); for i in 0..3 { - assert_matches!(tx.send(i).await, Ok(_)); + assert!(tx.send(i).await.is_ok()); assert_eq!(Some(i), s.next().await); } diff --git a/datafusion/physical-expr/src/expressions/dynamic_filters/mod.rs b/datafusion/physical-expr/src/expressions/dynamic_filters/mod.rs index 15544d09e2b56..4d384e952e0c0 100644 --- a/datafusion/physical-expr/src/expressions/dynamic_filters/mod.rs +++ b/datafusion/physical-expr/src/expressions/dynamic_filters/mod.rs @@ -606,13 +606,22 @@ impl PhysicalExpr for DynamicFilterPhysicalExpr { use datafusion_proto_models::protobuf; use datafusion_proto_models::protobuf::physical_expr_node::ExprType; - let children = self - .children + let Self { + children, + remapped_children, + current_cache: _, // Runtime cache, repopulated by current(). + inner, + state_watch: _, // Runtime channel, recreated from inner state by from_parts(). + data_type: _, // Cached test invariant, recomputed from the expression. + nullable: _, // Cached test invariant, recomputed from the expression. + } = self; + + let children = children .iter() .map(|c| ctx.encode_child(c)) .collect::>>()?; - let remapped_children = match &self.remapped_children { + let remapped_children = match remapped_children { Some(remapped) => remapped .iter() .map(|c| ctx.encode_child(c)) @@ -620,18 +629,23 @@ impl PhysicalExpr for DynamicFilterPhysicalExpr { None => vec![], }; - let inner = self.inner.read().clone(); - let inner_expr = Box::new(ctx.encode_child(&inner.expr)?); + let Inner { + expression_id, + generation, + expr, + is_complete, + } = inner.read().clone(); + let inner_expr = Box::new(ctx.encode_child(&expr)?); Ok(Some(protobuf::PhysicalExprNode { - expr_id: Some(inner.expression_id), + expr_id: Some(expression_id), expr_type: Some(ExprType::DynamicFilter(Box::new( protobuf::PhysicalDynamicFilterNode { children, remapped_children, - generation: inner.generation, + generation, inner_expr: Some(inner_expr), - is_complete: inner.is_complete, + is_complete, }, ))), })) @@ -649,28 +663,36 @@ impl DynamicFilterPhysicalExpr { proto: &datafusion_proto_models::protobuf::PhysicalExprNode, ctx: &datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx<'_>, ) -> Result> { + use datafusion_proto_models::protobuf; use datafusion_proto_models::protobuf::physical_expr_node::ExprType; - let ExprType::DynamicFilter(df) = proto.expr_type.as_ref().ok_or_else(|| { + let protobuf::PhysicalExprNode { expr_id, expr_type } = proto; + let ExprType::DynamicFilter(df) = expr_type.as_ref().ok_or_else(|| { internal_datafusion_err!("Missing expr_type in PhysicalExprNode") })? else { return Err(internal_datafusion_err!("Expected DynamicFilter expr_type")); }; + let protobuf::PhysicalDynamicFilterNode { + children, + remapped_children, + generation, + inner_expr, + is_complete, + } = df.as_ref(); // Decode original children - let children = df - .children + let children = children .iter() .map(|c| ctx.decode(c)) .collect::>>()?; // Decode remapped children (empty vec means None) - let remapped_children = if df.remapped_children.is_empty() { + let remapped_children = if remapped_children.is_empty() { None } else { Some( - df.remapped_children + remapped_children .iter() .map(|c| ctx.decode(c)) .collect::>>()?, @@ -678,13 +700,13 @@ impl DynamicFilterPhysicalExpr { }; // Decode the inner expression - let inner_expr_proto = df.inner_expr.as_ref().ok_or_else(|| { + let inner_expr_proto = inner_expr.as_ref().ok_or_else(|| { internal_datafusion_err!("Missing inner_expr in PhysicalDynamicFilterNode") })?; let inner_expr = ctx.decode(inner_expr_proto)?; // Restore the expression_id from the outer PhysicalExprNode - let expression_id = proto.expr_id.ok_or_else(|| { + let expression_id = expr_id.ok_or_else(|| { internal_datafusion_err!( "Missing expr_id in PhysicalExprNode for DynamicFilter" ) @@ -692,9 +714,9 @@ impl DynamicFilterPhysicalExpr { let inner = Inner { expression_id, - generation: df.generation, + generation: *generation, expr: inner_expr, - is_complete: df.is_complete, + is_complete: *is_complete, }; Ok(Arc::new(Self::from_parts( diff --git a/datafusion/physical-expr/src/scalar_subquery.rs b/datafusion/physical-expr/src/scalar_subquery.rs index 389baf3505279..927df8c2eeb9b 100644 --- a/datafusion/physical-expr/src/scalar_subquery.rs +++ b/datafusion/physical-expr/src/scalar_subquery.rs @@ -171,18 +171,25 @@ impl PhysicalExpr for ScalarSubqueryExpr { ) -> Result> { use datafusion_common::utils::usize_to_wire; use datafusion_proto_models::protobuf; + + let Self { + field, + index, + results: _, // Runtime state, supplied by ScalarSubqueryExec on decode. + } = self; + Ok(Some(protobuf::PhysicalExprNode { expr_id: None, expr_type: Some(protobuf::physical_expr_node::ExprType::ScalarSubquery( protobuf::PhysicalScalarSubqueryExprNode { - data_type: Some(self.field.data_type().try_into()?), - nullable: self.field.is_nullable(), + data_type: Some(field.data_type().try_into()?), + nullable: field.is_nullable(), index: usize_to_wire( - self.index.as_usize(), + index.as_usize(), "ScalarSubqueryExpr", "index", )?, - metadata: self.field.metadata().clone(), + metadata: field.metadata().clone(), }, )), })) @@ -214,22 +221,25 @@ impl ScalarSubqueryExpr { protobuf::physical_expr_node::ExprType::ScalarSubquery, "ScalarSubqueryExpr", ); - let data_type = require_proto_field( - sq.data_type.as_ref(), - "ScalarSubqueryExpr", - "data_type", - )? - .try_into()?; - let metadata = if sq.metadata.is_empty() { + let protobuf::PhysicalScalarSubqueryExprNode { + data_type, + nullable, + index, + metadata, + } = sq; + let data_type = + require_proto_field(data_type.as_ref(), "ScalarSubqueryExpr", "data_type")? + .try_into()?; + let metadata = if metadata.is_empty() { None } else { - Some(FieldMetadata::from(sq.metadata.clone())) + Some(FieldMetadata::from(metadata.clone())) }; Ok(Arc::new(ScalarSubqueryExpr::new_with_metadata( data_type, - sq.nullable, + *nullable, metadata, - SubqueryIndex::new(sq.index as usize), + SubqueryIndex::new(*index as usize), results.clone(), ))) }