diff --git a/src/distributed_ext.rs b/src/distributed_ext.rs index c8ccf042..63218f25 100644 --- a/src/distributed_ext.rs +++ b/src/distributed_ext.rs @@ -386,11 +386,14 @@ pub trait DistributedExt: Sized { /// Enables or disables the use of distributed dynamic filters across network boundaries. This /// does not affect dynamic filter pushdown intra-stage. - fn with_distributed_dynamic_filters_used(self, enabled: bool) -> Result; + fn with_distributed_remote_dynamic_filters( + self, + enabled: bool, + ) -> Result; - /// Same as [`DistributedExt::with_distributed_dynamic_filters_used`] but with an in-place + /// Same as [`DistributedExt::with_distributed_remote_dynamic_filters`] but with an in-place /// mutation. - fn set_distributed_dynamic_filters_used( + fn set_distributed_remote_dynamic_filters( &mut self, enabled: bool, ) -> Result<(), DataFusionError>; @@ -852,7 +855,7 @@ impl DistributedExt for SessionConfig { Ok(()) } - fn set_distributed_dynamic_filters_used( + fn set_distributed_remote_dynamic_filters( &mut self, enabled: bool, ) -> Result<(), DataFusionError> { @@ -1028,9 +1031,9 @@ impl DistributedExt for SessionConfig { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; - #[call(set_distributed_dynamic_filters_used)] + #[call(set_distributed_remote_dynamic_filters)] #[expr($?;Ok(self))] - fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + fn with_distributed_remote_dynamic_filters(mut self, enabled: bool) -> Result; #[call(set_distributed_children_isolator_unions)] #[expr($?;Ok(self))] @@ -1167,10 +1170,10 @@ impl DistributedExt for SessionStateBuilder { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; - fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; - #[call(set_distributed_dynamic_filters_used)] + fn set_distributed_remote_dynamic_filters(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_remote_dynamic_filters)] #[expr($?;Ok(self))] - fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + fn with_distributed_remote_dynamic_filters(mut self, enabled: bool) -> Result; fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] @@ -1332,10 +1335,10 @@ impl DistributedExt for SessionState { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result; - fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; - #[call(set_distributed_dynamic_filters_used)] + fn set_distributed_remote_dynamic_filters(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_remote_dynamic_filters)] #[expr($?;Ok(self))] - fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result; + fn with_distributed_remote_dynamic_filters(mut self, enabled: bool) -> Result; fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] @@ -1490,10 +1493,10 @@ impl DistributedExt for SessionContext { #[expr($?;Ok(self))] fn with_distributed_dynamic_filter_collection(self, enabled: bool) -> Result; - fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>; - #[call(set_distributed_dynamic_filters_used)] + fn set_distributed_remote_dynamic_filters(&mut self, enabled: bool) -> Result<(), DataFusionError>; + #[call(set_distributed_remote_dynamic_filters)] #[expr($?;Ok(self))] - fn with_distributed_dynamic_filters_used(self, enabled: bool) -> Result; + fn with_distributed_remote_dynamic_filters(self, enabled: bool) -> Result; fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>; #[call(set_distributed_children_isolator_unions)] diff --git a/src/dynamic_filtering/discovery.rs b/src/dynamic_filtering/discovery.rs index ccddd5cf..b1a7f05a 100644 --- a/src/dynamic_filtering/discovery.rs +++ b/src/dynamic_filtering/discovery.rs @@ -371,7 +371,7 @@ mod tests { let (ctx, _guard, _) = start_localhost_context(2, DefaultSessionBuilder).await; let ctx = ctx .with_distributed_broadcast_joins(false)? - .with_distributed_dynamic_filters_used(remote_dynamic_filters)? + .with_distributed_remote_dynamic_filters(remote_dynamic_filters)? .with_distributed_route_task_handler(captured_plans.clone()); { let state = ctx.state_ref(); diff --git a/tests/dynamic_filtering/common.rs b/tests/dynamic_filtering/common.rs index c0c4a6d5..104a76eb 100644 --- a/tests/dynamic_filtering/common.rs +++ b/tests/dynamic_filtering/common.rs @@ -88,7 +88,7 @@ impl<'a> TestQuery<'a> { let mut ctx = ctx .with_distributed_broadcast_joins(self.broadcast_joins)? .with_distributed_dynamic_filter_collection(self.collect_dynamic_filters)? - .with_distributed_dynamic_filters_used(self.remote_dynamic_filters)?; + .with_distributed_remote_dynamic_filters(self.remote_dynamic_filters)?; if self.one_task_per_leaf { ctx = ctx.with_distributed_desired_task_count_handler(1usize); }