From a83123c983bf89e062e707e5db3257e3397dd3b5 Mon Sep 17 00:00:00 2001 From: XL Liang Date: Sat, 18 Jul 2026 19:19:57 +0800 Subject: [PATCH 1/9] perf: prewarm projection-aware listing statistics --- .../src/formats/parquet/read.rs | 107 +++++++++++++++++- .../sail-data-source/src/listing/planner.rs | 64 ++++++----- crates/sail-data-source/src/listing/source.rs | 8 +- crates/sail-data-source/src/listing/table.rs | 70 +++++++++++- 4 files changed, 216 insertions(+), 33 deletions(-) diff --git a/crates/sail-data-source/src/formats/parquet/read.rs b/crates/sail-data-source/src/formats/parquet/read.rs index a05c4e6938..75bd0da763 100644 --- a/crates/sail-data-source/src/formats/parquet/read.rs +++ b/crates/sail-data-source/src/formats/parquet/read.rs @@ -9,7 +9,7 @@ use datafusion::datasource::physical_plan::parquet::metadata::{ }; use datafusion_common::config::TableParquetOptions; use datafusion_common::parsers::CompressionTypeVariant; -use datafusion_common::{DataFusionError, Result}; +use datafusion_common::{DataFusionError, Result, Statistics}; use datafusion_datasource::file_scan_config::{FileScanConfig, FileScanConfigBuilder}; use futures::{StreamExt, TryStreamExt}; use object_store::{ObjectMeta, ObjectStore}; @@ -115,6 +115,7 @@ impl ReadFormat for ParquetReadFormat { store: &Arc, object: &ObjectMeta, file_schema: SchemaRef, + statistics_columns: Option<&[usize]>, _compression: CompressionTypeVariant, ) -> Result { let options = self.options.clone().into_table_options(); @@ -124,8 +125,17 @@ impl ReadFormat for ParquetReadFormat { .with_file_metadata_cache(Some(metadata_cache)) .fetch_metadata() .await?; - let statistics = - DFParquetMetadata::statistics_from_parquet_metadata(&metadata, &file_schema)?; + let statistics = match statistics_columns { + None => DFParquetMetadata::statistics_from_parquet_metadata(&metadata, &file_schema)?, + Some(columns) => { + let statistics_schema = Arc::new(file_schema.project(columns)?); + let projected_statistics = DFParquetMetadata::statistics_from_parquet_metadata( + &metadata, + &statistics_schema, + )?; + expand_projected_statistics(&file_schema, columns, projected_statistics) + } + }; let ordering = ordering_from_parquet_metadata(&metadata, &file_schema)?; Ok(ListingFileMeta { statistics, @@ -168,6 +178,24 @@ impl ReadFormat for ParquetReadFormat { } } +fn expand_projected_statistics( + file_schema: &Schema, + columns: &[usize], + projected_statistics: Statistics, +) -> Statistics { + let mut statistics = Statistics::new_unknown(file_schema); + statistics.num_rows = projected_statistics.num_rows; + statistics.total_byte_size = projected_statistics.total_byte_size; + for (index, column_statistics) in columns + .iter() + .copied() + .zip(projected_statistics.column_statistics) + { + statistics.column_statistics[index] = column_statistics; + } + statistics +} + /// Clears all metadata (Schema level and field level) for a schema. fn clear_metadata(schema: Schema) -> Schema { let fields = schema @@ -194,3 +222,76 @@ fn parse_coerce_int96_string(setting: &str) -> Result { ))), } } + +#[cfg(test)] +mod tests { + use datafusion_common::ScalarValue; + use datafusion_common::stats::Precision; + + use super::*; + + #[test] + fn projected_statistics_keep_full_schema_positions() { + let file_schema = Schema::new(vec![ + Arc::new(datafusion::arrow::datatypes::Field::new( + "a", + datafusion::arrow::datatypes::DataType::Int32, + false, + )), + Arc::new(datafusion::arrow::datatypes::Field::new( + "b", + datafusion::arrow::datatypes::DataType::Int32, + false, + )), + Arc::new(datafusion::arrow::datatypes::Field::new( + "c", + datafusion::arrow::datatypes::DataType::Int32, + false, + )), + ]); + let mut selected_column = datafusion_common::ColumnStatistics::new_unknown(); + selected_column.min_value = Precision::Exact(ScalarValue::Int32(Some(10))); + selected_column.max_value = Precision::Exact(ScalarValue::Int32(Some(20))); + let projected_statistics = Statistics { + num_rows: Precision::Exact(100), + total_byte_size: Precision::Absent, + column_statistics: vec![selected_column.clone()], + }; + + let statistics = expand_projected_statistics(&file_schema, &[1], projected_statistics); + + assert_eq!(statistics.num_rows, Precision::Exact(100)); + assert_eq!(statistics.column_statistics.len(), 3); + assert_eq!( + statistics.column_statistics[0], + datafusion_common::ColumnStatistics::new_unknown() + ); + assert_eq!(statistics.column_statistics[1], selected_column); + assert_eq!( + statistics.column_statistics[2], + datafusion_common::ColumnStatistics::new_unknown() + ); + } + + #[test] + fn row_count_statistics_do_not_require_columns() { + let file_schema = Schema::new(vec![Arc::new(datafusion::arrow::datatypes::Field::new( + "a", + datafusion::arrow::datatypes::DataType::Int32, + false, + ))]); + let projected_statistics = Statistics { + num_rows: Precision::Exact(100), + total_byte_size: Precision::Absent, + column_statistics: vec![], + }; + + let statistics = expand_projected_statistics(&file_schema, &[], projected_statistics); + + assert_eq!(statistics.num_rows, Precision::Exact(100)); + assert_eq!( + statistics.column_statistics, + vec![datafusion_common::ColumnStatistics::new_unknown()] + ); + } +} diff --git a/crates/sail-data-source/src/listing/planner.rs b/crates/sail-data-source/src/listing/planner.rs index 03f6603802..36f9044301 100644 --- a/crates/sail-data-source/src/listing/planner.rs +++ b/crates/sail-data-source/src/listing/planner.rs @@ -52,6 +52,16 @@ struct ListFilesResult { #[derive(Debug, Default)] pub struct ListingPhysicalPlanner; +pub(crate) async fn prewarm_file_statistics( + source: &ListingTableSource, + ctx: &dyn Session, +) -> datafusion_common::Result<()> { + if source.config().collect_stat { + let _ = list_files_for_scan(source, ctx, &[], None, Some(&[])).await?; + } + Ok(()) +} + #[async_trait] impl ExtensionPlanner for ListingPhysicalPlanner { async fn plan_extension( @@ -104,6 +114,17 @@ impl ExtensionPlanner for ListingPhysicalPlanner { }); let statistic_file_limit = if filters.is_empty() { limit } else { None }; + let file_column_count = source.config().schema.file_schema().fields().len(); + let statistics_columns = projection.as_ref().map(|projection| { + let mut columns = projection + .iter() + .copied() + .filter(|index| *index < file_column_count) + .collect::>(); + columns.sort_unstable(); + columns.dedup(); + columns + }); let ListFilesResult { mut file_groups, @@ -114,6 +135,7 @@ impl ExtensionPlanner for ListingPhysicalPlanner { session_state, &partition_filters, statistic_file_limit, + statistics_columns.as_deref(), ) .await?; @@ -331,6 +353,7 @@ async fn list_files_for_scan<'a>( ctx: &'a dyn Session, filters: &'a [Expr], limit: Option, + statistics_columns: Option<&'a [usize]>, ) -> datafusion_common::Result { let store = if let Some(url) = source.config().table_paths.first() { ctx.runtime_env().object_store(url)? @@ -364,12 +387,21 @@ async fn list_files_for_scan<'a>( let meta_fetch_concurrency = ctx.config_options().execution.meta_fetch_concurrency; let file_list = stream::iter(file_list).flatten_unordered(meta_fetch_concurrency); + let statistics_cache_table = source.statistics_cache_table(statistics_columns); let files = file_list .map(|part_file| async { let part_file = part_file?; let (statistics, ordering) = if source.config().collect_stat { - do_collect_statistics_and_ordering(source, ctx, &store, &part_file).await? + do_collect_statistics_and_ordering( + source, + ctx, + &store, + &part_file, + statistics_columns, + &statistics_cache_table, + ) + .await? } else { ( Arc::new(Statistics::new_unknown( @@ -427,15 +459,14 @@ async fn do_collect_statistics_and_ordering( ctx: &dyn Session, store: &Arc, part_file: &datafusion_datasource::PartitionedFile, + statistics_columns: Option<&[usize]>, + statistics_cache_table: &datafusion_common::TableReference, ) -> datafusion_common::Result<(Arc, Option)> { let meta = &part_file.object_meta; let file_schema = source.config().schema.file_schema(); let file_statistic_cache = ctx.runtime_env().cache_manager.get_file_statistic_cache(); let cache_key = TableScopedPath { - table: Some(statistics_cache_table_ref( - part_file.table_reference.as_ref(), - file_schema.as_ref(), - )), + table: Some(statistics_cache_table.clone()), path: meta.location.clone(), }; @@ -459,6 +490,7 @@ async fn do_collect_statistics_and_ordering( store, meta, source.config().schema.file_schema().clone(), + statistics_columns, source.config().compression, ) .await?; @@ -478,28 +510,6 @@ async fn do_collect_statistics_and_ordering( Ok((statistics, file_meta.ordering)) } -fn statistics_cache_table_ref( - table: Option<&datafusion_common::TableReference>, - schema: &arrow::datatypes::Schema, -) -> datafusion_common::TableReference { - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - - for field in schema.fields() { - std::hash::Hash::hash(field.name(), &mut hasher); - std::hash::Hash::hash(&field.data_type().to_string(), &mut hasher); - std::hash::Hash::hash(&field.is_nullable(), &mut hasher); - } - - let table = table - .map(ToString::to_string) - .unwrap_or_else(|| "listing".to_string()); - - datafusion_common::TableReference::bare(format!( - "{table}_{:016x}", - std::hash::Hasher::finish(&hasher) - )) -} - async fn get_files_with_limit( files: impl Stream>, limit: Option, diff --git a/crates/sail-data-source/src/listing/source.rs b/crates/sail-data-source/src/listing/source.rs index 423eb0023d..5a2de9e721 100644 --- a/crates/sail-data-source/src/listing/source.rs +++ b/crates/sail-data-source/src/listing/source.rs @@ -24,6 +24,7 @@ use sail_common_datafusion::datasource::{ }; use url::Url; +use crate::listing::planner::prewarm_file_statistics; use crate::listing::table::{ListingTableSource, ListingTableSourceConfig}; use crate::listing::utils::{ infer_partitions, rewrite_utf8view_fields, sample_listing_files, validate_partitions, @@ -66,15 +67,17 @@ pub trait ReadFormat: Debug + Send + Sync + 'static { /// Infer file-level metadata needed for planning. /// The metadata includes statistics and ordering. + /// `statistics_columns` is `None` for all file columns and `Some` for a selected subset. async fn infer_file_meta( &self, ctx: &dyn Session, store: &Arc, object: &ObjectMeta, file_schema: SchemaRef, + statistics_columns: Option<&[usize]>, compression: CompressionTypeVariant, ) -> Result { - let _ = (ctx, store, object, compression); + let _ = (ctx, store, object, statistics_columns, compression); Ok(ListingFileMeta { statistics: Statistics::new_unknown(&file_schema), ordering: None, @@ -239,6 +242,9 @@ impl TableFormat for ListingTableFormat { read_format: Arc::new(read_format), compression, })?; + if let Err(error) = prewarm_file_statistics(&source, ctx).await { + log::warn!("failed to prewarm listing file statistics: {error}"); + } Ok(Arc::new(source)) } diff --git a/crates/sail-data-source/src/listing/table.rs b/crates/sail-data-source/src/listing/table.rs index 2ef187aec2..1d261ea31b 100644 --- a/crates/sail-data-source/src/listing/table.rs +++ b/crates/sail-data-source/src/listing/table.rs @@ -1,13 +1,14 @@ // The listing table source is adapted from the DataFusion `ListingTable` implementation. // [CREDIT]: https://github.com/apache/datafusion/blob/53.1.0/datafusion/catalog-listing/src/table.rs +use std::hash::{Hash, Hasher}; use std::sync::Arc; use datafusion::arrow::datatypes::SchemaRef; use datafusion::logical_expr::expr::Sort; use datafusion::logical_expr::{Expr, TableProviderFilterPushDown, TableSource, TableType}; use datafusion_common::parsers::CompressionTypeVariant; -use datafusion_common::{Constraints, Result}; +use datafusion_common::{Constraints, Result, TableReference}; use datafusion_datasource::{ListingTableUrl, TableSchema}; use crate::listing::source::ReadFormat; @@ -32,16 +33,81 @@ pub struct ListingTableSourceConfig { #[derive(Clone, Debug)] pub struct ListingTableSource { config: ListingTableSourceConfig, + statistics_cache_namespace: String, } impl ListingTableSource { pub fn try_new(config: ListingTableSourceConfig) -> Result { - Ok(Self { config }) + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + for path in &config.table_paths { + path.to_string().hash(&mut hasher); + } + for field in config.schema.file_schema().fields() { + field.name().hash(&mut hasher); + field.data_type().to_string().hash(&mut hasher); + field.is_nullable().hash(&mut hasher); + } + format!("{:?}", config.read_format).hash(&mut hasher); + format!("{:?}", config.compression).hash(&mut hasher); + let statistics_cache_namespace = format!("listing_{:016x}", hasher.finish()); + + Ok(Self { + config, + statistics_cache_namespace, + }) } pub fn config(&self) -> &ListingTableSourceConfig { &self.config } + + pub(crate) fn statistics_cache_table( + &self, + statistics_columns: Option<&[usize]>, + ) -> TableReference { + let coverage = statistics_cache_coverage(statistics_columns); + TableReference::bare(format!("{}_{}", self.statistics_cache_namespace, coverage)) + } +} + +fn statistics_cache_coverage(statistics_columns: Option<&[usize]>) -> String { + match statistics_columns { + None => "all".to_string(), + Some([]) => "rows".to_string(), + Some(columns) => { + let mut columns = columns.to_vec(); + columns.sort_unstable(); + columns.dedup(); + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + columns.hash(&mut hasher); + format!("columns_{:016x}", hasher.finish()) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn statistics_cache_coverage_is_canonical_and_isolated() { + assert_ne!( + statistics_cache_coverage(None), + statistics_cache_coverage(Some(&[])) + ); + assert_ne!( + statistics_cache_coverage(Some(&[])), + statistics_cache_coverage(Some(&[1])) + ); + assert_ne!( + statistics_cache_coverage(Some(&[1])), + statistics_cache_coverage(Some(&[2])) + ); + assert_eq!( + statistics_cache_coverage(Some(&[2, 1])), + statistics_cache_coverage(Some(&[1, 2, 2])) + ); + } } impl TableSource for ListingTableSource { From 9cc72c9054bf8a13e6c59c9f31af00f462d89967 Mon Sep 17 00:00:00 2001 From: XL Liang Date: Sat, 18 Jul 2026 19:20:10 +0800 Subject: [PATCH 2/9] perf: skip unused plan descriptions --- crates/sail-flight/src/service.rs | 2 +- crates/sail-plan/src/lib.rs | 48 +++++++++++++++---- .../sail-spark-connect/src/proto/function.rs | 2 +- .../src/service/plan_executor.rs | 10 ++-- 4 files changed, 47 insertions(+), 15 deletions(-) diff --git a/crates/sail-flight/src/service.rs b/crates/sail-flight/src/service.rs index dc806189ad..89f2c6b0fd 100644 --- a/crates/sail-flight/src/service.rs +++ b/crates/sail-flight/src/service.rs @@ -104,7 +104,7 @@ impl FlightSqlService for SailFlightSqlService { }; let ctx = self.get_session_context().await?; - let (plan, _) = resolve_and_execute_plan(&ctx, self.config.clone(), plan) + let plan = resolve_and_execute_plan(&ctx, self.config.clone(), plan) .await .map_err(|e| Status::internal(format!("plan error: {e}")))?; let schema = plan.schema(); diff --git a/crates/sail-plan/src/lib.rs b/crates/sail-plan/src/lib.rs index bb87e6d393..db753da6a4 100644 --- a/crates/sail-plan/src/lib.rs +++ b/crates/sail-plan/src/lib.rs @@ -35,11 +35,39 @@ pub async fn resolve_and_execute_plan( ctx: &SessionContext, config: Arc, plan: spec::Plan, +) -> PlanResult> { + let (plan, _) = resolve_execution_plan(ctx, config, plan, PlanDescriptionMode::Omit).await?; + Ok(plan) +} + +pub async fn resolve_and_describe_execution_plan( + ctx: &SessionContext, + config: Arc, + plan: spec::Plan, +) -> PlanResult<(Arc, Vec)> { + resolve_execution_plan(ctx, config, plan, PlanDescriptionMode::Collect).await +} + +enum PlanDescriptionMode { + Omit, + Collect, +} + +async fn resolve_execution_plan( + ctx: &SessionContext, + config: Arc, + plan: spec::Plan, + description_mode: PlanDescriptionMode, ) -> PlanResult<(Arc, Vec)> { - let mut info = vec![]; + let mut descriptions = match description_mode { + PlanDescriptionMode::Omit => None, + PlanDescriptionMode::Collect => Some(vec![]), + }; let resolver = PlanResolver::new(ctx, config); let NamedPlan { plan, fields } = resolver.resolve_named_plan(plan).await?; - info.push(plan.to_stringified(PlanType::InitialLogicalPlan)); + if let Some(descriptions) = &mut descriptions { + descriptions.push(plan.to_stringified(PlanType::InitialLogicalPlan)); + } let df = execute_logical_plan(ctx, plan).await?; let (session_state, plan) = df.into_parts(); let plan = session_state.optimize(&plan)?; @@ -48,7 +76,9 @@ pub async fn resolve_and_execute_plan( } else { plan }; - info.push(plan.to_stringified(PlanType::FinalLogicalPlan)); + if let Some(descriptions) = &mut descriptions { + descriptions.push(plan.to_stringified(PlanType::FinalLogicalPlan)); + } let plan = session_state .query_planner() .create_physical_plan(&plan, &session_state) @@ -58,9 +88,11 @@ pub async fn resolve_and_execute_plan( } else { plan }; - info.push(StringifiedPlan::new( - PlanType::FinalPhysicalPlan, - displayable(plan.as_ref()).indent(true).to_string(), - )); - Ok((plan, info)) + if let Some(descriptions) = &mut descriptions { + descriptions.push(StringifiedPlan::new( + PlanType::FinalPhysicalPlan, + displayable(plan.as_ref()).indent(true).to_string(), + )); + } + Ok((plan, descriptions.unwrap_or_default())) } diff --git a/crates/sail-spark-connect/src/proto/function.rs b/crates/sail-spark-connect/src/proto/function.rs index 0eef1bff7a..0199c21d4a 100644 --- a/crates/sail-spark-connect/src/proto/function.rs +++ b/crates/sail-spark-connect/src/proto/function.rs @@ -87,7 +87,7 @@ mod tests { let result = handle.primary().block_on(async { let spark = context.extension::()?; let service = context.extension::()?; - let (plan, _) = + let plan = resolve_and_execute_plan(&context, spark.plan_config()?, plan).await?; let stream = service.runner().execute(&context, plan).await?; read_stream(stream).await diff --git a/crates/sail-spark-connect/src/service/plan_executor.rs b/crates/sail-spark-connect/src/service/plan_executor.rs index d372544a63..387da56bfe 100644 --- a/crates/sail-spark-connect/src/service/plan_executor.rs +++ b/crates/sail-spark-connect/src/service/plan_executor.rs @@ -11,7 +11,7 @@ use log::{debug, warn}; use sail_common::spec; use sail_common_datafusion::extension::SessionExtensionAccessor; use sail_common_datafusion::session::job::JobService; -use sail_plan::resolve_and_execute_plan; +use sail_plan::{resolve_and_describe_execution_plan, resolve_and_execute_plan}; use tonic::Status; use tonic::codegen::tokio_stream::Stream; use tonic::codegen::tokio_stream::wrappers::ReceiverStream; @@ -64,7 +64,7 @@ impl Stream for ExecutePlanResponseStream { let mut response = ExecutePlanResponse::default(); response.session_id.clone_from(&self.session_id); response.server_side_session_id.clone_from(&self.session_id); - response.operation_id.clone_from(&self.operation_id.clone()); + response.operation_id.clone_from(&self.operation_id); response.response_id = item.id; match item.batch { ExecutorBatch::ArrowBatch(batch) => { @@ -126,7 +126,7 @@ async fn handle_execute_plan( let spark = ctx.extension::()?; let service = ctx.extension::()?; let operation_id = metadata.operation_id.clone(); - let (plan, _) = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; + let plan = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; let stream = { let span = Span::enter_with_parent("JobRunner::execute", &span); service.runner().execute(ctx, plan).in_span(span).await? @@ -247,7 +247,7 @@ pub(crate) async fn handle_execute_sql_command( let relation = match plan { spec::Plan::Query(_) => relation, command @ spec::Plan::Command(_) => { - let (plan, _) = resolve_and_execute_plan(ctx, spark.plan_config()?, command).await?; + let plan = resolve_and_execute_plan(ctx, spark.plan_config()?, command).await?; let stream = service.runner().execute(ctx, plan).await?; let schema = stream.schema(); let data = read_stream(stream).await?; @@ -286,7 +286,7 @@ pub(crate) async fn handle_execute_write_stream_operation_start( let reattachable = metadata.reattachable; let query_name = start.query_name.clone(); let plan = spec::Plan::Command(spec::CommandPlan::new(start.try_into()?)); - let (plan, info) = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; + let (plan, info) = resolve_and_describe_execution_plan(ctx, spark.plan_config()?, plan).await?; let stream = service.runner().execute(ctx, plan).await?; let id = spark.start_streaming_query(query_name.clone(), info, stream)?; let result = WriteStreamOperationStartResult { From a003933b84eb00dad70181ffc947ae6904cb6a6e Mon Sep 17 00:00:00 2001 From: XL Liang Date: Sat, 18 Jul 2026 19:20:16 +0800 Subject: [PATCH 3/9] perf: cache Spark plan configuration --- crates/sail-spark-connect/src/session.rs | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/crates/sail-spark-connect/src/session.rs b/crates/sail-spark-connect/src/session.rs index ef70593f6f..4db0041841 100644 --- a/crates/sail-spark-connect/src/session.rs +++ b/crates/sail-spark-connect/src/session.rs @@ -81,10 +81,15 @@ impl SparkSession { } pub(crate) fn plan_config(&self) -> SparkResult> { - let state = self.state.lock()?; + let mut state = self.state.lock()?; + if let Some(config) = &state.plan_config { + return Ok(Arc::clone(config)); + } let mut config = PlanConfig::try_from(&state.config)?; config.session_user_id = self.user_id().to_string(); - Ok(Arc::new(config)) + let config = Arc::new(config); + state.plan_config = Some(Arc::clone(&config)); + Ok(config) } pub(crate) fn get_config(&self, keys: Vec) -> SparkResult> { @@ -129,6 +134,7 @@ impl SparkSession { pub(crate) fn set_config(&self, kv: Vec) -> SparkResult<()> { let mut state = self.state.lock()?; + state.plan_config = None; for ConfigKeyValue { key, value } in kv { if let Some(value) = value { state.config.set(key, value)?; @@ -143,6 +149,7 @@ impl SparkSession { pub(crate) fn unset_config(&self, keys: Vec) -> SparkResult<()> { let mut state = self.state.lock()?; + state.plan_config = None; for key in keys { state.config.unset(&key)? } @@ -292,6 +299,7 @@ impl SparkSession { struct SparkSessionState { config: SparkRuntimeConfig, + plan_config: Option>, executors: HashMap>, streaming_queries: StreamingQueryManager, } @@ -300,6 +308,7 @@ impl SparkSessionState { fn new() -> Self { Self { config: SparkRuntimeConfig::new(), + plan_config: None, executors: HashMap::new(), streaming_queries: StreamingQueryManager::new(), } From 2c9e0d07fb756dc8fef842c3fd1bc09451e5a5f0 Mon Sep 17 00:00:00 2001 From: XL Liang Date: Thu, 23 Jul 2026 22:58:50 +0800 Subject: [PATCH 4/9] Revert "perf: cache Spark plan configuration" This reverts commit a003933b84eb00dad70181ffc947ae6904cb6a6e. --- crates/sail-spark-connect/src/session.rs | 13 ++----------- 1 file changed, 2 insertions(+), 11 deletions(-) diff --git a/crates/sail-spark-connect/src/session.rs b/crates/sail-spark-connect/src/session.rs index 4db0041841..ef70593f6f 100644 --- a/crates/sail-spark-connect/src/session.rs +++ b/crates/sail-spark-connect/src/session.rs @@ -81,15 +81,10 @@ impl SparkSession { } pub(crate) fn plan_config(&self) -> SparkResult> { - let mut state = self.state.lock()?; - if let Some(config) = &state.plan_config { - return Ok(Arc::clone(config)); - } + let state = self.state.lock()?; let mut config = PlanConfig::try_from(&state.config)?; config.session_user_id = self.user_id().to_string(); - let config = Arc::new(config); - state.plan_config = Some(Arc::clone(&config)); - Ok(config) + Ok(Arc::new(config)) } pub(crate) fn get_config(&self, keys: Vec) -> SparkResult> { @@ -134,7 +129,6 @@ impl SparkSession { pub(crate) fn set_config(&self, kv: Vec) -> SparkResult<()> { let mut state = self.state.lock()?; - state.plan_config = None; for ConfigKeyValue { key, value } in kv { if let Some(value) = value { state.config.set(key, value)?; @@ -149,7 +143,6 @@ impl SparkSession { pub(crate) fn unset_config(&self, keys: Vec) -> SparkResult<()> { let mut state = self.state.lock()?; - state.plan_config = None; for key in keys { state.config.unset(&key)? } @@ -299,7 +292,6 @@ impl SparkSession { struct SparkSessionState { config: SparkRuntimeConfig, - plan_config: Option>, executors: HashMap>, streaming_queries: StreamingQueryManager, } @@ -308,7 +300,6 @@ impl SparkSessionState { fn new() -> Self { Self { config: SparkRuntimeConfig::new(), - plan_config: None, executors: HashMap::new(), streaming_queries: StreamingQueryManager::new(), } From 1e35725521b2af74efb22e4ebda1680d051bf40b Mon Sep 17 00:00:00 2001 From: XL Liang Date: Thu, 23 Jul 2026 22:58:50 +0800 Subject: [PATCH 5/9] Revert "perf: skip unused plan descriptions" This reverts commit 9cc72c9054bf8a13e6c59c9f31af00f462d89967. --- crates/sail-flight/src/service.rs | 2 +- crates/sail-plan/src/lib.rs | 48 ++++--------------- .../sail-spark-connect/src/proto/function.rs | 2 +- .../src/service/plan_executor.rs | 10 ++-- 4 files changed, 15 insertions(+), 47 deletions(-) diff --git a/crates/sail-flight/src/service.rs b/crates/sail-flight/src/service.rs index 89f2c6b0fd..dc806189ad 100644 --- a/crates/sail-flight/src/service.rs +++ b/crates/sail-flight/src/service.rs @@ -104,7 +104,7 @@ impl FlightSqlService for SailFlightSqlService { }; let ctx = self.get_session_context().await?; - let plan = resolve_and_execute_plan(&ctx, self.config.clone(), plan) + let (plan, _) = resolve_and_execute_plan(&ctx, self.config.clone(), plan) .await .map_err(|e| Status::internal(format!("plan error: {e}")))?; let schema = plan.schema(); diff --git a/crates/sail-plan/src/lib.rs b/crates/sail-plan/src/lib.rs index db753da6a4..bb87e6d393 100644 --- a/crates/sail-plan/src/lib.rs +++ b/crates/sail-plan/src/lib.rs @@ -35,39 +35,11 @@ pub async fn resolve_and_execute_plan( ctx: &SessionContext, config: Arc, plan: spec::Plan, -) -> PlanResult> { - let (plan, _) = resolve_execution_plan(ctx, config, plan, PlanDescriptionMode::Omit).await?; - Ok(plan) -} - -pub async fn resolve_and_describe_execution_plan( - ctx: &SessionContext, - config: Arc, - plan: spec::Plan, -) -> PlanResult<(Arc, Vec)> { - resolve_execution_plan(ctx, config, plan, PlanDescriptionMode::Collect).await -} - -enum PlanDescriptionMode { - Omit, - Collect, -} - -async fn resolve_execution_plan( - ctx: &SessionContext, - config: Arc, - plan: spec::Plan, - description_mode: PlanDescriptionMode, ) -> PlanResult<(Arc, Vec)> { - let mut descriptions = match description_mode { - PlanDescriptionMode::Omit => None, - PlanDescriptionMode::Collect => Some(vec![]), - }; + let mut info = vec![]; let resolver = PlanResolver::new(ctx, config); let NamedPlan { plan, fields } = resolver.resolve_named_plan(plan).await?; - if let Some(descriptions) = &mut descriptions { - descriptions.push(plan.to_stringified(PlanType::InitialLogicalPlan)); - } + info.push(plan.to_stringified(PlanType::InitialLogicalPlan)); let df = execute_logical_plan(ctx, plan).await?; let (session_state, plan) = df.into_parts(); let plan = session_state.optimize(&plan)?; @@ -76,9 +48,7 @@ async fn resolve_execution_plan( } else { plan }; - if let Some(descriptions) = &mut descriptions { - descriptions.push(plan.to_stringified(PlanType::FinalLogicalPlan)); - } + info.push(plan.to_stringified(PlanType::FinalLogicalPlan)); let plan = session_state .query_planner() .create_physical_plan(&plan, &session_state) @@ -88,11 +58,9 @@ async fn resolve_execution_plan( } else { plan }; - if let Some(descriptions) = &mut descriptions { - descriptions.push(StringifiedPlan::new( - PlanType::FinalPhysicalPlan, - displayable(plan.as_ref()).indent(true).to_string(), - )); - } - Ok((plan, descriptions.unwrap_or_default())) + info.push(StringifiedPlan::new( + PlanType::FinalPhysicalPlan, + displayable(plan.as_ref()).indent(true).to_string(), + )); + Ok((plan, info)) } diff --git a/crates/sail-spark-connect/src/proto/function.rs b/crates/sail-spark-connect/src/proto/function.rs index 0199c21d4a..0eef1bff7a 100644 --- a/crates/sail-spark-connect/src/proto/function.rs +++ b/crates/sail-spark-connect/src/proto/function.rs @@ -87,7 +87,7 @@ mod tests { let result = handle.primary().block_on(async { let spark = context.extension::()?; let service = context.extension::()?; - let plan = + let (plan, _) = resolve_and_execute_plan(&context, spark.plan_config()?, plan).await?; let stream = service.runner().execute(&context, plan).await?; read_stream(stream).await diff --git a/crates/sail-spark-connect/src/service/plan_executor.rs b/crates/sail-spark-connect/src/service/plan_executor.rs index 387da56bfe..d372544a63 100644 --- a/crates/sail-spark-connect/src/service/plan_executor.rs +++ b/crates/sail-spark-connect/src/service/plan_executor.rs @@ -11,7 +11,7 @@ use log::{debug, warn}; use sail_common::spec; use sail_common_datafusion::extension::SessionExtensionAccessor; use sail_common_datafusion::session::job::JobService; -use sail_plan::{resolve_and_describe_execution_plan, resolve_and_execute_plan}; +use sail_plan::resolve_and_execute_plan; use tonic::Status; use tonic::codegen::tokio_stream::Stream; use tonic::codegen::tokio_stream::wrappers::ReceiverStream; @@ -64,7 +64,7 @@ impl Stream for ExecutePlanResponseStream { let mut response = ExecutePlanResponse::default(); response.session_id.clone_from(&self.session_id); response.server_side_session_id.clone_from(&self.session_id); - response.operation_id.clone_from(&self.operation_id); + response.operation_id.clone_from(&self.operation_id.clone()); response.response_id = item.id; match item.batch { ExecutorBatch::ArrowBatch(batch) => { @@ -126,7 +126,7 @@ async fn handle_execute_plan( let spark = ctx.extension::()?; let service = ctx.extension::()?; let operation_id = metadata.operation_id.clone(); - let plan = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; + let (plan, _) = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; let stream = { let span = Span::enter_with_parent("JobRunner::execute", &span); service.runner().execute(ctx, plan).in_span(span).await? @@ -247,7 +247,7 @@ pub(crate) async fn handle_execute_sql_command( let relation = match plan { spec::Plan::Query(_) => relation, command @ spec::Plan::Command(_) => { - let plan = resolve_and_execute_plan(ctx, spark.plan_config()?, command).await?; + let (plan, _) = resolve_and_execute_plan(ctx, spark.plan_config()?, command).await?; let stream = service.runner().execute(ctx, plan).await?; let schema = stream.schema(); let data = read_stream(stream).await?; @@ -286,7 +286,7 @@ pub(crate) async fn handle_execute_write_stream_operation_start( let reattachable = metadata.reattachable; let query_name = start.query_name.clone(); let plan = spec::Plan::Command(spec::CommandPlan::new(start.try_into()?)); - let (plan, info) = resolve_and_describe_execution_plan(ctx, spark.plan_config()?, plan).await?; + let (plan, info) = resolve_and_execute_plan(ctx, spark.plan_config()?, plan).await?; let stream = service.runner().execute(ctx, plan).await?; let id = spark.start_streaming_query(query_name.clone(), info, stream)?; let result = WriteStreamOperationStartResult { From 0f2998d77bcc32a5a215114dd4cabddd1c637b1f Mon Sep 17 00:00:00 2001 From: XL Liang Date: Thu, 23 Jul 2026 23:10:53 +0800 Subject: [PATCH 6/9] perf: reuse fully prewarmed listing statistics --- .../src/formats/parquet/read.rs | 107 +---------- .../sail-data-source/src/listing/planner.rs | 181 +++++++++++++++--- crates/sail-data-source/src/listing/source.rs | 4 +- crates/sail-data-source/src/listing/table.rs | 70 +------ 4 files changed, 158 insertions(+), 204 deletions(-) diff --git a/crates/sail-data-source/src/formats/parquet/read.rs b/crates/sail-data-source/src/formats/parquet/read.rs index 75bd0da763..a05c4e6938 100644 --- a/crates/sail-data-source/src/formats/parquet/read.rs +++ b/crates/sail-data-source/src/formats/parquet/read.rs @@ -9,7 +9,7 @@ use datafusion::datasource::physical_plan::parquet::metadata::{ }; use datafusion_common::config::TableParquetOptions; use datafusion_common::parsers::CompressionTypeVariant; -use datafusion_common::{DataFusionError, Result, Statistics}; +use datafusion_common::{DataFusionError, Result}; use datafusion_datasource::file_scan_config::{FileScanConfig, FileScanConfigBuilder}; use futures::{StreamExt, TryStreamExt}; use object_store::{ObjectMeta, ObjectStore}; @@ -115,7 +115,6 @@ impl ReadFormat for ParquetReadFormat { store: &Arc, object: &ObjectMeta, file_schema: SchemaRef, - statistics_columns: Option<&[usize]>, _compression: CompressionTypeVariant, ) -> Result { let options = self.options.clone().into_table_options(); @@ -125,17 +124,8 @@ impl ReadFormat for ParquetReadFormat { .with_file_metadata_cache(Some(metadata_cache)) .fetch_metadata() .await?; - let statistics = match statistics_columns { - None => DFParquetMetadata::statistics_from_parquet_metadata(&metadata, &file_schema)?, - Some(columns) => { - let statistics_schema = Arc::new(file_schema.project(columns)?); - let projected_statistics = DFParquetMetadata::statistics_from_parquet_metadata( - &metadata, - &statistics_schema, - )?; - expand_projected_statistics(&file_schema, columns, projected_statistics) - } - }; + let statistics = + DFParquetMetadata::statistics_from_parquet_metadata(&metadata, &file_schema)?; let ordering = ordering_from_parquet_metadata(&metadata, &file_schema)?; Ok(ListingFileMeta { statistics, @@ -178,24 +168,6 @@ impl ReadFormat for ParquetReadFormat { } } -fn expand_projected_statistics( - file_schema: &Schema, - columns: &[usize], - projected_statistics: Statistics, -) -> Statistics { - let mut statistics = Statistics::new_unknown(file_schema); - statistics.num_rows = projected_statistics.num_rows; - statistics.total_byte_size = projected_statistics.total_byte_size; - for (index, column_statistics) in columns - .iter() - .copied() - .zip(projected_statistics.column_statistics) - { - statistics.column_statistics[index] = column_statistics; - } - statistics -} - /// Clears all metadata (Schema level and field level) for a schema. fn clear_metadata(schema: Schema) -> Schema { let fields = schema @@ -222,76 +194,3 @@ fn parse_coerce_int96_string(setting: &str) -> Result { ))), } } - -#[cfg(test)] -mod tests { - use datafusion_common::ScalarValue; - use datafusion_common::stats::Precision; - - use super::*; - - #[test] - fn projected_statistics_keep_full_schema_positions() { - let file_schema = Schema::new(vec![ - Arc::new(datafusion::arrow::datatypes::Field::new( - "a", - datafusion::arrow::datatypes::DataType::Int32, - false, - )), - Arc::new(datafusion::arrow::datatypes::Field::new( - "b", - datafusion::arrow::datatypes::DataType::Int32, - false, - )), - Arc::new(datafusion::arrow::datatypes::Field::new( - "c", - datafusion::arrow::datatypes::DataType::Int32, - false, - )), - ]); - let mut selected_column = datafusion_common::ColumnStatistics::new_unknown(); - selected_column.min_value = Precision::Exact(ScalarValue::Int32(Some(10))); - selected_column.max_value = Precision::Exact(ScalarValue::Int32(Some(20))); - let projected_statistics = Statistics { - num_rows: Precision::Exact(100), - total_byte_size: Precision::Absent, - column_statistics: vec![selected_column.clone()], - }; - - let statistics = expand_projected_statistics(&file_schema, &[1], projected_statistics); - - assert_eq!(statistics.num_rows, Precision::Exact(100)); - assert_eq!(statistics.column_statistics.len(), 3); - assert_eq!( - statistics.column_statistics[0], - datafusion_common::ColumnStatistics::new_unknown() - ); - assert_eq!(statistics.column_statistics[1], selected_column); - assert_eq!( - statistics.column_statistics[2], - datafusion_common::ColumnStatistics::new_unknown() - ); - } - - #[test] - fn row_count_statistics_do_not_require_columns() { - let file_schema = Schema::new(vec![Arc::new(datafusion::arrow::datatypes::Field::new( - "a", - datafusion::arrow::datatypes::DataType::Int32, - false, - ))]); - let projected_statistics = Statistics { - num_rows: Precision::Exact(100), - total_byte_size: Precision::Absent, - column_statistics: vec![], - }; - - let statistics = expand_projected_statistics(&file_schema, &[], projected_statistics); - - assert_eq!(statistics.num_rows, Precision::Exact(100)); - assert_eq!( - statistics.column_statistics, - vec![datafusion_common::ColumnStatistics::new_unknown()] - ); - } -} diff --git a/crates/sail-data-source/src/listing/planner.rs b/crates/sail-data-source/src/listing/planner.rs index 979b793820..c16631c38a 100644 --- a/crates/sail-data-source/src/listing/planner.rs +++ b/crates/sail-data-source/src/listing/planner.rs @@ -56,8 +56,14 @@ pub(crate) async fn prewarm_file_statistics( source: &ListingTableSource, ctx: &dyn Session, ) -> datafusion_common::Result<()> { - if source.config().collect_stat { - let _ = list_files_for_scan(source, ctx, &[], None, Some(&[])).await?; + if source.config().collect_stat + && ctx + .runtime_env() + .cache_manager + .get_file_statistic_cache_limit() + > 0 + { + let _ = list_files_for_scan(source, ctx, &[], None).await?; } Ok(()) } @@ -114,17 +120,6 @@ impl ExtensionPlanner for ListingPhysicalPlanner { }); let statistic_file_limit = if filters.is_empty() { limit } else { None }; - let file_column_count = source.config().schema.file_schema().fields().len(); - let statistics_columns = projection.as_ref().map(|projection| { - let mut columns = projection - .iter() - .copied() - .filter(|index| *index < file_column_count) - .collect::>(); - columns.sort_unstable(); - columns.dedup(); - columns - }); let ListFilesResult { mut file_groups, @@ -135,7 +130,6 @@ impl ExtensionPlanner for ListingPhysicalPlanner { session_state, &partition_filters, statistic_file_limit, - statistics_columns.as_deref(), ) .await?; @@ -353,7 +347,6 @@ async fn list_files_for_scan<'a>( ctx: &'a dyn Session, filters: &'a [Expr], limit: Option, - statistics_columns: Option<&'a [usize]>, ) -> datafusion_common::Result { let store = if let Some(url) = source.config().table_paths.first() { ctx.runtime_env().object_store(url)? @@ -394,21 +387,12 @@ async fn list_files_for_scan<'a>( let meta_fetch_concurrency = ctx.config_options().execution.meta_fetch_concurrency; let file_list = stream::iter(file_list).flatten_unordered(meta_fetch_concurrency); - let statistics_cache_table = source.statistics_cache_table(statistics_columns); let files = file_list .map(|part_file| async { let part_file = part_file?; let (statistics, ordering) = if source.config().collect_stat { - do_collect_statistics_and_ordering( - source, - ctx, - &store, - &part_file, - statistics_columns, - &statistics_cache_table, - ) - .await? + do_collect_statistics_and_ordering(source, ctx, &store, &part_file).await? } else { ( Arc::new(Statistics::new_unknown( @@ -466,14 +450,15 @@ async fn do_collect_statistics_and_ordering( ctx: &dyn Session, store: &Arc, part_file: &datafusion_datasource::PartitionedFile, - statistics_columns: Option<&[usize]>, - statistics_cache_table: &datafusion_common::TableReference, ) -> datafusion_common::Result<(Arc, Option)> { let meta = &part_file.object_meta; let file_schema = source.config().schema.file_schema(); let file_statistic_cache = ctx.runtime_env().cache_manager.get_file_statistic_cache(); let cache_key = TableScopedPath { - table: Some(statistics_cache_table.clone()), + table: Some(statistics_cache_table_ref( + part_file.table_reference.as_ref(), + file_schema.as_ref(), + )), path: meta.location.clone(), }; @@ -497,7 +482,6 @@ async fn do_collect_statistics_and_ordering( store, meta, source.config().schema.file_schema().clone(), - statistics_columns, source.config().compression, ) .await?; @@ -517,6 +501,28 @@ async fn do_collect_statistics_and_ordering( Ok((statistics, file_meta.ordering)) } +fn statistics_cache_table_ref( + table: Option<&datafusion_common::TableReference>, + schema: &arrow::datatypes::Schema, +) -> datafusion_common::TableReference { + let mut hasher = std::collections::hash_map::DefaultHasher::new(); + + for field in schema.fields() { + std::hash::Hash::hash(field.name(), &mut hasher); + std::hash::Hash::hash(&field.data_type().to_string(), &mut hasher); + std::hash::Hash::hash(&field.is_nullable(), &mut hasher); + } + + let table = table + .map(ToString::to_string) + .unwrap_or_else(|| "listing".to_string()); + + datafusion_common::TableReference::bare(format!( + "{table}_{:016x}", + std::hash::Hasher::finish(&hasher) + )) +} + async fn get_files_with_limit( files: impl Stream>, limit: Option, @@ -560,3 +566,120 @@ async fn get_files_with_limit( let inexact_stats = all_files.next().await.is_some(); Ok((file_group, inexact_stats)) } + +#[cfg(test)] +#[expect(clippy::unwrap_used)] +mod tests { + use std::sync::atomic::{AtomicUsize, Ordering}; + + use bytes::Bytes; + use datafusion::arrow::datatypes::{Field, Schema, SchemaRef}; + use datafusion::prelude::SessionContext; + use datafusion_common::Constraints; + use datafusion_common::parsers::CompressionTypeVariant; + use datafusion_datasource::TableSchema; + use object_store::ObjectStoreExt; + use object_store::memory::InMemory; + use object_store::path::Path; + use url::Url; + + use super::*; + use crate::listing::source::{ + ListingFileMeta, ListingFileSample, ListingScanInput, ReadFormat, + }; + use crate::listing::table::ListingTableSourceConfig; + + #[derive(Debug)] + struct CountingReadFormat { + inference_count: Arc, + } + + #[async_trait] + impl ReadFormat for CountingReadFormat { + async fn infer_compression( + &self, + _ctx: &dyn Session, + _files: &[ListingFileSample<'_>], + ) -> datafusion_common::Result { + unreachable!() + } + + async fn infer_schema( + &self, + _ctx: &dyn Session, + _files: &[ListingFileSample<'_>], + _compression: CompressionTypeVariant, + ) -> datafusion_common::Result { + unreachable!() + } + + async fn infer_file_meta( + &self, + _ctx: &dyn Session, + _store: &Arc, + _object: &object_store::ObjectMeta, + file_schema: SchemaRef, + _compression: CompressionTypeVariant, + ) -> datafusion_common::Result { + self.inference_count.fetch_add(1, Ordering::Relaxed); + let mut statistics = Statistics::new_unknown(&file_schema); + statistics.num_rows = Precision::Exact(1); + Ok(ListingFileMeta { + statistics, + ordering: None, + }) + } + + async fn scan( + &self, + _ctx: &dyn Session, + _input: ListingScanInput, + ) -> datafusion_common::Result { + unreachable!() + } + } + + #[tokio::test] + async fn prewarmed_statistics_are_reused_by_scan() { + let object_store: Arc = Arc::new(InMemory::new()); + object_store + .put( + &Path::from("table/data"), + Bytes::from_static(b"data").into(), + ) + .await + .unwrap(); + + let context = SessionContext::new(); + context.register_object_store(&Url::parse("memory://").unwrap(), Arc::clone(&object_store)); + + let inference_count = Arc::new(AtomicUsize::new(0)); + let source = ListingTableSource::try_new(ListingTableSourceConfig { + table_paths: vec![ListingTableUrl::parse("memory:///table/").unwrap()], + schema: TableSchema::from_file_schema(Arc::new(Schema::new(vec![Field::new( + "value", + DataType::Int64, + false, + )]))), + constraints: Constraints::default(), + file_sort_order: vec![], + collect_stat: true, + target_partitions: 1, + read_format: Arc::new(CountingReadFormat { + inference_count: Arc::clone(&inference_count), + }), + compression: CompressionTypeVariant::UNCOMPRESSED, + }) + .unwrap(); + let state = context.state(); + + prewarm_file_statistics(&source, &state).await.unwrap(); + assert_eq!(inference_count.load(Ordering::Relaxed), 1); + + let result = list_files_for_scan(&source, &state, &[], None) + .await + .unwrap(); + assert_eq!(inference_count.load(Ordering::Relaxed), 1); + assert_eq!(result.statistics.num_rows, Precision::Exact(1)); + } +} diff --git a/crates/sail-data-source/src/listing/source.rs b/crates/sail-data-source/src/listing/source.rs index a8b6fcd836..25647fed97 100644 --- a/crates/sail-data-source/src/listing/source.rs +++ b/crates/sail-data-source/src/listing/source.rs @@ -67,17 +67,15 @@ pub trait ReadFormat: Debug + Send + Sync + 'static { /// Infer file-level metadata needed for planning. /// The metadata includes statistics and ordering. - /// `statistics_columns` is `None` for all file columns and `Some` for a selected subset. async fn infer_file_meta( &self, ctx: &dyn Session, store: &Arc, object: &ObjectMeta, file_schema: SchemaRef, - statistics_columns: Option<&[usize]>, compression: CompressionTypeVariant, ) -> Result { - let _ = (ctx, store, object, statistics_columns, compression); + let _ = (ctx, store, object, compression); Ok(ListingFileMeta { statistics: Statistics::new_unknown(&file_schema), ordering: None, diff --git a/crates/sail-data-source/src/listing/table.rs b/crates/sail-data-source/src/listing/table.rs index 1d261ea31b..2ef187aec2 100644 --- a/crates/sail-data-source/src/listing/table.rs +++ b/crates/sail-data-source/src/listing/table.rs @@ -1,14 +1,13 @@ // The listing table source is adapted from the DataFusion `ListingTable` implementation. // [CREDIT]: https://github.com/apache/datafusion/blob/53.1.0/datafusion/catalog-listing/src/table.rs -use std::hash::{Hash, Hasher}; use std::sync::Arc; use datafusion::arrow::datatypes::SchemaRef; use datafusion::logical_expr::expr::Sort; use datafusion::logical_expr::{Expr, TableProviderFilterPushDown, TableSource, TableType}; use datafusion_common::parsers::CompressionTypeVariant; -use datafusion_common::{Constraints, Result, TableReference}; +use datafusion_common::{Constraints, Result}; use datafusion_datasource::{ListingTableUrl, TableSchema}; use crate::listing::source::ReadFormat; @@ -33,81 +32,16 @@ pub struct ListingTableSourceConfig { #[derive(Clone, Debug)] pub struct ListingTableSource { config: ListingTableSourceConfig, - statistics_cache_namespace: String, } impl ListingTableSource { pub fn try_new(config: ListingTableSourceConfig) -> Result { - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - for path in &config.table_paths { - path.to_string().hash(&mut hasher); - } - for field in config.schema.file_schema().fields() { - field.name().hash(&mut hasher); - field.data_type().to_string().hash(&mut hasher); - field.is_nullable().hash(&mut hasher); - } - format!("{:?}", config.read_format).hash(&mut hasher); - format!("{:?}", config.compression).hash(&mut hasher); - let statistics_cache_namespace = format!("listing_{:016x}", hasher.finish()); - - Ok(Self { - config, - statistics_cache_namespace, - }) + Ok(Self { config }) } pub fn config(&self) -> &ListingTableSourceConfig { &self.config } - - pub(crate) fn statistics_cache_table( - &self, - statistics_columns: Option<&[usize]>, - ) -> TableReference { - let coverage = statistics_cache_coverage(statistics_columns); - TableReference::bare(format!("{}_{}", self.statistics_cache_namespace, coverage)) - } -} - -fn statistics_cache_coverage(statistics_columns: Option<&[usize]>) -> String { - match statistics_columns { - None => "all".to_string(), - Some([]) => "rows".to_string(), - Some(columns) => { - let mut columns = columns.to_vec(); - columns.sort_unstable(); - columns.dedup(); - let mut hasher = std::collections::hash_map::DefaultHasher::new(); - columns.hash(&mut hasher); - format!("columns_{:016x}", hasher.finish()) - } - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn statistics_cache_coverage_is_canonical_and_isolated() { - assert_ne!( - statistics_cache_coverage(None), - statistics_cache_coverage(Some(&[])) - ); - assert_ne!( - statistics_cache_coverage(Some(&[])), - statistics_cache_coverage(Some(&[1])) - ); - assert_ne!( - statistics_cache_coverage(Some(&[1])), - statistics_cache_coverage(Some(&[2])) - ); - assert_eq!( - statistics_cache_coverage(Some(&[2, 1])), - statistics_cache_coverage(Some(&[1, 2, 2])) - ); - } } impl TableSource for ListingTableSource { From c21fc777fb7dcc5561239d71c016cd5586e8e7b7 Mon Sep 17 00:00:00 2001 From: XL Liang Date: Fri, 24 Jul 2026 00:13:57 +0800 Subject: [PATCH 7/9] perf: prewarm listing statistics after catalog registration --- crates/sail-catalog/src/command.rs | 172 ++++++++++++++- .../sail-common-datafusion/src/datasource.rs | 14 ++ .../sail-common-datafusion/src/extension.rs | 18 ++ .../sail-data-source/src/listing/planner.rs | 118 +++++++---- crates/sail-data-source/src/listing/source.rs | 199 ++++++++++++++++-- .../src/resolver/command/catalog/table.rs | 5 +- .../sail-plan/src/resolver/command/write.rs | 1 + 7 files changed, 465 insertions(+), 62 deletions(-) diff --git a/crates/sail-catalog/src/command.rs b/crates/sail-catalog/src/command.rs index ba99cf2dfd..ff7700343e 100644 --- a/crates/sail-catalog/src/command.rs +++ b/crates/sail-catalog/src/command.rs @@ -1,10 +1,11 @@ use datafusion::arrow::array::RecordBatch; -use datafusion::arrow::datatypes::SchemaRef; +use datafusion::arrow::datatypes::{Schema, SchemaRef}; +use datafusion_common::Constraints; use sail_common_datafusion::array::serde::ArrowSerializer; -use sail_common_datafusion::catalog::{FunctionStatus, LakehouseOperation}; +use sail_common_datafusion::catalog::{FunctionStatus, LakehouseOperation, TableKind}; use sail_common_datafusion::datasource::{ - TableFormatAlterTableOperation, TableFormatCreateTableColumn, TableFormatCreateTableInfo, - TableFormatRegistry, is_lakehouse_format, + OptionLayer, SourceInfo, TableFormatAlterTableOperation, TableFormatCreateTableColumn, + TableFormatCreateTableInfo, TableFormatRegistry, is_lakehouse_format, }; use sail_common_datafusion::extension::SessionExtensionAccessor; use sail_common_datafusion::session::plan::PlanService; @@ -23,6 +24,11 @@ use crate::provider::{ }; use crate::utils::{quote_names_if_needed, quote_namespace_if_needed}; +#[derive(Debug, Clone, Eq, PartialEq, PartialOrd, Hash, Serialize, Deserialize)] +pub struct CreateTableStatisticsWarmup { + pub read_case_sensitive: bool, +} + #[derive(Debug, Clone, Eq, PartialEq, PartialOrd, Hash, Serialize, Deserialize)] pub enum CatalogCommand { CurrentCatalog, @@ -57,6 +63,8 @@ pub enum CatalogCommand { CreateTable { table: Vec, options: CreateTableOptions, + #[serde(default)] + statistics_warmup: Option, }, TableExists { table: Vec, @@ -326,7 +334,11 @@ impl CatalogCommand { manager.drop_database(&database, options).await?; display.bools().to_record_batch(vec![true])? } - CatalogCommand::CreateTable { table, options } => { + CatalogCommand::CreateTable { + table, + options, + statistics_warmup, + } => { let existed_before = if options.mode.ignore_if_exists() { match manager.get_table_or_view(&table).await { Ok(_) => true, @@ -351,6 +363,15 @@ impl CatalogCommand { &create_plan, ) .await?; + if let Some(warmup) = statistics_warmup + && let Err(error) = + prewarm_created_table_statistics(ctx, &status, warmup).await + { + log::warn!( + "failed to prewarm file statistics for catalog table '{}': {error}", + table.join(".") + ); + } } display.bools().to_record_batch(vec![true])? } @@ -731,6 +752,74 @@ impl CatalogCommand { } } +async fn prewarm_created_table_statistics( + ctx: &C, + status: &sail_common_datafusion::catalog::TableStatus, + warmup: CreateTableStatisticsWarmup, +) -> CatalogResult<()> { + let session_config = ctx.session_config(); + let runtime_env = ctx.runtime_env(); + if !session_config.collect_statistics() + || runtime_env.cache_manager.get_file_statistic_cache_limit() == 0 + { + return Ok(()); + } + + let TableKind::Table { + columns, + constraints: _, + location: Some(location), + format, + partition_by, + sort_by, + bucket_by, + properties, + .. + } = &status.kind + else { + return Ok(()); + }; + + let registry = ctx.extension::().map_err(|error| { + CatalogError::External(format!( + "missing TableFormatRegistry for statistics warmup on format '{format}': {error}" + )) + })?; + let table_format = registry.get(format).map_err(|error| { + CatalogError::External(format!( + "unknown table format '{format}' for statistics warmup: {error}" + )) + })?; + table_format + .prewarm_statistics( + session_config, + runtime_env, + SourceInfo { + paths: vec![location.clone()], + lakehouse_table: None, + schema: Some(Schema::new( + columns + .iter() + .map(|column| column.field()) + .collect::>(), + )), + constraints: Constraints::default(), + partition_by: partition_by + .iter() + .map(|field| field.column.clone()) + .collect(), + bucket_by: bucket_by.clone().map(Into::into), + sort_order: sort_by.iter().cloned().map(Into::into).collect(), + options: vec![OptionLayer::TablePropertyList { + items: properties.clone(), + }], + read_case_sensitive: warmup.read_case_sensitive, + }, + ) + .await + .map_err(CatalogError::DataFusionError) +} + async fn prepare_create_table_storage_metadata( ctx: &C, manager: &CatalogManager, @@ -1078,6 +1167,7 @@ struct DescribeFunctionRow { #[cfg(test)] mod tests { use std::sync::Arc; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use async_trait::async_trait; use datafusion::catalog::Session; @@ -1212,7 +1302,7 @@ mod tests { _table: &str, _options: CreateTableOptions, ) -> CatalogResult { - unreachable!() + Ok(self.table_status.clone()) } async fn get_table( @@ -1301,7 +1391,10 @@ mod tests { } } - struct TestTableFormat; + struct TestTableFormat { + warmup_count: Arc, + read_case_sensitive: Arc, + } #[async_trait] impl TableFormat for TestTableFormat { @@ -1325,6 +1418,18 @@ mod tests { not_impl_err!("unused in test") } + async fn prewarm_statistics( + &self, + _session_config: SessionConfig, + _runtime_env: Arc, + info: SourceInfo, + ) -> datafusion_common::Result<()> { + self.warmup_count.fetch_add(1, Ordering::Relaxed); + self.read_case_sensitive + .store(info.read_case_sensitive, Ordering::Relaxed); + Ok(()) + } + async fn alter_table( &self, _runtime_env: Arc, @@ -1337,8 +1442,18 @@ mod tests { } fn test_session_context() -> SessionContext { + test_session_context_with_statistics_warmup().0 + } + + fn test_session_context_with_statistics_warmup() + -> (SessionContext, Arc, Arc) { + let warmup_count = Arc::new(AtomicUsize::new(0)); + let read_case_sensitive = Arc::new(AtomicBool::new(false)); let registry = Arc::new(TableFormatRegistry::new()); - let register_result = registry.register(Arc::new(TestTableFormat)); + let register_result = registry.register(Arc::new(TestTableFormat { + warmup_count: Arc::clone(&warmup_count), + read_case_sensitive: Arc::clone(&read_case_sensitive), + })); assert!( register_result.is_ok(), "failed to register test table format: {register_result:?}" @@ -1350,7 +1465,11 @@ mod tests { let config = SessionConfig::new() .with_extension(registry) .with_extension(plan_service); - SessionContext::new_with_config(config) + ( + SessionContext::new_with_config(config), + warmup_count, + read_case_sensitive, + ) } fn test_manager(alter_error: Option<&str>) -> CatalogManager { @@ -1420,6 +1539,41 @@ mod tests { manager } + #[tokio::test] + async fn create_table_prewarms_statistics_after_registration() { + let (ctx, warmup_count, read_case_sensitive) = + test_session_context_with_statistics_warmup(); + let manager = test_manager_with_catalog_behavior(None, false, false); + let command = CatalogCommand::CreateTable { + table: vec!["items".to_string()], + options: CreateTableOptions { + columns: vec![], + comment: None, + constraints: vec![], + location: Some("s3://bucket/items".to_string()), + format: "delta".to_string(), + partition_by: vec![], + sort_by: vec![], + bucket_by: None, + mode: Default::default(), + properties: vec![], + is_external: true, + is_write_precondition: false, + }, + statistics_warmup: Some(CreateTableStatisticsWarmup { + read_case_sensitive: true, + }), + }; + + let result = command.execute(&ctx, &manager).await; + assert!( + result.is_ok(), + "expected CREATE TABLE to succeed, got {result:?}" + ); + assert_eq!(warmup_count.load(Ordering::Relaxed), 1); + assert!(read_case_sensitive.load(Ordering::Relaxed)); + } + #[tokio::test] async fn table_or_view_lookup_treats_unsupported_persistent_views_as_missing() { let manager = test_manager_with_catalog_behavior(None, false, false); diff --git a/crates/sail-common-datafusion/src/datasource.rs b/crates/sail-common-datafusion/src/datasource.rs index 768bd4b62f..0802729db3 100644 --- a/crates/sail-common-datafusion/src/datasource.rs +++ b/crates/sail-common-datafusion/src/datasource.rs @@ -6,6 +6,7 @@ use chrono::{DateTime, Utc}; use datafusion::arrow::datatypes::{DataType, FieldRef, Schema, SchemaRef}; use datafusion::catalog::Session; use datafusion::common::plan_datafusion_err; +use datafusion::execution::context::SessionConfig; use datafusion::execution::runtime_env::RuntimeEnv; use datafusion::logical_expr::LogicalPlan; use datafusion::physical_expr::{ @@ -487,6 +488,19 @@ pub trait TableFormat: Send + Sync { info: SourceInfo, ) -> Result>; + /// Prewarms reusable physical statistics after a catalog table is registered. + /// + /// Formats without file-level statistics can keep the default no-op. + async fn prewarm_statistics( + &self, + session_config: SessionConfig, + runtime_env: Arc, + info: SourceInfo, + ) -> Result<()> { + let _ = (session_config, runtime_env, info); + Ok(()) + } + /// Infers the logical schema for planning without requiring callers to construct a read source. async fn infer_schema(&self, ctx: &dyn Session, info: SourceInfo) -> Result { Ok(self.create_source(ctx, info).await?.schema()) diff --git a/crates/sail-common-datafusion/src/extension.rs b/crates/sail-common-datafusion/src/extension.rs index 30fbbd03a9..cc9bca0aa6 100644 --- a/crates/sail-common-datafusion/src/extension.rs +++ b/crates/sail-common-datafusion/src/extension.rs @@ -1,6 +1,7 @@ use std::sync::Arc; use datafusion::catalog::Session; +use datafusion::execution::context::SessionConfig; use datafusion::execution::runtime_env::RuntimeEnv; use datafusion::execution::{SessionState, TaskContext}; use datafusion::prelude::SessionContext; @@ -12,6 +13,7 @@ pub trait SessionExtension: Send + Sync + 'static { pub trait SessionExtensionAccessor { fn extension(&self) -> Result>; + fn session_config(&self) -> SessionConfig; fn runtime_env(&self) -> Arc; } @@ -24,6 +26,10 @@ impl SessionExtensionAccessor for SessionContext { .ok_or_else(|| internal_datafusion_err!("session extension not found: {}", T::name())) } + fn session_config(&self) -> SessionConfig { + self.state_ref().read().config().clone() + } + fn runtime_env(&self) -> Arc { self.state_ref().read().runtime_env().clone() } @@ -36,6 +42,10 @@ impl SessionExtensionAccessor for SessionState { .ok_or_else(|| internal_datafusion_err!("session extension not found: {}", T::name())) } + fn session_config(&self) -> SessionConfig { + self.config().clone() + } + fn runtime_env(&self) -> Arc { self.runtime_env().clone() } @@ -48,6 +58,10 @@ impl SessionExtensionAccessor for &dyn Session { .ok_or_else(|| internal_datafusion_err!("session extension not found: {}", T::name())) } + fn session_config(&self) -> SessionConfig { + self.config().clone() + } + fn runtime_env(&self) -> Arc { Session::runtime_env(*self).clone() } @@ -60,6 +74,10 @@ impl SessionExtensionAccessor for TaskContext { .ok_or_else(|| internal_datafusion_err!("session extension not found: {}", T::name())) } + fn session_config(&self) -> SessionConfig { + TaskContext::session_config(self).clone() + } + fn runtime_env(&self) -> Arc { self.runtime_env().clone() } diff --git a/crates/sail-data-source/src/listing/planner.rs b/crates/sail-data-source/src/listing/planner.rs index c16631c38a..9f7646e7a6 100644 --- a/crates/sail-data-source/src/listing/planner.rs +++ b/crates/sail-data-source/src/listing/planner.rs @@ -56,16 +56,34 @@ pub(crate) async fn prewarm_file_statistics( source: &ListingTableSource, ctx: &dyn Session, ) -> datafusion_common::Result<()> { - if source.config().collect_stat - && ctx + if !source.config().collect_stat + || ctx .runtime_env() .cache_manager .get_file_statistic_cache_limit() - > 0 + == 0 { - let _ = list_files_for_scan(source, ctx, &[], None).await?; + return Ok(()); } - Ok(()) + + let Some(url) = source.config().table_paths.first() else { + return Ok(()); + }; + let store = ctx.runtime_env().object_store(url)?; + let partition_cols = listing_partition_columns(source); + let file_list = + list_partitioned_files(source, ctx, store.as_ref(), &[], &partition_cols).await?; + let meta_fetch_concurrency = ctx.config_options().execution.meta_fetch_concurrency; + + file_list + .map(|part_file| async { + let part_file = part_file?; + do_collect_statistics_and_ordering(source, ctx, &store, &part_file).await?; + Ok(()) + }) + .buffer_unordered(meta_fetch_concurrency) + .try_for_each(|()| future::ready(Ok(()))) + .await } #[async_trait] @@ -358,35 +376,10 @@ async fn list_files_for_scan<'a>( }); }; - let partition_cols: Vec<(String, DataType)> = source - .config() - .schema - .table_partition_cols() - .iter() - .map(|field| (field.name().clone(), field.data_type().clone())) - .collect(); - - let store_ref = store.as_ref(); - let partition_cols_ref = partition_cols.as_slice(); - let file_list = future::try_join_all(source.config().table_paths.iter().map( - |table_path| async move { - let files = - pruned_partition_list(ctx, store_ref, table_path, filters, "", partition_cols_ref) - .await?; - // Skip hidden files so scans agree with `inputFiles` and schema sampling. - let files = files.try_filter(move |file| { - futures::future::ready(!has_hidden_path_component( - table_path, - &file.object_meta.location, - )) - }); - Ok::<_, datafusion_common::DataFusionError>(files.boxed()) - }, - )) - .await?; - + let partition_cols = listing_partition_columns(source); + let file_list = + list_partitioned_files(source, ctx, store.as_ref(), filters, &partition_cols).await?; let meta_fetch_concurrency = ctx.config_options().execution.meta_fetch_concurrency; - let file_list = stream::iter(file_list).flatten_unordered(meta_fetch_concurrency); let files = file_list .map(|part_file| async { @@ -445,6 +438,50 @@ async fn list_files_for_scan<'a>( }) } +fn listing_partition_columns(source: &ListingTableSource) -> Vec<(String, DataType)> { + source + .config() + .schema + .table_partition_cols() + .iter() + .map(|field| (field.name().clone(), field.data_type().clone())) + .collect() +} + +async fn list_partitioned_files<'a>( + source: &'a ListingTableSource, + ctx: &'a dyn Session, + store: &'a dyn ObjectStore, + filters: &'a [Expr], + partition_cols: &'a [(String, DataType)], +) -> datafusion_common::Result< + futures::stream::BoxStream< + 'a, + datafusion_common::Result, + >, +> { + let file_list = future::try_join_all(source.config().table_paths.iter().map( + |table_path| async move { + let files = + pruned_partition_list(ctx, store, table_path, filters, "", partition_cols).await?; + // Skip hidden files so scans agree with `inputFiles` and schema sampling. + let files = files.try_filter(move |file| { + futures::future::ready(!has_hidden_path_component( + table_path, + &file.object_meta.location, + )) + }); + Ok::<_, datafusion_common::DataFusionError>(files.boxed()) + }, + )) + .await?; + + let meta_fetch_concurrency = ctx.config_options().execution.meta_fetch_concurrency; + Ok(stream::iter(file_list) + .flatten_unordered(meta_fetch_concurrency) + .boxed()) +} + async fn do_collect_statistics_and_ordering( source: &ListingTableSource, ctx: &dyn Session, @@ -640,7 +677,7 @@ mod tests { } #[tokio::test] - async fn prewarmed_statistics_are_reused_by_scan() { + async fn prewarmed_statistics_are_reused_after_source_recreation() { let object_store: Arc = Arc::new(InMemory::new()); object_store .put( @@ -654,7 +691,7 @@ mod tests { context.register_object_store(&Url::parse("memory://").unwrap(), Arc::clone(&object_store)); let inference_count = Arc::new(AtomicUsize::new(0)); - let source = ListingTableSource::try_new(ListingTableSourceConfig { + let config = ListingTableSourceConfig { table_paths: vec![ListingTableUrl::parse("memory:///table/").unwrap()], schema: TableSchema::from_file_schema(Arc::new(Schema::new(vec![Field::new( "value", @@ -669,14 +706,17 @@ mod tests { inference_count: Arc::clone(&inference_count), }), compression: CompressionTypeVariant::UNCOMPRESSED, - }) - .unwrap(); + }; + let warmup_source = ListingTableSource::try_new(config.clone()).unwrap(); + let scan_source = ListingTableSource::try_new(config).unwrap(); let state = context.state(); - prewarm_file_statistics(&source, &state).await.unwrap(); + prewarm_file_statistics(&warmup_source, &state) + .await + .unwrap(); assert_eq!(inference_count.load(Ordering::Relaxed), 1); - let result = list_files_for_scan(&source, &state, &[], None) + let result = list_files_for_scan(&scan_source, &state, &[], None) .await .unwrap(); assert_eq!(inference_count.load(Ordering::Relaxed), 1); diff --git a/crates/sail-data-source/src/listing/source.rs b/crates/sail-data-source/src/listing/source.rs index 25647fed97..430a6accd5 100644 --- a/crates/sail-data-source/src/listing/source.rs +++ b/crates/sail-data-source/src/listing/source.rs @@ -7,6 +7,8 @@ use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; use datafusion::catalog::Session; use datafusion::datasource::physical_plan::FileSinkConfig; use datafusion::execution::object_store::ObjectStoreUrl; +use datafusion::execution::runtime_env::RuntimeEnv; +use datafusion::execution::{SessionStateBuilder, context::SessionConfig}; use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder, TableSource}; use datafusion::physical_expr::LexRequirement; use datafusion::physical_expr_common::sort_expr::LexOrdering; @@ -141,17 +143,12 @@ pub struct ListingTableFormat { phantom: PhantomData, } -#[async_trait] -impl TableFormat for ListingTableFormat { - fn name(&self) -> &str { - T::name() - } - - async fn create_source( +impl ListingTableFormat { + async fn create_listing_source( &self, ctx: &dyn Session, info: SourceInfo, - ) -> Result> { + ) -> Result { let SourceInfo { paths, lakehouse_table: _, @@ -235,7 +232,7 @@ impl TableFormat for ListingTableFormat { validate_partitions(&sampled_files, &partition_fields)?; - let source = ListingTableSource::try_new(ListingTableSourceConfig { + ListingTableSource::try_new(ListingTableSourceConfig { table_paths: urls, schema: TableSchema::new(schema, partition_fields), constraints, @@ -244,11 +241,47 @@ impl TableFormat for ListingTableFormat { target_partitions: ctx.config().target_partitions(), read_format: Arc::new(read_format), compression, - })?; - if let Err(error) = prewarm_file_statistics(&source, ctx).await { - log::warn!("failed to prewarm listing file statistics: {error}"); + }) + } +} + +#[async_trait] +impl TableFormat for ListingTableFormat { + fn name(&self) -> &str { + T::name() + } + + async fn create_source( + &self, + ctx: &dyn Session, + info: SourceInfo, + ) -> Result> { + Ok(Arc::new(self.create_listing_source(ctx, info).await?)) + } + + async fn prewarm_statistics( + &self, + session_config: SessionConfig, + runtime_env: Arc, + info: SourceInfo, + ) -> Result<()> { + if !session_config.collect_statistics() + || runtime_env.cache_manager.get_file_statistic_cache_limit() == 0 + { + return Ok(()); } - Ok(Arc::new(source)) + + let session_state = SessionStateBuilder::new() + .with_config(session_config) + .with_runtime_env(runtime_env) + .build(); + let source = self.create_listing_source(&session_state, info).await?; + prewarm_file_statistics(&source, &session_state).await?; + log::debug!( + "prewarmed listing file statistics cache for {} path(s)", + source.config().table_paths.len() + ); + Ok(()) } async fn create_writer(&self, ctx: &dyn Session, info: SinkInfo) -> Result { @@ -354,3 +387,143 @@ fn reconcile_schema_names_case_insensitive(schema: Schema, physical: &Schema) -> } Ok(Schema::new_with_metadata(fields, schema.metadata().clone())) } + +#[cfg(test)] +#[expect(clippy::unwrap_used)] +mod tests { + use std::sync::atomic::{AtomicUsize, Ordering}; + + use bytes::Bytes; + use datafusion::arrow::datatypes::{Field, SchemaRef}; + use datafusion::prelude::SessionContext; + use datafusion_common::Constraints; + use object_store::memory::InMemory; + use object_store::path::Path; + use object_store::{ObjectStore, ObjectStoreExt}; + use url::Url; + + use super::*; + + static FILE_META_INFERENCE_COUNT: AtomicUsize = AtomicUsize::new(0); + + #[derive(Debug, Default)] + struct TestFormatFactory; + + #[derive(Debug)] + struct TestReadFormat; + + #[derive(Debug)] + struct TestWriteFormat; + + impl FormatFactory for TestFormatFactory { + type Read = TestReadFormat; + type Write = TestWriteFormat; + + fn name() -> &'static str { + "test" + } + + fn read(_ctx: &dyn Session, _options: Vec) -> Result { + Ok(TestReadFormat) + } + + fn write(_ctx: &dyn Session, _options: Vec) -> Result { + Ok(TestWriteFormat) + } + } + + #[async_trait] + impl ReadFormat for TestReadFormat { + async fn infer_compression( + &self, + _ctx: &dyn Session, + _files: &[ListingFileSample<'_>], + ) -> Result { + Ok(CompressionTypeVariant::UNCOMPRESSED) + } + + async fn infer_schema( + &self, + _ctx: &dyn Session, + _files: &[ListingFileSample<'_>], + _compression: CompressionTypeVariant, + ) -> Result { + unreachable!() + } + + async fn infer_file_meta( + &self, + _ctx: &dyn Session, + _store: &Arc, + _object: &ObjectMeta, + file_schema: SchemaRef, + _compression: CompressionTypeVariant, + ) -> Result { + FILE_META_INFERENCE_COUNT.fetch_add(1, Ordering::Relaxed); + Ok(ListingFileMeta { + statistics: Statistics::new_unknown(&file_schema), + ordering: None, + }) + } + + async fn scan( + &self, + _ctx: &dyn Session, + _input: ListingScanInput, + ) -> Result { + unreachable!() + } + } + + #[async_trait] + impl WriteFormat for TestWriteFormat { + async fn sink( + &self, + _ctx: &dyn Session, + _input: ListingSinkInput, + ) -> Result> { + unreachable!() + } + } + + #[tokio::test] + async fn create_source_does_not_prewarm_file_statistics() { + FILE_META_INFERENCE_COUNT.store(0, Ordering::Relaxed); + let object_store: Arc = Arc::new(InMemory::new()); + object_store + .put( + &Path::from("table/data"), + Bytes::from_static(b"data").into(), + ) + .await + .unwrap(); + let context = SessionContext::new(); + context.register_object_store(&Url::parse("memory://").unwrap(), object_store); + let state = context.state(); + let format = ListingTableFormat::::default(); + let info = SourceInfo { + paths: vec!["memory:///table/".to_string()], + lakehouse_table: None, + schema: Some(Schema::new(vec![Field::new( + "value", + DataType::Int64, + false, + )])), + constraints: Constraints::default(), + partition_by: vec![], + bucket_by: None, + sort_order: vec![], + options: vec![], + read_case_sensitive: true, + }; + + format.create_source(&state, info.clone()).await.unwrap(); + assert_eq!(FILE_META_INFERENCE_COUNT.load(Ordering::Relaxed), 0); + + format + .prewarm_statistics(state.config().clone(), state.runtime_env().clone(), info) + .await + .unwrap(); + assert_eq!(FILE_META_INFERENCE_COUNT.load(Ordering::Relaxed), 1); + } +} diff --git a/crates/sail-plan/src/resolver/command/catalog/table.rs b/crates/sail-plan/src/resolver/command/catalog/table.rs index 469b843f6b..a1ca8ef006 100644 --- a/crates/sail-plan/src/resolver/command/catalog/table.rs +++ b/crates/sail-plan/src/resolver/command/catalog/table.rs @@ -1,5 +1,5 @@ use datafusion_expr::LogicalPlan; -use sail_catalog::command::CatalogCommand; +use sail_catalog::command::{CatalogCommand, CreateTableStatisticsWarmup}; use sail_catalog::manager::CatalogManager; use sail_catalog::provider::{ AlterTableOptions, CatalogPartitionField, CreateTableColumnOptions, CreateTableOptions, @@ -90,6 +90,9 @@ impl PlanResolver<'_> { is_external, is_write_precondition: false, }, + statistics_warmup: Some(CreateTableStatisticsWarmup { + read_case_sensitive: self.config.case_sensitive, + }), }; self.resolve_catalog_command(command) } diff --git a/crates/sail-plan/src/resolver/command/write.rs b/crates/sail-plan/src/resolver/command/write.rs index 4abb3943d1..321a8d2dac 100644 --- a/crates/sail-plan/src/resolver/command/write.rs +++ b/crates/sail-plan/src/resolver/command/write.rs @@ -505,6 +505,7 @@ impl PlanResolver<'_> { let command = CatalogCommand::CreateTable { table: table.clone().into(), options: create_options, + statistics_warmup: None, }; preconditions.push(Arc::new(self.resolve_catalog_command(command)?)); } From 25e1ea913506a2990bd245020f478f896aa0413c Mon Sep 17 00:00:00 2001 From: XL Liang Date: Fri, 24 Jul 2026 00:17:03 +0800 Subject: [PATCH 8/9] fmt --- crates/sail-data-source/src/listing/source.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/crates/sail-data-source/src/listing/source.rs b/crates/sail-data-source/src/listing/source.rs index 430a6accd5..f3af011781 100644 --- a/crates/sail-data-source/src/listing/source.rs +++ b/crates/sail-data-source/src/listing/source.rs @@ -6,9 +6,10 @@ use async_trait::async_trait; use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; use datafusion::catalog::Session; use datafusion::datasource::physical_plan::FileSinkConfig; +use datafusion::execution::SessionStateBuilder; +use datafusion::execution::context::SessionConfig; use datafusion::execution::object_store::ObjectStoreUrl; use datafusion::execution::runtime_env::RuntimeEnv; -use datafusion::execution::{SessionStateBuilder, context::SessionConfig}; use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder, TableSource}; use datafusion::physical_expr::LexRequirement; use datafusion::physical_expr_common::sort_expr::LexOrdering; From 6934b88575e3662bcfa077d89d9571b8dc9cbe84 Mon Sep 17 00:00:00 2001 From: XL Liang Date: Fri, 24 Jul 2026 16:32:46 +0800 Subject: [PATCH 9/9] df --- crates/sail-common/src/config/application.rs | 1 + .../sail-common/src/config/application.yaml | 11 +++ crates/sail-data-source/src/listing/source.rs | 67 +++++++++++++++---- crates/sail-session/src/formats.rs | 39 ++++++++--- .../src/session_factory/server.rs | 6 +- 5 files changed, 100 insertions(+), 24 deletions(-) diff --git a/crates/sail-common/src/config/application.rs b/crates/sail-common/src/config/application.rs index 7e25629a8b..66e6910994 100644 --- a/crates/sail-common/src/config/application.rs +++ b/crates/sail-common/src/config/application.rs @@ -290,6 +290,7 @@ pub struct ExecutionConfig { pub batch_size: usize, pub default_parallelism: usize, pub collect_statistics: bool, + pub prewarm_file_statistics_on_source_creation: bool, pub use_row_number_estimates_to_optimize_partitioning: bool, pub file_listing_cache: FileListingCacheConfig, } diff --git a/crates/sail-common/src/config/application.yaml b/crates/sail-common/src/config/application.yaml index cd59a557f3..7288ac9526 100644 --- a/crates/sail-common/src/config/application.yaml +++ b/crates/sail-common/src/config/application.yaml @@ -275,6 +275,17 @@ Has no effect after the table is created. experimental: true +- key: execution.prewarm_file_statistics_on_source_creation + type: boolean + default: "false" + description: | + Whether to eagerly populate the file statistics cache when resolving a listing data source. + This applies to direct reads such as `spark.read.parquet()` as well as catalog table reads. + It can make data source resolution more expensive but avoids collecting file statistics + during the first query. This setting requires `execution.collect_statistics` and an enabled + file statistics cache. + experimental: true + - key: execution.use_row_number_estimates_to_optimize_partitioning type: boolean default: "false" diff --git a/crates/sail-data-source/src/listing/source.rs b/crates/sail-data-source/src/listing/source.rs index f3af011781..21135a61d2 100644 --- a/crates/sail-data-source/src/listing/source.rs +++ b/crates/sail-data-source/src/listing/source.rs @@ -139,12 +139,20 @@ pub trait WriteFormat: Debug + Send + Sync + 'static { ) -> Result>; } -#[derive(Debug, Default)] +#[derive(Debug)] pub struct ListingTableFormat { + prewarm_file_statistics_on_source_creation: bool, phantom: PhantomData, } impl ListingTableFormat { + pub fn new(prewarm_file_statistics_on_source_creation: bool) -> Self { + Self { + prewarm_file_statistics_on_source_creation, + phantom: PhantomData, + } + } + async fn create_listing_source( &self, ctx: &dyn Session, @@ -244,6 +252,25 @@ impl ListingTableFormat { compression, }) } + + async fn prewarm_listing_source( + &self, + source: &ListingTableSource, + ctx: &dyn Session, + ) -> Result<()> { + prewarm_file_statistics(source, ctx).await?; + log::debug!( + "prewarmed listing file statistics cache for {} path(s)", + source.config().table_paths.len() + ); + Ok(()) + } +} + +impl Default for ListingTableFormat { + fn default() -> Self { + Self::new(false) + } } #[async_trait] @@ -257,7 +284,19 @@ impl TableFormat for ListingTableFormat { ctx: &dyn Session, info: SourceInfo, ) -> Result> { - Ok(Arc::new(self.create_listing_source(ctx, info).await?)) + let source = self.create_listing_source(ctx, info).await?; + if self.prewarm_file_statistics_on_source_creation + && ctx.config().collect_statistics() + && ctx + .runtime_env() + .cache_manager + .get_file_statistic_cache_limit() + > 0 + && let Err(error) = self.prewarm_listing_source(&source, ctx).await + { + log::warn!("failed to prewarm listing file statistics: {error}"); + } + Ok(Arc::new(source)) } async fn prewarm_statistics( @@ -277,12 +316,7 @@ impl TableFormat for ListingTableFormat { .with_runtime_env(runtime_env) .build(); let source = self.create_listing_source(&session_state, info).await?; - prewarm_file_statistics(&source, &session_state).await?; - log::debug!( - "prewarmed listing file statistics cache for {} path(s)", - source.config().table_paths.len() - ); - Ok(()) + self.prewarm_listing_source(&source, &session_state).await } async fn create_writer(&self, ctx: &dyn Session, info: SinkInfo) -> Result { @@ -488,7 +522,7 @@ mod tests { } #[tokio::test] - async fn create_source_does_not_prewarm_file_statistics() { + async fn create_source_prewarms_file_statistics_when_enabled() { FILE_META_INFERENCE_COUNT.store(0, Ordering::Relaxed); let object_store: Arc = Arc::new(InMemory::new()); object_store @@ -501,7 +535,6 @@ mod tests { let context = SessionContext::new(); context.register_object_store(&Url::parse("memory://").unwrap(), object_store); let state = context.state(); - let format = ListingTableFormat::::default(); let info = SourceInfo { paths: vec!["memory:///table/".to_string()], lakehouse_table: None, @@ -518,13 +551,21 @@ mod tests { read_case_sensitive: true, }; - format.create_source(&state, info.clone()).await.unwrap(); + let disabled_format = ListingTableFormat::::new(false); + disabled_format + .create_source(&state, info.clone()) + .await + .unwrap(); assert_eq!(FILE_META_INFERENCE_COUNT.load(Ordering::Relaxed), 0); - format - .prewarm_statistics(state.config().clone(), state.runtime_env().clone(), info) + let enabled_format = ListingTableFormat::::new(true); + enabled_format + .create_source(&state, info.clone()) .await .unwrap(); assert_eq!(FILE_META_INFERENCE_COUNT.load(Ordering::Relaxed), 1); + + enabled_format.create_source(&state, info).await.unwrap(); + assert_eq!(FILE_META_INFERENCE_COUNT.load(Ordering::Relaxed), 1); } } diff --git a/crates/sail-session/src/formats.rs b/crates/sail-session/src/formats.rs index daa9bd558f..3a13ca780c 100644 --- a/crates/sail-session/src/formats.rs +++ b/crates/sail-session/src/formats.rs @@ -17,21 +17,40 @@ use sail_data_source::formats::text::TextTableFormat; use sail_delta_lake::DeltaTableFormat; use sail_iceberg::IcebergTableFormat; -pub fn create_table_format_registry() -> Result> { +pub fn create_table_format_registry( + prewarm_file_statistics_on_source_creation: bool, +) -> Result> { let registry = Arc::new(TableFormatRegistry::new()); - register_builtin_formats(®istry)?; + register_builtin_formats(®istry, prewarm_file_statistics_on_source_creation)?; register_external_formats(®istry)?; Ok(registry) } -fn register_builtin_formats(registry: &Arc) -> Result<()> { - registry.register(Arc::new(ArrowTableFormat::default()))?; - registry.register(Arc::new(AvroTableFormat::default()))?; - registry.register(Arc::new(BinaryTableFormat::default()))?; - registry.register(Arc::new(CsvTableFormat::default()))?; - registry.register(Arc::new(JsonTableFormat::default()))?; - registry.register(Arc::new(ParquetTableFormat::default()))?; - registry.register(Arc::new(TextTableFormat::default()))?; +fn register_builtin_formats( + registry: &Arc, + prewarm_file_statistics_on_source_creation: bool, +) -> Result<()> { + registry.register(Arc::new(ArrowTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(AvroTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(BinaryTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(CsvTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(JsonTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(ParquetTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; + registry.register(Arc::new(TextTableFormat::new( + prewarm_file_statistics_on_source_creation, + )))?; registry.register(Arc::new(SocketTableFormat))?; registry.register(Arc::new(RateTableFormat))?; registry.register(Arc::new(ConsoleTableFormat))?; diff --git a/crates/sail-session/src/session_factory/server.rs b/crates/sail-session/src/session_factory/server.rs index 3cc0d453a0..2893c7b85e 100644 --- a/crates/sail-session/src/session_factory/server.rs +++ b/crates/sail-session/src/session_factory/server.rs @@ -108,7 +108,11 @@ impl ServerSessionFactory { // We do not use the DataFusion catalog and schema since we manage catalogs ourselves. .with_create_default_catalog_and_schema(false) .with_information_schema(false) - .with_extension(create_table_format_registry()?) + .with_extension(create_table_format_registry( + self.config + .execution + .prewarm_file_statistics_on_source_creation, + )?) .with_extension(Arc::new(create_catalog_manager( &self.config, self.runtime.clone(),