From a1457c1dc4138b4fdf990ef284ab5a0ba250f900 Mon Sep 17 00:00:00 2001 From: discord9 <55937128+discord9@users.noreply.github.com> Date: Thu, 10 Sep 2026 19:40:36 +0800 Subject: [PATCH] feat(query): expose per-query dynamic filter controls Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com> --- src/operator/src/statement.rs | 9 +- src/operator/src/statement/set.rs | 78 +++++++- src/query/src/datafusion.rs | 150 ++++++++++++++- src/query/src/dist_plan/merge_scan.rs | 174 +++++++++++++++-- src/query/src/options.rs | 175 +++++++++++++++++- src/query/src/query_engine/context.rs | 5 + .../src/query_engine/default_serializer.rs | 93 +++++++++- src/query/src/sql.rs | 23 +++ src/servers/src/grpc/greptime_handler.rs | 16 +- src/servers/src/http/hints.rs | 16 +- src/session/src/context.rs | 116 +++++++++++- 11 files changed, 830 insertions(+), 25 deletions(-) diff --git a/src/operator/src/statement.rs b/src/operator/src/statement.rs index 21c3ed6146..8cf8e73d8d 100644 --- a/src/operator/src/statement.rs +++ b/src/operator/src/statement.rs @@ -63,7 +63,7 @@ use query::QueryEngineRef; use query::parser::QueryStatement; use session::context::{Channel, QueryContextBuilder, QueryContextRef}; use session::table_name::table_idents_to_full_name; -use set::{set_query_timeout, set_read_preference}; +use set::{set_dynamic_filter_pushdown, set_query_timeout, set_read_preference}; use snafu::{OptionExt, ResultExt, ensure}; use sql::ast::ObjectNamePartExt; use sql::statements::OptionMap; @@ -550,6 +550,13 @@ impl StatementExecutor { // Allow query to fallback when failed to push down. "ALLOW_QUERY_FALLBACK" => set_allow_query_fallback(set_var.value, query_ctx)?, + "ENABLE_DYNAMIC_FILTER_PUSHDOWN" + | "ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN" + | "ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN" + | "ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN" => { + set_dynamic_filter_pushdown(&var_name.to_lowercase(), set_var.value, query_ctx)? + } + "CLIENT_ENCODING" => validate_client_encoding(set_var)?, "@@SESSION.MAX_EXECUTION_TIME" | "MAX_EXECUTION_TIME" => match query_ctx.channel() { Channel::Mysql => set_query_timeout(set_var.value, query_ctx)?, diff --git a/src/operator/src/statement/set.rs b/src/operator/src/statement/set.rs index b0305dbdae..f3ad92018a 100644 --- a/src/operator/src/statement/set.rs +++ b/src/operator/src/statement/set.rs @@ -20,7 +20,10 @@ use lazy_static::lazy_static; use regex::Regex; use session::ReadPreference; use session::context::Channel::Postgres; -use session::context::QueryContextRef; +use session::context::{ + ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN, ENABLE_DYNAMIC_FILTER_PUSHDOWN, + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN, ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, QueryContextRef, +}; use session::session_config::{PGByteaOutputValue, PGDateOrder, PGDateTimeStyle, PGIntervalStyle}; use snafu::{OptionExt, ResultExt, ensure}; use sql::ast::{Expr, Ident, Value}; @@ -249,6 +252,38 @@ pub fn set_allow_query_fallback(exprs: Vec, ctx: QueryContextRef) -> Resul } } +pub fn set_dynamic_filter_pushdown( + name: &str, + exprs: Vec, + ctx: QueryContextRef, +) -> Result<()> { + let Some((Expr::Value(value), [])) = exprs.split_first() else { + return NotSupportedSnafu { + feat: format!("Set variable value must be one boolean for {name}"), + } + .fail(); + }; + let value = match &value.value { + Value::Boolean(value) => *value, + _ => { + return NotSupportedSnafu { + feat: format!("Set variable value must be a boolean for {name}"), + } + .fail(); + } + }; + debug_assert!(matches!( + name, + ENABLE_DYNAMIC_FILTER_PUSHDOWN + | ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN + | ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN + | ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN + )); + ctx.configuration_parameter() + .set_dynamic_filter_pushdown(name, value); + Ok(()) +} + pub fn set_intervalstyle(exprs: Vec, ctx: QueryContextRef) -> Result<()> { let Some((var_value, [])) = exprs.split_first() else { return NotSupportedSnafu { @@ -381,7 +416,46 @@ fn parse_pg_query_timeout_input(input: &str) -> Result { #[cfg(test)] mod test { - use crate::statement::set::parse_pg_query_timeout_input; + use std::sync::Arc; + + use session::context::{ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, QueryContextBuilder}; + use sql::ast::{Expr, Value}; + + use super::{parse_pg_query_timeout_input, set_dynamic_filter_pushdown}; + + #[test] + fn test_set_dynamic_filter_pushdown_requires_boolean() { + let ctx = Arc::new(QueryContextBuilder::default().build()); + set_dynamic_filter_pushdown( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, + vec![Expr::Value(Value::Boolean(false).into())], + ctx.clone(), + ) + .unwrap(); + assert_eq!( + ctx.configuration_parameter() + .dynamic_filter_pushdown() + .enable_topk_dynamic_filter_pushdown, + Some(false) + ); + + assert!( + set_dynamic_filter_pushdown( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, + vec![Expr::Value( + Value::SingleQuotedString("false".to_string()).into() + )], + ctx.clone(), + ) + .is_err() + ); + assert_eq!( + ctx.configuration_parameter() + .dynamic_filter_pushdown() + .enable_topk_dynamic_filter_pushdown, + Some(false) + ); + } #[test] fn test_parse_pg_query_timeout_input() { diff --git a/src/query/src/datafusion.rs b/src/query/src/datafusion.rs index f9d87782a0..ec28d48eb1 100644 --- a/src/query/src/datafusion.rs +++ b/src/query/src/datafusion.rs @@ -68,7 +68,7 @@ use crate::metrics::{ OnDone, QUERY_STAGE_ELAPSED, maybe_attach_region_watermark_metrics, should_collect_region_watermark_from_query_ctx, }; -use crate::options::ScheduledTimeExtension; +use crate::options::{ScheduledTimeExtension, apply_dynamic_filter_pushdown_options}; use crate::physical_wrapper::PhysicalPlanWrapperRef; use crate::planner::{DfLogicalPlanner, LogicalPlanner}; use crate::query_engine::{DescribeResult, QueryEngineContext, QueryEngineState}; @@ -457,10 +457,18 @@ impl DatafusionQueryEngine { }; let _timer = metrics::CREATE_PHYSICAL_ELAPSED.start_timer(); - let state = ctx.state(); + let query_ctx = ctx.query_ctx(); + let state = ctx.state_mut(); common_telemetry::debug!("Create physical plan, input plan: {logical_plan}"); + apply_dynamic_filter_pushdown_options(state.config_mut().options_mut(), &query_ctx)?; + let config_options = state.config_options().clone(); + let _ = state + .execution_props_mut() + .config_options + .insert(config_options); + // special handle EXPLAIN plan if matches!(logical_plan, DfLogicalPlan::Explain(_)) { return state @@ -1247,6 +1255,144 @@ mod tests { ); } + fn assert_dynamic_filter_pushdown_disabled(state: &datafusion::execution::SessionState) { + let optimizer = &state.config_options().optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + + let snapshot = state.execution_props().config_options.as_ref().unwrap(); + let optimizer = &snapshot.optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + } + + #[tokio::test] + async fn test_dynamic_filter_pushdown_hint_overrides_session_options() { + let engine = create_test_engine().await; + let engine = engine + .as_any() + .downcast_ref::() + .unwrap(); + let configuration_parameter = Arc::new(session::context::ConfigurationVariables::default()); + configuration_parameter + .set_dynamic_filter_pushdown("enable_topk_dynamic_filter_pushdown", false); + let query_ctx = Arc::new( + QueryContextBuilder::default() + .configuration_parameter(configuration_parameter) + .set_extension( + "enable_topk_dynamic_filter_pushdown".to_string(), + "true".to_string(), + ) + .build(), + ); + let mut engine_ctx = engine.engine_context(query_ctx); + let plan = datafusion_expr::LogicalPlanBuilder::empty(true) + .build() + .unwrap(); + let explain = datafusion_expr::LogicalPlanBuilder::from(plan) + .explain(false, false) + .unwrap() + .build() + .unwrap(); + engine + .create_physical_plan(&mut engine_ctx, &explain) + .await + .unwrap(); + let optimizer = &engine_ctx.state().config_options().optimizer; + + assert!(optimizer.enable_dynamic_filter_pushdown); + assert!(optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(optimizer.enable_join_dynamic_filter_pushdown); + assert!(optimizer.enable_topk_dynamic_filter_pushdown); + } + + #[tokio::test] + async fn test_dynamic_filter_pushdown_master_hint_applies_to_regular_and_explain_plans() { + let engine = create_test_engine().await; + let engine = engine + .as_any() + .downcast_ref::() + .unwrap(); + + for explain in [false, true] { + let query_ctx = Arc::new( + QueryContextBuilder::default() + .set_extension( + "enable_dynamic_filter_pushdown".to_string(), + "false".to_string(), + ) + .set_extension( + "enable_topk_dynamic_filter_pushdown".to_string(), + "true".to_string(), + ) + .build(), + ); + let mut engine_ctx = engine.engine_context(query_ctx); + let plan = datafusion_expr::LogicalPlanBuilder::empty(true) + .build() + .unwrap(); + let plan = if explain { + datafusion_expr::LogicalPlanBuilder::from(plan) + .explain(false, false) + .unwrap() + .build() + .unwrap() + } else { + plan + }; + + engine + .create_physical_plan(&mut engine_ctx, &plan) + .await + .unwrap(); + assert_dynamic_filter_pushdown_disabled(engine_ctx.state()); + } + } + + #[tokio::test] + async fn test_invalid_dynamic_filter_hint_rejects_regular_and_explain_plans() { + let engine = create_test_engine().await; + let engine = engine + .as_any() + .downcast_ref::() + .unwrap(); + + for explain in [false, true] { + let query_ctx = Arc::new( + QueryContextBuilder::default() + .set_extension( + "enable_dynamic_filter_pushdown".to_string(), + "invalid".to_string(), + ) + .build(), + ); + let mut engine_ctx = engine.engine_context(query_ctx); + let plan = datafusion_expr::LogicalPlanBuilder::empty(true) + .build() + .unwrap(); + let plan = if explain { + datafusion_expr::LogicalPlanBuilder::from(plan) + .explain(false, false) + .unwrap() + .build() + .unwrap() + } else { + plan + }; + + assert!( + engine + .create_physical_plan(&mut engine_ctx, &plan) + .await + .is_err() + ); + } + } + #[tokio::test] async fn test_topk_dynamic_filter_pushdown_reaches_region_scan() { let engine = create_test_engine().await; diff --git a/src/query/src/dist_plan/merge_scan.rs b/src/query/src/dist_plan/merge_scan.rs index 8d3580ea45..7ac31c961d 100644 --- a/src/query/src/dist_plan/merge_scan.rs +++ b/src/query/src/dist_plan/merge_scan.rs @@ -56,6 +56,8 @@ use greptime_proto::v1::region::RegionRequestHeader; use meter_core::data::ReadItem; use meter_macros::read_meter; use session::context::{ + ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN, ENABLE_DYNAMIC_FILTER_PUSHDOWN, + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN, ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, FLIGHT_METRICS_HEARTBEAT_INTERVAL, QueryContextRef, SUPPORT_FLIGHT_METRICS_BEFORE_BATCH_EXTENSION_KEY, }; @@ -77,7 +79,10 @@ use crate::dist_plan::{ FilterId, RemoteDynFilterProducerId, RemoteDynFilterRegistryLease, Subscriber, }; use crate::metrics::{MERGE_SCAN_ERRORS_TOTAL, MERGE_SCAN_POLL_ELAPSED, MERGE_SCAN_REGIONS}; -use crate::options::{FlowQueryExtensions, remote_dyn_filter_pushdown_enabled_from_extensions}; +use crate::options::{ + FlowQueryExtensions, dynamic_filter_pushdown_options, + remote_dyn_filter_pushdown_enabled_from_extensions, +}; use crate::query_engine::QueryEngineState; use crate::region_query::RegionQueryHandlerRef; @@ -130,6 +135,39 @@ fn remote_plan_row_bound(plan: &LogicalPlan) -> Option { } } +fn materialize_dynamic_filter_pushdown_extensions( + query_ctx: &mut session::context::QueryContext, +) -> Result<()> { + let options = dynamic_filter_pushdown_options( + query_ctx + .configuration_parameter() + .dynamic_filter_pushdown(), + &query_ctx.extensions(), + ) + .map_err(|err| DataFusionError::External(Box::new(err)))?; + for (name, enabled) in [ + ( + ENABLE_DYNAMIC_FILTER_PUSHDOWN, + options.enable_dynamic_filter_pushdown, + ), + ( + ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN, + options.enable_aggregate_dynamic_filter_pushdown, + ), + ( + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN, + options.enable_join_dynamic_filter_pushdown, + ), + ( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, + options.enable_topk_dynamic_filter_pushdown, + ), + ] { + query_ctx.set_extension(name, enabled.to_string()); + } + Ok(()) +} + fn remote_dyn_filter_enabled(query_ctx: &QueryContextRef) -> Result { remote_dyn_filter_pushdown_enabled_from_extensions(&query_ctx.extensions()) .map_err(|err| DataFusionError::External(Box::new(err))) @@ -667,6 +705,7 @@ impl MergeScanExec { region_id, &captured_remote_dyn_filters, ); + materialize_dynamic_filter_pushdown_extensions(&mut region_query_ctx)?; if live_analyze_metrics { let remote_query_id = region_query_ctx.remote_query_id().map(str::to_string); if let Some(remote_query_id) = remote_query_id { @@ -1431,7 +1470,7 @@ mod tests { use datatypes::vectors::{Int64Vector, StringVector, TimestampMillisecondVector}; use futures_util::{Stream, TryStreamExt}; use session::ReadPreference; - use session::context::QueryContext; + use session::context::{QueryContext, QueryContextBuilder}; use session::query_id::QueryId; use table::table::scan::REGION_SCAN_EXEC_NAME; use table::table_name::TableName; @@ -1441,7 +1480,7 @@ mod tests { use super::*; use crate::dist_plan::DynFilterRegistryManager; use crate::options::QueryOptions; - use crate::query_engine::{QueryEngineContext, QueryEngineState}; + use crate::query_engine::{DefaultPlanDecoder, QueryEngineContext, QueryEngineState}; use crate::region_query::RegionQueryHandler; fn test_target(id: u64) -> crate::region_query::RegionQueryTarget { @@ -2039,7 +2078,18 @@ mod tests { #[tokio::test] async fn immediate_eof_do_get_receives_refreshed_remote_dyn_filter_snapshot() { let handler = Arc::new(ImmediateEofRegionQueryHandler::default()); - let query_ctx = QueryContext::arc(); + let configuration_parameter = Arc::new(session::context::ConfigurationVariables::default()); + configuration_parameter + .set_dynamic_filter_pushdown(ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, false); + let query_ctx = Arc::new( + session::context::QueryContextBuilder::default() + .configuration_parameter(configuration_parameter) + .set_extension( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "true".to_string(), + ) + .build(), + ); let state = query_engine_state(handler.clone()); let exec = remote_dyn_filter_test_exec(handler.clone(), query_ctx.clone()); let dyn_filter = install_remote_dyn_filter(&exec); @@ -2055,6 +2105,37 @@ mod tests { let snapshot = registrations.regs[0].initial_snapshot.as_ref().unwrap(); assert!(snapshot.generation > 0); assert!(!snapshot.is_complete); + assert_eq!( + handler + .dynamic_filter_extensions() + .get(ENABLE_DYNAMIC_FILTER_PUSHDOWN), + Some(&"true".to_string()) + ); + assert_eq!( + handler + .dynamic_filter_extensions() + .get(ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN), + Some(&"true".to_string()) + ); + let dn_ctx = Arc::new(QueryContext::from(api::v1::QueryContext { + extensions: handler.dynamic_filter_extensions(), + ..Default::default() + })); + let decoder = DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &dn_ctx).unwrap(); + assert!( + decoder + .session_state() + .config_options() + .optimizer + .enable_dynamic_filter_pushdown + ); + assert!( + decoder + .session_state() + .config_options() + .optimizer + .enable_topk_dynamic_filter_pushdown + ); assert_eq!( snapshot .payload @@ -2073,6 +2154,63 @@ mod tests { ); } + #[tokio::test] + async fn session_dynamic_filter_settings_reach_dn_decoder_without_hint() { + let handler = Arc::new(ImmediateEofRegionQueryHandler::default()); + let configuration_parameter = Arc::new(session::context::ConfigurationVariables::default()); + configuration_parameter.set_dynamic_filter_pushdown(ENABLE_DYNAMIC_FILTER_PUSHDOWN, false); + let query_ctx = Arc::new( + QueryContextBuilder::default() + .configuration_parameter(configuration_parameter) + .build(), + ); + let state = query_engine_state(handler.clone()); + let exec = remote_dyn_filter_test_exec(handler.clone(), query_ctx.clone()); + let dyn_filter = install_remote_dyn_filter(&exec); + dyn_filter.update(physical_lit(false) as _).unwrap(); + + let mut stream = exec + .to_stream(task_context_with_engine_state(state, query_ctx), 0) + .unwrap(); + assert!(stream.next().await.is_none()); + + let dn_ctx = Arc::new(QueryContext::from(api::v1::QueryContext { + extensions: handler.dynamic_filter_extensions(), + ..Default::default() + })); + let decoder = DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &dn_ctx).unwrap(); + let optimizer = &decoder.session_state().config_options().optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + let optimizer = &decoder + .session_state() + .execution_props() + .config_options + .as_ref() + .unwrap() + .optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + } + + #[test] + fn materialized_dynamic_filter_extensions_reject_invalid_value_on_dn() { + let dn_ctx = Arc::new( + QueryContextBuilder::default() + .set_extension( + ENABLE_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "invalid".to_string(), + ) + .build(), + ); + + assert!(DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &dn_ctx).is_err()); + } + #[test] fn remote_dyn_filter_registry_cleanup_waits_for_last_query_scoped_stream_drop() { let registry_manager = Arc::new(DynFilterRegistryManager::default()); @@ -2467,12 +2605,21 @@ mod tests { #[derive(Default)] struct ImmediateEofRegionQueryHandler { registrations: Mutex>, + dynamic_filter_extensions: Mutex>>, } impl ImmediateEofRegionQueryHandler { fn registrations(&self) -> InitialDynFilterRegs { self.registrations.lock().unwrap().clone().unwrap() } + + fn dynamic_filter_extensions(&self) -> std::collections::HashMap { + self.dynamic_filter_extensions + .lock() + .unwrap() + .clone() + .unwrap() + } } #[async_trait] @@ -2490,19 +2637,18 @@ mod tests { _target: &crate::region_query::RegionQueryTarget, request: common_query::request::QueryRequest, ) -> crate::error::Result { - let registrations = request + let extensions = request .header - .clone() - .and_then(|header| header.query_context) - .and_then(|query_context| { - query_context - .extensions - .get(INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY) - .cloned() - }) - .map(|serialized| InitialDynFilterRegs::from_extension_value(&serialized).unwrap()) + .as_ref() + .and_then(|header| header.query_context.as_ref()) + .map(|query_context| query_context.extensions.clone()) + .unwrap(); + let registrations = extensions + .get(INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY) + .map(|serialized| InitialDynFilterRegs::from_extension_value(serialized).unwrap()) .unwrap(); *self.registrations.lock().unwrap() = Some(registrations); + *self.dynamic_filter_extensions.lock().unwrap() = Some(extensions); Ok(empty_record_batch_stream(&request)) } diff --git a/src/query/src/options.rs b/src/query/src/options.rs index c998839bf2..f8a1fe920b 100644 --- a/src/query/src/options.rs +++ b/src/query/src/options.rs @@ -18,7 +18,11 @@ use chrono::{DateTime, Utc}; use common_base::memory_limit::MemoryLimit; use datafusion::config::{ConfigEntry, ConfigExtension, ExtensionOptions}; use serde::{Deserialize, Serialize}; -use session::context::QueryContextRef; +use session::context::{ + DynamicFilterPushdownOptions, ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN, + ENABLE_DYNAMIC_FILTER_PUSHDOWN, ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN, + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, QueryContextRef, +}; use store_api::storage::RegionId; use table::metadata::TableId; @@ -38,6 +42,82 @@ pub const QUERY_ENABLE_REMOTE_DYNAMIC_FILTER_PUSHDOWN: &str = pub const FLOW_INCREMENTAL_MODE_MEMTABLE_ONLY: &str = "memtable_only"; +/// Resolves dynamic-filter settings from a session and per-query extensions. +/// Query extensions take precedence over session values; disabling the master +/// switch disables all dynamic-filter variants. +pub fn dynamic_filter_pushdown_options( + session_options: DynamicFilterPushdownOptions, + extensions: &HashMap, +) -> Result { + let defaults = datafusion_common::config::ConfigOptions::default().optimizer; + let option = |name, session_value, default| { + extensions + .get(name) + .map(|value| parse_bool(name, value)) + .transpose() + .map(|value| value.or(session_value).unwrap_or(default)) + }; + let enable_dynamic_filter_pushdown = option( + ENABLE_DYNAMIC_FILTER_PUSHDOWN, + session_options.enable_dynamic_filter_pushdown, + defaults.enable_dynamic_filter_pushdown, + )?; + let enable_aggregate_dynamic_filter_pushdown = option( + ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN, + session_options.enable_aggregate_dynamic_filter_pushdown, + defaults.enable_aggregate_dynamic_filter_pushdown, + )?; + let enable_join_dynamic_filter_pushdown = option( + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN, + session_options.enable_join_dynamic_filter_pushdown, + defaults.enable_join_dynamic_filter_pushdown, + )?; + let enable_topk_dynamic_filter_pushdown = option( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, + session_options.enable_topk_dynamic_filter_pushdown, + defaults.enable_topk_dynamic_filter_pushdown, + )?; + + Ok(ResolvedDynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown, + enable_aggregate_dynamic_filter_pushdown: enable_dynamic_filter_pushdown + && enable_aggregate_dynamic_filter_pushdown, + enable_join_dynamic_filter_pushdown: enable_dynamic_filter_pushdown + && enable_join_dynamic_filter_pushdown, + enable_topk_dynamic_filter_pushdown: enable_dynamic_filter_pushdown + && enable_topk_dynamic_filter_pushdown, + }) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ResolvedDynamicFilterPushdownOptions { + pub enable_dynamic_filter_pushdown: bool, + pub enable_aggregate_dynamic_filter_pushdown: bool, + pub enable_join_dynamic_filter_pushdown: bool, + pub enable_topk_dynamic_filter_pushdown: bool, +} + +/// Resolves and applies dynamic-filter settings to a DataFusion query state. +pub fn apply_dynamic_filter_pushdown_options( + config: &mut datafusion_common::config::ConfigOptions, + query_ctx: &QueryContextRef, +) -> Result<()> { + let options = dynamic_filter_pushdown_options( + query_ctx + .configuration_parameter() + .dynamic_filter_pushdown(), + &query_ctx.extensions(), + )?; + config.optimizer.enable_dynamic_filter_pushdown = options.enable_dynamic_filter_pushdown; + config.optimizer.enable_aggregate_dynamic_filter_pushdown = + options.enable_aggregate_dynamic_filter_pushdown; + config.optimizer.enable_join_dynamic_filter_pushdown = + options.enable_join_dynamic_filter_pushdown; + config.optimizer.enable_topk_dynamic_filter_pushdown = + options.enable_topk_dynamic_filter_pushdown; + Ok(()) +} + /// Query engine config #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(default)] @@ -388,6 +468,99 @@ mod flow_extension_tests { assert_eq!(parsed, None); } + #[test] + fn test_dynamic_filter_pushdown_options_defaults_and_overrides() { + let defaults = DynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: None, + enable_aggregate_dynamic_filter_pushdown: None, + enable_join_dynamic_filter_pushdown: None, + enable_topk_dynamic_filter_pushdown: None, + }; + assert_eq!( + dynamic_filter_pushdown_options(defaults, &HashMap::new()).unwrap(), + ResolvedDynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: true, + enable_aggregate_dynamic_filter_pushdown: true, + enable_join_dynamic_filter_pushdown: true, + enable_topk_dynamic_filter_pushdown: true, + } + ); + + let extensions = HashMap::from([( + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "false".to_string(), + )]); + let options = dynamic_filter_pushdown_options(defaults, &extensions).unwrap(); + assert!(options.enable_dynamic_filter_pushdown); + assert!(!options.enable_join_dynamic_filter_pushdown); + assert!(options.enable_aggregate_dynamic_filter_pushdown); + assert!(options.enable_topk_dynamic_filter_pushdown); + } + + #[test] + fn test_dynamic_filter_pushdown_master_switch_suppresses_variants() { + let extensions = HashMap::from([ + ( + ENABLE_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "false".to_string(), + ), + ( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "true".to_string(), + ), + ]); + let options = dynamic_filter_pushdown_options( + DynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: None, + enable_aggregate_dynamic_filter_pushdown: None, + enable_join_dynamic_filter_pushdown: None, + enable_topk_dynamic_filter_pushdown: None, + }, + &extensions, + ) + .unwrap(); + + assert!(!options.enable_dynamic_filter_pushdown); + assert!(!options.enable_aggregate_dynamic_filter_pushdown); + assert!(!options.enable_join_dynamic_filter_pushdown); + assert!(!options.enable_topk_dynamic_filter_pushdown); + } + + #[test] + fn test_dynamic_filter_pushdown_options_session_master_switch_suppresses_variants() { + let session_options = DynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: Some(false), + enable_aggregate_dynamic_filter_pushdown: None, + enable_join_dynamic_filter_pushdown: None, + enable_topk_dynamic_filter_pushdown: None, + }; + let options = dynamic_filter_pushdown_options(session_options, &HashMap::new()).unwrap(); + + assert!(!options.enable_dynamic_filter_pushdown); + assert!(!options.enable_aggregate_dynamic_filter_pushdown); + assert!(!options.enable_join_dynamic_filter_pushdown); + assert!(!options.enable_topk_dynamic_filter_pushdown); + } + + #[test] + fn test_dynamic_filter_pushdown_options_reject_invalid_value() { + let extensions = HashMap::from([( + ENABLE_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "not-a-bool".to_string(), + )]); + let err = dynamic_filter_pushdown_options( + DynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: None, + enable_aggregate_dynamic_filter_pushdown: None, + enable_join_dynamic_filter_pushdown: None, + enable_topk_dynamic_filter_pushdown: None, + }, + &extensions, + ) + .unwrap_err(); + assert!(format!("{err}").contains(ENABLE_DYNAMIC_FILTER_PUSHDOWN)); + } + #[test] fn test_remote_dyn_filter_pushdown_enabled_from_extensions_defaults_true() { assert!(remote_dyn_filter_pushdown_enabled_from_extensions(&HashMap::new()).unwrap()); diff --git a/src/query/src/query_engine/context.rs b/src/query/src/query_engine/context.rs index c967fc38bd..67aa85fe4a 100644 --- a/src/query/src/query_engine/context.rs +++ b/src/query/src/query_engine/context.rs @@ -37,6 +37,11 @@ impl QueryEngineContext { &self.state } + #[inline] + pub fn state_mut(&mut self) -> &mut SessionState { + &mut self.state + } + #[inline] pub fn query_ctx(&self) -> QueryContextRef { self.query_ctx.clone() diff --git a/src/query/src/query_engine/default_serializer.rs b/src/query/src/query_engine/default_serializer.rs index e009d80954..6ba426d05c 100644 --- a/src/query/src/query_engine/default_serializer.rs +++ b/src/query/src/query_engine/default_serializer.rs @@ -50,6 +50,7 @@ use substrait::extension_serializer::ExtensionSerializer; use substrait::{DFLogicalSubstraitConvertor, SubstraitPlan}; use crate::dist_plan::MergeScanLogicalPlan; +use crate::options::apply_dynamic_filter_pushdown_options; /// Extended [`substrait::extension_serializer::ExtensionSerializer`] but supports [`MergeScanLogicalPlan`] serialization. #[derive(Debug)] @@ -104,10 +105,22 @@ pub struct DefaultPlanDecoder { } impl DefaultPlanDecoder { + #[cfg(test)] + pub(crate) fn session_state(&self) -> &SessionState { + &self.session_state + } + pub fn new( session_state: SessionState, query_ctx: &QueryContextRef, ) -> crate::error::Result { + let mut session_state = session_state; + apply_dynamic_filter_pushdown_options(session_state.config_mut().options_mut(), query_ctx)?; + let config_options = session_state.config_options().clone(); + let _ = session_state + .execution_props_mut() + .config_options + .insert(config_options); Ok(Self { session_state, query_ctx: query_ctx.clone(), @@ -252,7 +265,10 @@ mod tests { }; use datatypes::data_type::DataType; use promql::extension_plan::RangeManipulate; - use session::context::QueryContext; + use session::context::{ + ConfigurationVariables, ENABLE_DYNAMIC_FILTER_PUSHDOWN, + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, QueryContext, QueryContextBuilder, + }; use super::*; use crate::QueryEngineFactory; @@ -260,6 +276,81 @@ mod tests { use crate::optimizer::test_util::mock_table_provider; use crate::options::QueryOptions; + fn assert_dynamic_filter_pushdown_disabled(session_state: &SessionState) { + let optimizer = &session_state.config_options().optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + + let snapshot = session_state + .execution_props() + .config_options + .as_ref() + .unwrap(); + let optimizer = &snapshot.optimizer; + assert!(!optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_aggregate_dynamic_filter_pushdown); + assert!(!optimizer.enable_join_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + } + + #[test] + fn test_plan_decoder_applies_dynamic_filter_extensions() { + let session_only = Arc::new(ConfigurationVariables::default()); + session_only.set_dynamic_filter_pushdown(ENABLE_DYNAMIC_FILTER_PUSHDOWN, false); + let session_only_ctx = Arc::new( + QueryContextBuilder::default() + .configuration_parameter(session_only) + .build(), + ); + let decoder = + DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &session_only_ctx).unwrap(); + assert_dynamic_filter_pushdown_disabled(&decoder.session_state); + + let session_with_hint = Arc::new(ConfigurationVariables::default()); + session_with_hint.set_dynamic_filter_pushdown(ENABLE_DYNAMIC_FILTER_PUSHDOWN, false); + let hinted_ctx = Arc::new( + QueryContextBuilder::default() + .configuration_parameter(session_with_hint) + .set_extension( + ENABLE_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "true".to_string(), + ) + .set_extension( + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "false".to_string(), + ) + .build(), + ); + let decoder = + DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &hinted_ctx).unwrap(); + let optimizer = &decoder.session_state.config_options().optimizer; + assert!(optimizer.enable_dynamic_filter_pushdown); + assert!(!optimizer.enable_topk_dynamic_filter_pushdown); + assert_eq!( + decoder + .session_state + .execution_props() + .config_options + .as_ref() + .unwrap() + .optimizer + .enable_topk_dynamic_filter_pushdown, + false + ); + + let invalid_ctx = Arc::new( + QueryContextBuilder::default() + .set_extension( + ENABLE_DYNAMIC_FILTER_PUSHDOWN.to_string(), + "invalid".to_string(), + ) + .build(), + ); + assert!(DefaultPlanDecoder::new(SessionStateBuilder::new().build(), &invalid_ctx).is_err()); + } + fn mock_plan(schema: SchemaRef) -> LogicalPlan { let table_source = LogicalTableSource::new(schema); let projection = None; diff --git a/src/query/src/sql.rs b/src/query/src/sql.rs index 93e74ca33e..ed364dbdc8 100644 --- a/src/query/src/sql.rs +++ b/src/query/src/sql.rs @@ -71,6 +71,7 @@ use table::requests::{FILE_TABLE_LOCATION_KEY, FILE_TABLE_PATTERN_KEY}; use crate::QueryEngineRef; use crate::error::{self, Result, UnsupportedVariableSnafu}; +use crate::options::dynamic_filter_pushdown_options; use crate::planner::DfLogicalPlanner; const SCHEMAS_COLUMN: &str = "Database"; @@ -845,10 +846,28 @@ pub async fn show_charsets_dataframe( pub fn show_variable(stmt: ShowVariables, query_ctx: QueryContextRef) -> Result { let variable = stmt.variable.to_string().to_uppercase(); + let dynamic_filter_pushdown = dynamic_filter_pushdown_options( + query_ctx + .configuration_parameter() + .dynamic_filter_pushdown(), + &query_ctx.extensions(), + )?; let value = match variable.as_str() { "SYSTEM_TIME_ZONE" | "SYSTEM_TIMEZONE" => get_timezone(None).to_string(), "TIME_ZONE" | "TIMEZONE" => query_ctx.timezone().to_string(), "READ_PREFERENCE" => query_ctx.read_preference().to_string(), + "ENABLE_DYNAMIC_FILTER_PUSHDOWN" => dynamic_filter_pushdown + .enable_dynamic_filter_pushdown + .to_string(), + "ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN" => dynamic_filter_pushdown + .enable_aggregate_dynamic_filter_pushdown + .to_string(), + "ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN" => dynamic_filter_pushdown + .enable_join_dynamic_filter_pushdown + .to_string(), + "ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN" => dynamic_filter_pushdown + .enable_topk_dynamic_filter_pushdown + .to_string(), "DATESTYLE" => { let (style, order) = *query_ctx.configuration_parameter().pg_datetime_style(); format!("{}, {}", style, order) @@ -1668,6 +1687,10 @@ mod test { exec_show_variable("TIMEZONE", "Asia/Shanghai").unwrap(), "Asia/Shanghai" ); + assert_eq!( + exec_show_variable("ENABLE_DYNAMIC_FILTER_PUSHDOWN", "Asia/Shanghai").unwrap(), + "true" + ); assert!(exec_show_variable("TIME ZONE", "Asia/Shanghai").is_err()); assert!(exec_show_variable("SYSTEM TIME ZONE", "Asia/Shanghai").is_err()); } diff --git a/src/servers/src/grpc/greptime_handler.rs b/src/servers/src/grpc/greptime_handler.rs index 98786d30ad..46b9aae595 100644 --- a/src/servers/src/grpc/greptime_handler.rs +++ b/src/servers/src/grpc/greptime_handler.rs @@ -312,7 +312,7 @@ mod tests { use common_error::ext::BoxedError; use common_error::{GREPTIME_DB_HEADER_ERROR_CODE, GREPTIME_DB_HEADER_ERROR_RETRY_HINT}; use common_time::Timezone; - use query::options::FLOW_SCHEDULED_TIME_MILLIS; + use query::options::{FLOW_SCHEDULED_TIME_MILLIS, dynamic_filter_pushdown_options}; use session::hints::{ INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY, REMOTE_QUERY_ID_EXTENSION_KEY, }; @@ -347,6 +347,10 @@ mod tests { FLOW_SCHEDULED_TIME_MILLIS.to_string(), "1700000000000".to_string(), ), + ( + "enable_dynamic_filter_pushdown".to_string(), + "false".to_string(), + ), ], HashMap::from([(7, 88)]), ) @@ -373,6 +377,16 @@ mod tests { query_context.extension(FLOW_SCHEDULED_TIME_MILLIS), Some("1700000000000") ); + assert!( + !dynamic_filter_pushdown_options( + query_context + .configuration_parameter() + .dynamic_filter_pushdown(), + &query_context.extensions(), + ) + .unwrap() + .enable_dynamic_filter_pushdown + ); } #[test] diff --git a/src/servers/src/http/hints.rs b/src/servers/src/http/hints.rs index 09e4879dd5..09dc910c5d 100644 --- a/src/servers/src/http/hints.rs +++ b/src/servers/src/http/hints.rs @@ -46,7 +46,7 @@ fn apply_hints(query_ctx: &mut QueryContext, hints: Vec<(String, String)>) { #[cfg(test)] mod tests { use common_query::request::INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY as COMMON_INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY; - use query::options::FLOW_SCHEDULED_TIME_MILLIS; + use query::options::{FLOW_SCHEDULED_TIME_MILLIS, dynamic_filter_pushdown_options}; use session::context::{QueryContextBuilder, generate_remote_query_id}; use session::hints::{ INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY, REMOTE_QUERY_ID_EXTENSION_KEY, @@ -79,6 +79,10 @@ mod tests { FLOW_SCHEDULED_TIME_MILLIS.to_string(), "1700000000000".to_string(), ), + ( + "enable_dynamic_filter_pushdown".to_string(), + "false".to_string(), + ), ("ttl".to_string(), "7d".to_string()), ], ); @@ -97,6 +101,16 @@ mod tests { Some("1700000000000") ); assert_eq!(query_ctx.extension("ttl"), Some("7d")); + assert!( + !dynamic_filter_pushdown_options( + query_ctx + .configuration_parameter() + .dynamic_filter_pushdown(), + &query_ctx.extensions(), + ) + .unwrap() + .enable_dynamic_filter_pushdown + ); } #[test] diff --git a/src/session/src/context.rs b/src/session/src/context.rs index d383aef477..3bfe410ebb 100644 --- a/src/session/src/context.rs +++ b/src/session/src/context.rs @@ -20,7 +20,7 @@ use std::time::Duration; use api::v1::ExplainOptions; use api::v1::region::RegionRequestHeader; -use arc_swap::ArcSwap; +use arc_swap::{ArcSwap, ArcSwapOption}; use auth::UserInfoRef; pub use common_base::protocol::Channel; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; @@ -47,6 +47,12 @@ pub type ConnInfoRef = Arc; pub const FLIGHT_METRICS_HEARTBEAT_INTERVAL: Duration = Duration::from_secs(1); +pub const ENABLE_DYNAMIC_FILTER_PUSHDOWN: &str = "enable_dynamic_filter_pushdown"; +pub const ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN: &str = + "enable_aggregate_dynamic_filter_pushdown"; +pub const ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN: &str = "enable_join_dynamic_filter_pushdown"; +pub const ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN: &str = "enable_topk_dynamic_filter_pushdown"; + const CURSOR_COUNT_WARNING_LIMIT: usize = 10; pub fn generate_remote_query_id() -> String { @@ -628,12 +634,34 @@ pub fn dialect_for_channel(channel: Channel) -> Arc { } } -#[derive(Default, Debug)] +#[derive(Debug)] pub struct ConfigurationVariables { postgres_bytea_output: ArcSwap, pg_datestyle_format: ArcSwap<(PGDateTimeStyle, PGDateOrder)>, pg_intervalstyle_format: ArcSwap, allow_query_fallback: ArcSwap, + enable_dynamic_filter_pushdown: ArcSwapOption, + enable_aggregate_dynamic_filter_pushdown: ArcSwapOption, + enable_join_dynamic_filter_pushdown: ArcSwapOption, + enable_topk_dynamic_filter_pushdown: ArcSwapOption, +} + +impl Default for ConfigurationVariables { + fn default() -> Self { + Self { + postgres_bytea_output: ArcSwap::from_pointee(PGByteaOutputValue::default()), + pg_datestyle_format: ArcSwap::from_pointee(( + PGDateTimeStyle::default(), + PGDateOrder::default(), + )), + pg_intervalstyle_format: ArcSwap::from_pointee(PGIntervalStyle::default()), + allow_query_fallback: ArcSwap::from_pointee(false), + enable_dynamic_filter_pushdown: ArcSwapOption::empty(), + enable_aggregate_dynamic_filter_pushdown: ArcSwapOption::empty(), + enable_join_dynamic_filter_pushdown: ArcSwapOption::empty(), + enable_topk_dynamic_filter_pushdown: ArcSwapOption::empty(), + } + } } impl Clone for ConfigurationVariables { @@ -643,6 +671,18 @@ impl Clone for ConfigurationVariables { pg_datestyle_format: ArcSwap::new(self.pg_datestyle_format.load().clone()), pg_intervalstyle_format: ArcSwap::new(self.pg_intervalstyle_format.load().clone()), allow_query_fallback: ArcSwap::new(self.allow_query_fallback.load().clone()), + enable_dynamic_filter_pushdown: ArcSwapOption::new( + self.enable_dynamic_filter_pushdown.load_full(), + ), + enable_aggregate_dynamic_filter_pushdown: ArcSwapOption::new( + self.enable_aggregate_dynamic_filter_pushdown.load_full(), + ), + enable_join_dynamic_filter_pushdown: ArcSwapOption::new( + self.enable_join_dynamic_filter_pushdown.load_full(), + ), + enable_topk_dynamic_filter_pushdown: ArcSwapOption::new( + self.enable_topk_dynamic_filter_pushdown.load_full(), + ), } } } @@ -683,6 +723,49 @@ impl ConfigurationVariables { pub fn set_allow_query_fallback(&self, allow: bool) { self.allow_query_fallback.swap(Arc::new(allow)); } + + pub fn dynamic_filter_pushdown(&self) -> DynamicFilterPushdownOptions { + DynamicFilterPushdownOptions { + enable_dynamic_filter_pushdown: self + .enable_dynamic_filter_pushdown + .load_full() + .map(|value| *value), + enable_aggregate_dynamic_filter_pushdown: self + .enable_aggregate_dynamic_filter_pushdown + .load_full() + .map(|value| *value), + enable_join_dynamic_filter_pushdown: self + .enable_join_dynamic_filter_pushdown + .load_full() + .map(|value| *value), + enable_topk_dynamic_filter_pushdown: self + .enable_topk_dynamic_filter_pushdown + .load_full() + .map(|value| *value), + } + } + + pub fn set_dynamic_filter_pushdown(&self, name: &str, value: bool) -> bool { + let option = match name { + ENABLE_DYNAMIC_FILTER_PUSHDOWN => &self.enable_dynamic_filter_pushdown, + ENABLE_AGGREGATE_DYNAMIC_FILTER_PUSHDOWN => { + &self.enable_aggregate_dynamic_filter_pushdown + } + ENABLE_JOIN_DYNAMIC_FILTER_PUSHDOWN => &self.enable_join_dynamic_filter_pushdown, + ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN => &self.enable_topk_dynamic_filter_pushdown, + _ => return false, + }; + option.store(Some(Arc::new(value))); + true + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct DynamicFilterPushdownOptions { + pub enable_dynamic_filter_pushdown: Option, + pub enable_aggregate_dynamic_filter_pushdown: Option, + pub enable_join_dynamic_filter_pushdown: Option, + pub enable_topk_dynamic_filter_pushdown: Option, } #[cfg(test)] @@ -728,6 +811,35 @@ mod test { assert_eq!("test", context.get_db_string()); } + #[test] + fn test_dynamic_filter_pushdown_options_are_isolated_by_session() { + let first = Session::new(None, Channel::Mysql, Default::default(), 1); + let second = Session::new(None, Channel::Mysql, Default::default(), 2); + + assert!( + first + .new_query_context() + .configuration_parameter() + .set_dynamic_filter_pushdown(ENABLE_TOPK_DYNAMIC_FILTER_PUSHDOWN, false) + ); + assert_eq!( + first + .new_query_context() + .configuration_parameter() + .dynamic_filter_pushdown() + .enable_topk_dynamic_filter_pushdown, + Some(false) + ); + assert_eq!( + second + .new_query_context() + .configuration_parameter() + .dynamic_filter_pushdown() + .enable_topk_dynamic_filter_pushdown, + None + ); + } + #[test] fn test_fork_has_independent_mutable_session_data() { let context = QueryContext::with(DEFAULT_CATALOG_NAME, "public");