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]

Reply via email to