viirya commented on code in PR #6500:
URL: https://github.com/apache/datafusion-comet/pull/6500#discussion_r4163437227


##########
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:
   Thanks, you're right. `budget / sorters` is not an upper bound: 
`FairSpillPool` divides the pool by the spillable consumers registered at the 
moment, so a sorter's share grows while other consumers have not started or 
after they finish, and its buffered batches can then exceed the threshold and 
take the unreserved merge path.
   
   Fixed in 70c3a310d by setting `sort_in_place_threshold_bytes` to 
`usize::MAX`, so sorters always sort their buffered batches in place and never 
reach that merge. The cost is that in-memory multi-column sorts concatenate a 
partition's buffered batches and sort them with the lexicographic comparator. 
The upstream fix is in DataFusion 56.0.0, and Comet depends on the released 
55.1 crates, so the override stays until that upgrade; the code comment and 
docs say so.
   
   To cover changing shares, `multi_column_sorts_spill_under_a_shared_budget` 
now uses uneven partitions, so sorters that finish early enlarge the others' 
shares. It fails 6 of 10 runs with the previous share-based threshold and 10 of 
10 without the override, and passes 20 of 20 with the fix. 105afe87a fixes the 
Scalafix failure.
   



-- 
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]

Reply via email to