Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions src/coordinator/distributed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use crate::coordinator::prepare_static_plan::prepare_static_plan;
use crate::coordinator::query_coordinator::QueryCoordinator;
use crate::coordinator::store::{Store, task_keys_for_plan};
use crate::dynamic_filtering::{
is_dynamic_filtering_enabled, sever_dynamic_filter_relationships_in_plan_for_display,
is_local_dynamic_filtering_enabled, sever_dynamic_filter_relationships_in_plan_for_display,
};
use crate::{DistributedConfig, TaskCompletedDynamicFilters, TaskKey, TaskMetrics};
use datafusion::common::internal_datafusion_err;
Expand Down Expand Up @@ -251,7 +251,8 @@ impl ExecutionPlan for DistributedExec {
false => prepare_static_plan(&query_coordinator, &base_plan).await?,
};

let dynamic_filtering_enabled = is_dynamic_filtering_enabled(context.session_config());
let dynamic_filtering_enabled =
is_local_dynamic_filtering_enabled(context.session_config());
prepared.plan_for_viz = match dynamic_filtering_enabled && collect_dynamic_filters {
true => sever_dynamic_filter_relationships_in_plan_for_display(
prepared.plan_for_viz,
Expand Down
11 changes: 9 additions & 2 deletions src/coordinator/prepare_dynamic_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,9 @@ use crate::distributed_planner::{
InjectNetworkBoundaryContext, NetworkBoundaryBuilderResult, ProducerHead, calculate_cost,
inject_network_boundaries,
};
use crate::dynamic_filtering::orphan_dynamic_filter_consumers;
use crate::dynamic_filtering::{
is_remote_dynamic_filtering_enabled, orphan_dynamic_filter_consumers,
};
use crate::events::TaskCountAnnotation::{Desired, Maximum};
use crate::execution_plans::SamplerExec;
use crate::stage::{LocalStage, RemoteStage};
Expand Down Expand Up @@ -82,7 +84,12 @@ pub(super) async fn prepare_dynamic_plan(
// In order to infer the compute the cost of the stage above this one, here a sampler
// is injected to gather runtime statistics.
input_stage.plan = ProducerHead::insert_sampler(input_stage.plan)?;
let dynamic_filter_anchors = orphan_dynamic_filter_consumers(&input_stage.plan)?;
let dynamic_filter_anchors =
if is_remote_dynamic_filtering_enabled(query_coordinator.session_config())? {
orphan_dynamic_filter_consumers(&input_stage.plan)?
} else {
vec![]
};

let mut load_info_rxs = Vec::with_capacity(input_stage.tasks);

Expand Down
11 changes: 9 additions & 2 deletions src/coordinator/prepare_static_plan.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
use crate::common::TreeNodeExt;
use crate::coordinator::distributed::PreparedPlan;
use crate::coordinator::query_coordinator::QueryCoordinator;
use crate::dynamic_filtering::orphan_dynamic_filter_consumers;
use crate::dynamic_filtering::{
is_remote_dynamic_filtering_enabled, orphan_dynamic_filter_consumers,
};
use crate::stage::RemoteStage;
use crate::{NetworkBoundaryExt, Stage};
use datafusion::common::tree_node::Transformed;
Expand Down Expand Up @@ -32,7 +34,12 @@ pub(super) async fn prepare_static_plan(
let Stage::Local(stage) = plan.input_stage() else {
return exec_err!("Input stage from network boundary was not in Local state");
};
let dynamic_filter_anchors = orphan_dynamic_filter_consumers(&stage.plan)?;
let dynamic_filter_anchors =
if is_remote_dynamic_filtering_enabled(query_coordinator.session_config())? {
orphan_dynamic_filter_consumers(&stage.plan)?
} else {
vec![]
};

let mut stage_coordinator = query_coordinator.stage_coordinator(stage);
let mut futures = Vec::with_capacity(stage.tasks);
Expand Down
23 changes: 14 additions & 9 deletions src/coordinator/query_coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@ use crate::coordinator::DynamicFilterRegistry;
use crate::coordinator::Store;
use crate::coordinator::latency_metric::LatencyMetric;
use crate::dynamic_filtering::{
dynamic_filter_remote_producer_ids, is_dynamic_filtering_enabled,
dynamic_filter_remote_producer_ids, is_local_dynamic_filtering_enabled,
is_remote_dynamic_filtering_enabled,
maybe_roundtrip_plan_to_sever_in_memory_dynamic_filter_relationships,
};
use crate::events::{
Expand Down Expand Up @@ -176,8 +177,10 @@ impl<'a> StageCoordinator<'a> {
task_number: task_i,
};

self.dynamic_filter_registry
.register_task(&plan, task_key)?;
if is_remote_dynamic_filtering_enabled(session_config)? {
self.dynamic_filter_registry
.register_task(&plan, task_key)?;
}

let mut headers = get_config_extension_propagation_headers(session_config)?;
headers.extend(get_passthrough_headers(session_config));
Expand Down Expand Up @@ -433,7 +436,8 @@ impl<'a> StageCoordinator<'a> {
let wuf_registry = session_config
.get_extension::<WorkUnitFeedRegistry>()
.unwrap_or_default();
let dynamic_filtering_enabled = is_dynamic_filtering_enabled(session_config);
let dynamic_filtering_enabled = is_local_dynamic_filtering_enabled(session_config);
let remote_dynamic_filtering_enabled = is_remote_dynamic_filtering_enabled(session_config)?;

let mut work_unit_feed_declarations = vec![];
let d_ctx = DistributedTaskContext {
Expand Down Expand Up @@ -491,11 +495,12 @@ impl<'a> StageCoordinator<'a> {
} else {
transformed.data
};
let dynamic_filter_remote_producer_ids = if dynamic_filtering_enabled {
dynamic_filter_remote_producer_ids(&plan)?
} else {
vec![]
};
let dynamic_filter_remote_producer_ids =
if dynamic_filtering_enabled && remote_dynamic_filtering_enabled {
dynamic_filter_remote_producer_ids(&plan)?
} else {
vec![]
};
Ok(TaskSpecializedPlan {
plan,
work_unit_feed_declarations,
Expand Down
39 changes: 39 additions & 0 deletions src/distributed_ext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -384,6 +384,17 @@ pub trait DistributedExt: Sized {
enabled: bool,
) -> Result<(), DataFusionError>;

/// 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<Self, DataFusionError>;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The convention is to name these after the the DistributedConfig property it's controlling.

As this is controlling remote_dynamic_filters this should be called:

  • with_distributed_remote_dynamic_filters
  • set_distributed_remote_dynamic_filters

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed by #755


/// Same as [`DistributedExt::with_distributed_dynamic_filters_used`] but with an in-place
/// mutation.
fn set_distributed_dynamic_filters_used(
&mut self,
enabled: bool,
) -> Result<(), DataFusionError>;

/// Enables children isolator unions for distributing UNION operations across as many tasks as
/// the sum of all the tasks required for each child.
///
Expand Down Expand Up @@ -841,6 +852,15 @@ impl DistributedExt for SessionConfig {
Ok(())
}

fn set_distributed_dynamic_filters_used(
&mut self,
enabled: bool,
) -> Result<(), DataFusionError> {
let d_cfg = DistributedConfig::from_config_options_mut(self.options_mut())?;
d_cfg.remote_dynamic_filters = enabled;
Ok(())
}

fn set_distributed_children_isolator_unions(
&mut self,
enabled: bool,
Expand Down Expand Up @@ -1008,6 +1028,10 @@ impl DistributedExt for SessionConfig {
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result<Self, DataFusionError>;

#[call(set_distributed_dynamic_filters_used)]
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result<Self, DataFusionError>;

#[call(set_distributed_children_isolator_unions)]
#[expr($?;Ok(self))]
fn with_distributed_children_isolator_unions(mut self, enabled: bool) -> Result<Self, DataFusionError>;
Expand Down Expand Up @@ -1143,6 +1167,11 @@ impl DistributedExt for SessionStateBuilder {
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_dynamic_filters_used)]
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_children_isolator_unions)]
#[expr($?;Ok(self))]
Expand Down Expand Up @@ -1303,6 +1332,11 @@ impl DistributedExt for SessionState {
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filter_collection(mut self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_dynamic_filters_used)]
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filters_used(mut self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_children_isolator_unions)]
#[expr($?;Ok(self))]
Expand Down Expand Up @@ -1456,6 +1490,11 @@ impl DistributedExt for SessionContext {
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filter_collection(self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_dynamic_filters_used(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_dynamic_filters_used)]
#[expr($?;Ok(self))]
fn with_distributed_dynamic_filters_used(self, enabled: bool) -> Result<Self, DataFusionError>;

fn set_distributed_children_isolator_unions(&mut self, enabled: bool) -> Result<(), DataFusionError>;
#[call(set_distributed_children_isolator_unions)]
#[expr($?;Ok(self))]
Expand Down
3 changes: 3 additions & 0 deletions src/distributed_planner/distributed_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,9 @@ extensions_options! {
/// the distributed plan. This does not control whether dynamic filtering is used during
/// query execution.
pub collect_dynamic_filters: bool, default = true
/// Control whether dynamic filters are distributed across network boundaries. This does
/// not control intra-stage filter behavior.
pub remote_dynamic_filters: bool, default = true
/// Enable broadcast joins for CollectLeft hash joins. When enabled, the build side of
/// a CollectLeft join is broadcast to all consumer tasks.
pub broadcast_joins: bool, default = true
Expand Down
55 changes: 54 additions & 1 deletion src/dynamic_filtering/discovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,7 @@ mod tests {
) other_build ON other_build.key = probe."RainTomorrow"
WHERE probe."MinTemp" > 0
"#,
true,
)
.await?;
assert_snapshot!(display, @r"
Expand Down Expand Up @@ -285,6 +286,7 @@ mod tests {
GROUP BY "RainTomorrow"
) probe ON build.key = probe.key
"#,
true,
)
.await?;
assert_snapshot!(display, @r"
Expand Down Expand Up @@ -314,11 +316,62 @@ mod tests {
Ok(())
}

async fn display_query(sql: &str) -> Result<String> {
/// Disabling remote dynamic filters removes the anchors and remote producers, but the
/// task-local producers and consumers are left intact.
#[tokio::test]
async fn skips_anchors_when_remote_dynamic_filters_are_disabled() -> Result<()> {
let display = display_query(
r#"
SELECT COUNT(*)
FROM (
SELECT DISTINCT "RainToday" AS key
FROM weather
) build
JOIN weather probe ON build.key = probe."RainToday"
JOIN (
SELECT DISTINCT "RainTomorrow" AS key
FROM weather
) other_build ON other_build.key = probe."RainTomorrow"
WHERE probe."MinTemp" > 0
"#,
false,
)
.await?;
assert_snapshot!(display, @r"
Stage 5
AggregateExec
HashJoinExec producers=[1]
NetworkShuffleExec
AggregateExec
NetworkShuffleExec
Stage 4
RepartitionExec
AggregateExec
DataSourceExec consumers=[1]
Stage 3
RepartitionExec
HashJoinExec producers=[2]
NetworkShuffleExec
AggregateExec
NetworkShuffleExec
Stage 2
RepartitionExec
AggregateExec
DataSourceExec consumers=[2]
Stage 1
RepartitionExec
FilterExec
DataSourceExec
");
Ok(())
}

async fn display_query(sql: &str, remote_dynamic_filters: bool) -> Result<String> {
let captured_plans = CapturePlans::default();
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_route_task_handler(captured_plans.clone());
{
let state = ctx.state_ref();
Expand Down
10 changes: 9 additions & 1 deletion src/dynamic_filtering/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
mod discovery;
mod display;

use crate::DistributedConfig;
use crate::codec::roundtrip_pb;
use datafusion::common::Result;
use datafusion::execution::TaskContext;
Expand All @@ -12,13 +13,20 @@ pub(crate) use discovery::*;
pub use display::rewrite_distributed_plan_with_dynamic_filters;
pub(crate) use display::sever_dynamic_filter_relationships_in_plan_for_display;

pub(crate) fn is_dynamic_filtering_enabled(session_config: &SessionConfig) -> bool {
pub(crate) fn is_local_dynamic_filtering_enabled(session_config: &SessionConfig) -> bool {
session_config
.options()
.optimizer
.enable_dynamic_filter_pushdown
}

pub(crate) fn is_remote_dynamic_filtering_enabled(session_config: &SessionConfig) -> Result<bool> {
let remote_enabled =
DistributedConfig::from_session_config(session_config)?.remote_dynamic_filters;
let local_enabled = is_local_dynamic_filtering_enabled(session_config);
Ok(remote_enabled && local_enabled)
}

/// Deepcopies the plan if it contains any dynamic filter producers or consumers. This isolates
/// any dynamic filters in this plan from dynamic filters *outside* the plan. Plan nodes
/// *within* this plan will share in-memory dynamic filter state with eachother, even after copying.
Expand Down
Loading
Loading