sunchao commented on code in PR #6500: URL: https://github.com/apache/datafusion-comet/pull/6500#discussion_r4162073498
########## 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: [P2] Make the sort workaround cover changing live pool shares. `memory_limit / sorters` is not an upper bound on a sorter's buffered reservation as consumers register and finish. A sorter can therefore exceed this threshold and enter DataFusion 55.1's unreserved merge path. With four partitions, two Int64 sort keys, 1,048,576 rows, batch size 4096 and an 8 MiB budget, execution fails with `ResourcesExhausted` in `ExternalSorterMerge` despite spilling being enabled. It should spill and return every row in order. This makes admitted local sorts fail under memory pressure. Please use a workaround that covers the possible live reservation sizes, or backport the upstream merge-reservation fix. Evidence: Exact-head CI job https://github.com/apache/datafusion-comet/actions/runs/36835871816/job/110284444406 fails the PR's `multi_column_sorts_spill_under_a_shared_budget` test at planner.rs:1003: allocating 456 B for `ExternalSorterMerge[0]` fails with zero pool space available. Independently reproduced using the locked dependencies and `UnionExec -> partition-preserving SortExec -> SortPreservingMergeExec`: rows have keys `(id * 37 % 1000003, id)`, distributed in 4096-row batches across four partitions. The current 512 KiB spill reservation and 2 MiB in-place threshold fail on iteration zero while requesting another 64 KiB. Keeping the budget and input identical but raising only the in-place threshold to 8 MiB passes ten iterations, checking row count and ordering. Disposable reproduction retained at `/tmp/comet6500-review-current-sort.rs`, with results at `/tmp/comet6500-review-current-sort-results.log`. -- 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]
