viirya commented on code in PR #6500: URL: https://github.com/apache/datafusion-comet/pull/6500#discussion_r4175710225
########## native/core/src/local/planner.rs: ########## @@ -0,0 +1,1049 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Local graph construction. Reuse Comet scan/expression builders without per-task plan execution. + +use std::sync::Arc; + +use arrow::compute::SortOptions; +use datafusion::common::{JoinType, NullEquality}; +use datafusion::execution::disk_manager::{DiskManagerBuilder, DiskManagerMode}; +use datafusion::execution::memory_pool::FairSpillPool; +use datafusion::execution::runtime_env::RuntimeEnvBuilder; +use datafusion::physical_expr::{LexOrdering, PhysicalSortExpr}; +use datafusion::physical_plan::aggregates::{AggregateExec, AggregateMode, PhysicalGroupBy}; +use datafusion::physical_plan::coalesce_partitions::CoalescePartitionsExec; +use datafusion::physical_plan::joins::{HashJoinExec, PartitionMode}; +use datafusion::physical_plan::limit::GlobalLimitExec; +use datafusion::physical_plan::repartition::RepartitionExec; +use datafusion::physical_plan::sorts::sort::SortExec; +use datafusion::physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec; +use datafusion::physical_plan::{ + filter::FilterExec, projection::ProjectionExec, union::UnionExec, ExecutionPlan, +}; +use datafusion::physical_plan::{ExecutionPlanProperties, Partitioning}; +use datafusion::prelude::{SessionConfig, SessionContext}; +use datafusion_comet_local::LocalQuery; +use datafusion_comet_proto::local::{LocalAggregate, LocalJoin, LocalOutput}; +use datafusion_comet_proto::spark_operator::{operator::OpStruct, Operator, SparkFilePartition}; +use prost::Message; + +use crate::execution::operators::ExecutionError; +use crate::execution::planner::PhysicalPlanner; +use crate::parquet::parquet_support::CometObjectStoreRegistry; + +pub(super) struct QuerySettings<'a> { + pub terminal: &'a [u8], + pub aggregate: &'a [u8], + pub memory_limit: usize, + pub spill_enabled: bool, +} + +pub(super) fn parquet_query( + bytes: &[u8], + partitions: &[Vec<u8>], + batch_size: usize, + columns: usize, + row_filter_pushdown: bool, + settings: QuerySettings<'_>, +) -> Result<LocalQuery, ExecutionError> { + let root = Operator::decode(bytes)?; + let groups = partitions + .iter() + .map(|b| SparkFilePartition::decode(b.as_slice())) + .collect::<Result<Vec<_>, _>>()?; + let context = query_context(batch_size, groups.len(), row_filter_pushdown, &settings)?; + let planner = PhysicalPlanner::new(Arc::clone(&context), 0).with_sql_text_pool(&root); + let plan = build(&root, &groups, &planner)?; + let plan = if settings.aggregate.is_empty() { + plan + } else { + aggregate_plan(plan, &LocalAggregate::decode(settings.aggregate)?, &planner)? + }; + let plan = output_plan(plan, settings.terminal, &planner)?; + if plan.schema().fields().len() != columns { + return Err(ExecutionError::GeneralError( + "Local output schema width mismatch".into(), + )); + } + size_sorters(&context, plan.as_ref(), settings.memory_limit); + Ok(LocalQuery::new(plan, context.task_ctx())) +} + +fn query_context( + batch_size: usize, + partitions: usize, + row_filter_pushdown: bool, + settings: &QuerySettings<'_>, +) -> Result<Arc<SessionContext>, ExecutionError> { + let mut config = SessionConfig::new() + .with_batch_size(batch_size) + .with_target_partitions(partitions.max(1)); + config.options_mut().execution.parquet.pushdown_filters = row_filter_pushdown; + config.options_mut().execution.parquet.reorder_filters = row_filter_pushdown; + // Registry and configuration are query-owned. Never inherit another query's credentials. + let runtime = RuntimeEnvBuilder::new() + .with_memory_pool(Arc::new(FairSpillPool::new(settings.memory_limit))) + .with_disk_manager_builder(DiskManagerBuilder::default().with_mode( + if settings.spill_enabled { + DiskManagerMode::OsTmpDirectory + } else { + DiskManagerMode::Disabled + }, + )) + .with_object_store_registry(Arc::new(CometObjectStoreRegistry::default())) + .build()?; + let context = Arc::new(SessionContext::new_with_config_rt( + config, + Arc::new(runtime), + )); + Ok(context) +} + +/// Sizes DataFusion's sort settings by the number of sorters that share the query budget. +/// Must run after the whole graph is built and before its task context is created. +fn size_sorters(context: &SessionContext, plan: &dyn ExecutionPlan, memory_limit: usize) { + fn sorters(plan: &dyn ExecutionPlan) -> usize { + let own = match plan.downcast_ref::<SortExec>() { + Some(sort) if sort.preserve_partitioning() => { + sort.input().output_partitioning().partition_count() + } + Some(_) => 1, + None => 0, + }; + own + plan + .children() + .into_iter() + .map(|child| sorters(child.as_ref())) + .sum::<usize>() + } + let sorters = sorters(plan); + if sorters == 0 { + return; + } + let share = memory_limit / sorters; + let state = context.state_ref(); + let mut state = state.write(); + let execution = &mut state.config_mut().options_mut().execution; + // Every sorter reserves this much for its final merge before sorting. With the default + // 10 MiB, a few concurrent sorters can exhaust a small query budget up front. + execution.sort_spill_reservation_bytes = + (share / 4).min(execution.sort_spill_reservation_bytes); + // Workaround for DataFusion 55.1's ExternalSorter, fixed upstream in DataFusion 56.0.0: + // before spilling, it frees its merge reservation and merges buffered batches with a new, + // empty, unspillable reservation. Once spillable sorters fill the fair pool, that merge + // cannot grow and the query fails instead of spilling. A sorter spills once its buffered + // batches reach its fair share, so a threshold of one share makes it sort them in place + // instead of merging. This costs unaccounted transient copies and slower multi-column + // sorts that fit in memory. Remove this override after upgrading to DataFusion 56.0.0; + // `multi_column_sorts_spill_under_a_shared_budget` must still pass without it. + execution.sort_in_place_threshold_bytes = share.max(execution.sort_in_place_threshold_bytes); Review Comment: You're right, and my earlier reply was wrong to call this fixed. `sort_and_spill_in_mem_batches` frees `merge_reservation` before it sorts and spills and only calls `reserve_memory_for_merge` once the spill finishes, so the `try_resize` back to `sort_spill_reservation_bytes` races with the other sorters' spillable growth. The in-place threshold never reaches that path, and `usize::MAX` adds the `concat_batches` offset overflow you describe. I'll drop global sort admission together with both sort setting overrides and `multi_column_sorts_spill_under_a_shared_budget`, and keep Top-K, which goes through `TopK` rather than `ExternalSorter`. Global sort can be admitted again after the DataFusion 56 upgrade (#6410), with a regression test for concurrent spilling sorts. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
