andygrove commented on code in PR #6500: URL: https://github.com/apache/datafusion-comet/pull/6500#discussion_r4169251319
########## spark-local/src/main/scala/org/apache/comet/local/CometLocalRule.scala: ########## @@ -0,0 +1,130 @@ +/* + * 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. + */ + +package org.apache.comet.local + +import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.catalyst.expressions.{Alias, AttributeReference, NamedExpression} +import org.apache.spark.sql.catalyst.rules.Rule +import org.apache.spark.sql.execution.{CollectLimitExec, ProjectExec, QueryExecution, RangeExec, SparkPlan} + +import org.apache.comet.CometConf +import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, withInfo} +import org.apache.comet.local.shims.LocalModeSupport +import org.apache.comet.shims.ShimCometStreaming + +private[comet] case class CometLocalRule(session: SparkSession) extends Rule[SparkPlan] { + override def apply(plan: SparkPlan): SparkPlan = { + if (!CometConf.COMET_EXEC_LOCAL_ENABLED.get(conf) || CometLocalRule.isLocal(plan)) { + return plan + } + val reason = LocalModeSupport + .environmentRejection(session.sparkContext.isLocal, conf.adaptiveExecutionEnabled) Review Comment: `environmentRejection` reads the session's AQE setting, so enabling local execution means turning AQE off for every query. That's what makes unadmitted TPC-H about five times slower on the benchmark page. It also declines queries AQE never touches. `InsertAdaptiveSparkPlan` only wraps plans with an exchange, a required distribution or a subquery, so a plain scan/filter/project runs without AQE even when it's enabled, yet admission still rejects it. Could admission run from the `CometRule(session, queryStagePrep = true)` instance that `CometSparkSessionExtensions` registers through `injectQueryStagePrepRule`? Under AQE it sees the whole initial plan, and Spark runs custom prep rules after `EnsureRequirements`, so the same exchanges would be there for these planners to match. Replacing the plan with a leaf leaves AQE no stages to create, and exchange-free plans still reach the columnar instance. That would drop the session-wide AQE requirement entirely. ########## spark-local/src/main/scala/org/apache/comet/local/CometLocalResultExec.scala: ########## @@ -0,0 +1,122 @@ +/* + * 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. + */ + +package org.apache.comet.local + +import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.atomic.AtomicLong + +import scala.collection.mutable.ArrayBuffer + +import org.apache.spark.SparkException +import org.apache.spark.rdd.RDD +import org.apache.spark.sql.catalyst.InternalRow +import org.apache.spark.sql.catalyst.expressions.{Attribute, SortOrder, UnsafeRow} +import org.apache.spark.sql.catalyst.plans.physical.Partitioning +import org.apache.spark.sql.execution.{SparkPlan, UnaryExecNode} + +/** + * Row result boundary of a local query. Spark's collect encodes and compresses every row of a + * result partition in its task, then decodes them on the driver. A local query has one result + * partition, so that work would run on a single thread. Admission requires an in-process local + * master, so collect and take still run one Spark task (keeping cancellation, job groups and SQL + * metrics) but hand copied rows to the driver in this JVM. + */ +case class CometLocalResultExec(child: SparkPlan) extends UnaryExecNode { + override def output: Seq[Attribute] = child.output + override def outputPartitioning: Partitioning = child.outputPartitioning Review Comment: What happens with `df.cache()` or `checkpoint()` on an admitted query? `CacheManager` prepares the cached plan through `sessionState.executePlan(...).executedPlan`, so it goes through admission and comes out as a `CometLocalResultExec` with `SinglePartition`. `InMemoryTableScanExec` reports `cachedPlan.outputPartitioning`, and `SinglePartition` satisfies every `ClusteredDistribution`, so `EnsureRequirements` won't add a shuffle above it. A later `groupBy(...).agg(sum(...))` over the cache, which isn't admitted because of `SUM`, would then run in a single task. Could you add a test that caches an admitted query and checks the partition count of a query over the cache? If it collapses to one partition, it may be worth declining admission for plans prepared for the cache. ########## native/core/src/local/planner.rs: ########## @@ -0,0 +1,1056 @@ +// 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 consumers fill the fair pool, that merge + // cannot grow and the query fails instead of spilling. A sorter's fair share changes as + // other consumers register and finish, so no smaller threshold bounds its buffered + // batches; always sorting them in place avoids that merge. 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 = usize::MAX; +} + +pub(super) fn join_query( + bytes: &[u8], + batch_size: usize, + columns: usize, + row_filter_pushdown: bool, + settings: QuerySettings<'_>, +) -> Result<LocalQuery, ExecutionError> { + let join = LocalJoin::decode(bytes)?; + if join.left_files.len() > 1024 || join.right_files.len() > 1024 { + return Err(ExecutionError::GeneralError( + "Too many local join file groups".into(), + )); + } + let context = query_context( + batch_size, + join.partitions as usize, + row_filter_pushdown, + &settings, + )?; + let invalid = || ExecutionError::GeneralError("Missing local join input".into()); + let left = join.left.as_ref().ok_or_else(invalid)?; + let right = join.right.as_ref().ok_or_else(invalid)?; + // Each input retains its own SQL text pool; neither borrows the other's scan metadata. + let left_planner = PhysicalPlanner::new(Arc::clone(&context), 0).with_sql_text_pool(left); + let right_planner = PhysicalPlanner::new(Arc::clone(&context), 0).with_sql_text_pool(right); + let left = build(left, &join.left_files, &left_planner)?; + let right = build(right, &join.right_files, &right_planner)?; + let planner = PhysicalPlanner::new(Arc::clone(&context), 0); + let plan = join_plan(left, right, &join, &planner)?; + let plan = output_plan(plan, settings.terminal, &planner)?; + if plan.schema().fields().len() != columns { + return Err(ExecutionError::GeneralError( + "Local join output schema width mismatch".into(), + )); + } + size_sorters(&context, plan.as_ref(), settings.memory_limit); + Ok(LocalQuery::new(plan, context.task_ctx())) +} + +fn join_plan( + left: Arc<dyn ExecutionPlan>, + right: Arc<dyn ExecutionPlan>, + join: &LocalJoin, + planner: &PhysicalPlanner, +) -> Result<Arc<dyn ExecutionPlan>, ExecutionError> { + use datafusion_comet_proto::spark_operator::JoinType as SparkJoinType; + if !(1..=1024).contains(&join.partitions) + || join.left_keys.is_empty() + || join.left_keys.len() != join.right_keys.len() + || join.result.is_empty() + { + return Err(ExecutionError::GeneralError("Invalid local join".into())); + } + let kind = match SparkJoinType::try_from(join.join_type) { + Ok(SparkJoinType::Inner) => JoinType::Inner, + Ok(SparkJoinType::LeftOuter) => JoinType::Left, + Ok(SparkJoinType::RightOuter) => JoinType::Right, + Ok(SparkJoinType::FullOuter) => JoinType::Full, + Ok(SparkJoinType::LeftSemi) => JoinType::LeftSemi, + Ok(SparkJoinType::LeftAnti) => JoinType::LeftAnti, + Err(_) => { + return Err(ExecutionError::GeneralError( + "Invalid local join type".into(), + )) + } + }; + let left_keys = join + .left_keys + .iter() + .map(|e| planner.create_expr(e, left.schema())) + .collect::<Result<Vec<_>, _>>()?; + let right_keys = join + .right_keys + .iter() + .map(|e| planner.create_expr(e, right.schema())) + .collect::<Result<Vec<_>, _>>()?; + // Both exchanges use the same DataFusion hash implementation and partition count. + // Spark's partition IDs and hash algorithm never cross this boundary. + let left = Arc::new(RepartitionExec::try_new( + left, + Partitioning::Hash(left_keys.clone(), join.partitions as usize), + )?); + let right = Arc::new(RepartitionExec::try_new( + right, + Partitioning::Hash(right_keys.clone(), join.partitions as usize), + )?); + let on = left_keys.into_iter().zip(right_keys).collect(); + let hash = HashJoinExec::try_new( + left, + right, + on, + None, + &kind, + None, + PartitionMode::Partitioned, + NullEquality::NullEqualsNothing, + false, + )?; + // swap_inputs restores Spark's logical output order with a projection when needed. + let plan: Arc<dyn ExecutionPlan> = if join.build_right { + hash.swap_inputs(PartitionMode::Partitioned)? + } else { + Arc::new(hash) + }; + let result = join + .result + .iter() + .enumerate() + .map(|(i, e)| { + planner + .create_expr(e, plan.schema()) + .map(|expr| (expr, format!("col_{i}"))) + }) + .collect::<Result<Vec<_>, _>>()?; + Ok(Arc::new(ProjectionExec::try_new(result, plan)?)) +} + +fn output_plan( + input: Arc<dyn ExecutionPlan>, + bytes: &[u8], + planner: &PhysicalPlanner, +) -> Result<Arc<dyn ExecutionPlan>, ExecutionError> { + if bytes.is_empty() { + return Ok(input); + } + let output = LocalOutput::decode(bytes)?; + let skip = usize::try_from(output.skip) + .map_err(|_| ExecutionError::GeneralError("Local offset overflow".into()))?; + let fetch = output + .fetch + .map(usize::try_from) + .transpose() + .map_err(|_| ExecutionError::GeneralError("Local fetch overflow".into()))?; + let top = fetch + .map(|n| { + n.checked_add(skip) + .ok_or_else(|| ExecutionError::GeneralError("Local Top-K overflow".into())) + }) + .transpose()?; + let expressions = output + .orders + .iter() + .map(|order| { + let child = order + .child + .as_ref() + .ok_or_else(|| ExecutionError::GeneralError("Missing local sort key".into()))?; + Ok(PhysicalSortExpr { + expr: planner.create_expr(child, input.schema())?, + options: SortOptions { + descending: order.descending, + nulls_first: order.nulls_first, + }, + }) + }) + .collect::<Result<Vec<_>, ExecutionError>>()?; + // Sort each input partition in parallel, then merge the sorted runs into one ordered + // partition. Never send independently sorted partitions through the unordered result + // coalescer: that would lose the global order. An unordered limit just gathers. + let mut plan: Arc<dyn ExecutionPlan> = match LexOrdering::new(expressions) { + Some(ordering) if input.output_partitioning().partition_count() > 1 => { + let sorted = Arc::new( + SortExec::new(ordering.clone(), input) + .with_preserve_partitioning(true) + .with_fetch(top), + ); + Arc::new(SortPreservingMergeExec::new(ordering, sorted).with_fetch(top)) + } + Some(ordering) => Arc::new(SortExec::new(ordering, input).with_fetch(top)), + None => Arc::new(CoalescePartitionsExec::new(input)), Review Comment: For an unordered `LIMIT` or `OFFSET` this gathers through `CoalescePartitionsExec`, so which rows survive depends on which partitions produce first. Spark's `CollectLimitExec` and Comet's `CometCollectLimitExec` both go through `executeTake`, which reads partitions in order, so `show()`, `head()` and `offset(n)` return the same rows on every run. Here they'd return a different subset from run to run. SQL allows that, but people will notice it in exactly the notebook setting this targets. The test at `CometLocalExecutionSuite.scala:861` only checks the row count. Is the divergence intended? If so, could `local-execution.md` say so? If not, could unordered limits stay on the ordinary path or gather partitions in order? ########## native/core/src/local/planner.rs: ########## @@ -0,0 +1,1056 @@ +// 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 Review Comment: Local queries spill into the OS temp directory, while the per-task path spills into Spark's local dirs (`SparkEnv.get.blockManager.getLocalDiskDirs` in `CometExecIterator.scala:105`) and caps them with `spark.comet.maxTempDirectorySize` (`prepare_datafusion_session_context` in `jni_api.rs`). Anyone who points `spark.local.dir` at a large scratch disk would find local queries spilling into `/tmp` instead, under DataFusion's default cap rather than Comet's. Where `/tmp` is tmpfs, a spill lands in RAM, which defeats the purpose under memory pressure. `LocalQueryIterator` runs inside the result task, so `SparkEnv` is available there. Could it pass the local dirs and the max size through `createParquet` and `createJoin`? ########## spark-local/src/main/scala/org/apache/comet/local/LocalAggregatePlanner.scala: ########## @@ -0,0 +1,120 @@ +/* + * 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. + */ + +package org.apache.comet.local + +import scala.jdk.CollectionConverters._ + +import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.catalyst.expressions.aggregate.{Count, Final, Max, Min, Partial} +import org.apache.spark.sql.catalyst.plans.physical.{HashPartitioning, SinglePartition} +import org.apache.spark.sql.execution.SparkPlan +import org.apache.spark.sql.execution.aggregate.HashAggregateExec +import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec +import org.apache.spark.sql.types.{DoubleType, FloatType} + +import org.apache.comet.CometConf +import org.apache.comet.serde.LocalOuterClass.LocalAggregate +import org.apache.comet.serde.QueryPlanSerde.{aggExprToProto, exprToProto} + +/** Recognize an entire Spark aggregation, then discard its task/buffer/exchange boundaries. */ +private[local] object LocalAggregatePlanner { + import LocalParquetPlanner.{nativeOnly, supportedExpression} + + def plan(root: SparkPlan, session: SparkSession): Option[LocalParquetSpec] = root match { + case last: HashAggregateExec + if CometConf.COMET_EXEC_AGGREGATE_ENABLED.get(last.conf) && + last.aggregateExpressions.nonEmpty && + last.aggregateExpressions.forall(a => a.mode == Final && !a.isDistinct) => + val input = last.child match { + case exchange: ShuffleExchangeExec => + exchange.child match { + case first: HashAggregateExec => + exchange.outputPartitioning match { + case hash: HashPartitioning + if first.groupingExpressions.nonEmpty && + hash.expressions.size == first.groupingExpressions.size && + hash.expressions.zip(first.groupingExpressions).forall { + case (key, group) => + key.semanticEquals(group.toAttribute) + } => + Some((first, hash.numPartitions)) Review Comment: The native partition count comes from Spark's exchange, so it follows `spark.sql.shuffle.partitions`. The benchmark sets that to 8, but with AQE off most users will be at the default of 200. In DataFusion 55.1 every `RepartitionExec` output partition and every grouped aggregate stream registers its own spillable consumer. A grouped aggregate at the default would split the 256 MiB `FairSpillPool` across roughly 400 consumers, about 640 KiB each, and spill early. Could you add a default-partitions run to the benchmark page? Since the exchange never leaves the graph, would it make sense to pick the native partition count independently, for example from `defaultParallelism`? ########## spark-local/src/main/scala/org/apache/comet/local/shims/LocalModeSupport.scala: ########## @@ -0,0 +1,31 @@ +/* + * 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. + */ + +package org.apache.comet.local.shims + +/** Admission only: the experimental bridge is validated against the Spark 4.1 line. */ +private[local] object LocalModeSupport { + def supported: Boolean = org.apache.spark.SPARK_VERSION.startsWith("4.1.") Review Comment: This gates on a version string in shared code, in a `shims` package that isn't one of the per-version shim source sets. What's 4.1-specific here? Everything compiles on every profile, and once 4.2 becomes supported (#6417) the mode would quietly switch off there and the whole suite would be cancelled. If the gate is needed, could it live in the `spark/src/main/spark-*/org/apache/comet/shims/` sources so each new Spark version makes an explicit choice? If it isn't, could the suite run on 4.0 and 4.2 as well? ########## docs/source/contributor-guide/local-execution.md: ########## @@ -0,0 +1,258 @@ +<!-- +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 Execution + +Local execution is an experimental mode in which a whole admitted query runs as one +DataFusion graph inside the Spark driver process. Spark still parses, analyzes and +optimizes the query, and Comet's Spark-compatible expressions and Parquet scan are +reused, but DataFusion schedules all partitions and exchanges itself. An admitted +query has no Spark shuffle and runs as a single Spark result task. + +`spark.comet.exec.local.enabled=true` opts in through the existing Comet extension. +The option is internal and defaults to false. Queries that are not admitted as a +whole use ordinary Comet/Spark planning. See [Local Execution Benchmark](local-execution-benchmark.md) +for the manual benchmark and current results. + +## Execution contract + +Each query execution owns one fresh physical graph. Every partition uses those same +operator instances, including the channels and state of `RepartitionExec`. Plans +deserialized separately for Spark tasks cannot implement this exchange, so sharing +is a consequence of query ownership, not a cache. Repeating an action on the same +DataFrame, or retrying the result task, creates a new graph; no stateful operator +tree is reused across executions. + +The native API accepts a fully planned graph and a query-scoped DataFusion +`TaskContext`. It does not optimize the graph or verify distribution requirements; +the local planner must satisfy them. Execution consumes the query once and returns +one result stream. DataFusion's `execute_stream` drives all root partitions +concurrently, so progress does not depend on Spark scheduling partitions. + +The result stream is unordered when the root has several partitions. Global +ordering must therefore be established by the planner with a single ordered root; +coalescing partitions never preserves order. Spark sees one output partition, and +the local node reports its ordering conservatively. + +EOF, the first stream error, explicit cancellation and dropping the stream release +the owner's references to the graph and context. DataFusion tasks abort +asynchronously, so dropping is not a cleanup barrier; the process-owned Tokio +runtime must stay alive to finish teardown. + +## Module boundaries + +- `native/local`: query execution lifecycle and result handoff, independent of JNI + and Spark tasks. It reuses DataFusion scheduling, exchange and streaming rather + than implementing another scheduler. +- `native/core/src/local.rs` and `native/core/src/local/planner.rs`: the JNI entry + points and the local planner. The planner reuses core's `PhysicalPlanner` for + scans and expressions and owns distribution, ordering and exchange translation + for the complete graph. Core may depend on `native/local`; not the reverse. +- `spark-local`: JVM mode selection, whole-query admission, query ownership, the + result bridge, cancellation and SQL metrics. It is compiled into the existing + Spark artifact, avoiding a cyclic dependency on `CometConf` and `NativeUtil`. +- `native/proto/src/proto/local.proto`: aggregation, join and terminal sort/limit + descriptions that have no counterpart in Comet's per-task operator protocol. + +Task-bound memory managers, JVM input iterators, per-task plan creation and Spark +shuffle readers/writers are not used by local execution. + +## Admission + +Admission inspects the complete physical plan before execution starts. Fallback is +a planning decision only: there is no fallback after native output has started. + +Environment requirements: + +- Spark 4.1. Local execution stays disabled on other Spark profiles. +- `SparkContext.isLocal`, as reported by the application's effective SparkContext, + not a session override of `spark.master`. `local-cluster` and remote masters + are rejected. +- AQE disabled, Comet and Comet native execution enabled, not in plan-only mode. + AQE is a session setting, so queries that are not admitted also run without AQE; + see [Local Execution Benchmark](local-execution-benchmark.md#queries-outside-admission) + for what that costs on TPC-H. +- Batch queries only. Streaming plans and subquery preparation are rejected. + +Admitted query shapes: + +- `RangeExec` with optional direct column or alias projections. +- DataSource V1 scans of Spark's built-in Parquet format on `file:` paths, with + filter and projection. `CometScanRule` validates a copy of the scan, and Spark's + static partition pruning and file splitting are reused. Encrypted reads, custom + or cloud filesystems, object-store options, bucketed or ordered scans and file + metadata columns are rejected, so credentials and encryption callbacks never + cross the local boundary. +- Grouped or global `COUNT`, `MIN` and `MAX` over an admitted scan, recognized from + Spark's partial/final hash aggregate pair and its exchange. +- A single `ShuffledHashJoinExec` over two admitted scans (for example selected + with a `SHUFFLE_HASH` hint): inner, left/right/full outer, left semi and left + anti, with equal-type, non-floating-point attribute keys and no residual + condition. Sort-merge and broadcast joins are not converted. +- A terminal global `SortExec`, `TakeOrderedAndProjectExec` or root `CollectLimitExec` + over any of the Parquet, aggregate or join shapes. A global sort's range exchange + is removed only when every sort expression, direction and null placement matches. + A range query keeps Spark's root `CollectLimitExec` instead. + +Expressions are limited to an explicit allowlist: references, literals, aliases, +arithmetic, comparisons, boolean and null predicates, casts, overflow checks and +conditionals. Both the Spark expression and its serialized form must be admitted, +because serde can choose JVM codegen callbacks for familiar expression classes. +Supported types are primitive numeric, boolean, string, binary, decimal, date, +timestamp and timestamp NTZ. Floating-point grouping, join and sort keys, nested +and collated types, UDFs, nondeterministic or partition-sensitive expressions, +subqueries, `DISTINCT`, `SUM`/`AVG`, `HAVING`, nested joins and aggregates around +joins fall back as a whole. Existing Comet operator enablement flags are honored. + +Limits are 1,024 file groups or native partitions, 1,024 output columns and batch +size 65,536. These bound execution overhead, not total memory. + +## Native plan shapes + +- Scans: one native Parquet scan per Spark file partition, combined with a union, + with shared filter and projection operators above it. +- Grouped aggregation: a partial aggregate per input partition, a DataFusion hash + `RepartitionExec` into the Spark shuffle partition count, and a + `FinalPartitioned` aggregate. Global aggregation, or a single shuffle partition, + coalesces the partial states into a `Final` aggregate instead; with only one input + partition there is a single aggregate. Partial states and hash buckets never leave + the graph. +- Joins: two DataFusion hash repartitions with the same partition count feed a + partitioned `HashJoinExec`. Spark's build side is retained; DataFusion's input + swap projection restores the logical output order. +- Global sort and Top-K: each input partition is sorted (keeping only the Top-K + rows, if any), and a `SortPreservingMergeExec` merges the runs into one ordered + partition. Offset is applied once by a `GlobalLimitExec`. An unordered limit + gathers the partitions before the limit. +- Range: the native range partitions are merged with a sort-preserving merge, since + Spark may already have removed a redundant sort based on range ordering. + +## Memory and spill Review Comment: Could `memory_management.md` get a matching update? Its "Who allocates what" table says Comet's native memory is bounded by `memory_limit`, and the paragraph after it says every native reservation is forwarded to Spark's off-heap execution pool. Neither holds for a local query, whose reservations go to a per-query pool sized by `spark.comet.exec.local.memoryLimit` that Spark never sees. A row in that table, or a short note linking here, would keep the memory guide accurate once someone turns this on. ########## 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: Picking this up since the 1 in 32 failure is still there at `105afe87a`. The 512 KiB request is exactly the `sort_spill_reservation_bytes` that `size_sorters` sets for that test (8 MiB / 4 sorters / 4). In DataFusion 55.1, `sort_and_spill_in_mem_batches` frees `merge_reservation` before it sorts and spills, then calls `reserve_memory_for_merge` again at the end. In between, `FairSpillPool` hands the freed unspillable bytes to the other sorters' spillable shares, so the `try_resize` back to 512 KiB fails with `ResourcesExhausted` instead of spilling. `sort_in_place_threshold_bytes` never reaches that path, so I don't think any threshold value can close it. The `usize::MAX` threshold also adds a failure of its own. Every spill and every final in-memory sort now runs `concat_batches` over the sorter's whole buffer. String and binary columns use i32 offsets, so once one sorter buffers more than 2 GiB of a single such column, the concat fails with `OffsetOverflowError`. The default threshold never takes that path. It needs a `memoryLimit` of a few GiB, which seems likely for driver-side queries on a big machine. Since #6404 and EPIC #6410 are already moving Comet to DataFusion 56, would it be simpler to admit full global sorts only after that upgrade? Top-K could stay, since `SortExec` with a fetch uses `TopK` rather than `ExternalSorter`. -- 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]
