Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -335,6 +335,23 @@ Native operators reserve through DataFusion's `MemoryConsumer` / `MemoryReservat

An operator that never calls `try_grow` is invisible to the pool no matter how much memory it uses.

### Sort and whole-partition windows

The sort merge reservation is capped at 1/32 of the configured off-heap budget per
concurrent Spark task (executor cores divided by task CPUs), up to DataFusion's default.
This leaves room for input batches on small executors; the spillable merge can grow its
reservation when it needs more. It does not increase the memory pool or suppress allocation
failures. An individual batch still has to fit the available execution budget.

`PartitionAggregateWindowExec` handles full-partition `sum`, `avg`, `count`, `min`, and
`max` frames. It updates the existing native accumulators incrementally and reserves the
retained input batches. On reservation failure it spills those rows through DataFusion's
spill manager, then replays one spill file at a time with the final aggregate columns.
Only the current window partition is retained, and small partitions avoid disk entirely.
Accumulator state is reserved separately. This preserves native execution without retaining
an entire wide partition in memory. Other window frames continue to use DataFusion's existing
window operators; this is not a general spill implementation for all window functions.

## Crossing the FFI boundary

Batches move between the JVM and native over the Arrow C Data and C Stream interfaces, which are
Expand Down
41 changes: 41 additions & 0 deletions native/core/src/execution/jni_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -666,6 +666,7 @@ pub unsafe extern "system" fn Java_org_apache_comet_Native_createPlan(
task_cpus as usize,
&spark_config,
&spark_plan,
(off_heap_mode != JNI_FALSE).then_some(memory_limit as usize),
)?;

let plan_creation_time = start.elapsed();
Expand Down Expand Up @@ -819,7 +820,24 @@ fn configure_skip_partial_aggregation(config: &mut SessionConfig, plan: &Operato
}
}

/// DataFusion's fixed 10 MiB merge reserve can consume most of a small Spark task's
/// share before the sorter admits its first batch. Cap this eager reservation at 1/32
/// of the per-task budget. The spillable merge can grow it as needed; larger executors
/// retain the upstream default. Explicit testing overrides are applied afterwards.
fn configure_sort_spill_reservation(
config: &mut SessionConfig,
off_heap_limit: usize,
executor_cores: usize,
task_cpus: usize,
) {
let concurrent_tasks = (executor_cores / task_cpus.max(1)).max(1);
let cap = (off_heap_limit / concurrent_tasks / 32).max(1);
let reservation = &mut config.options_mut().execution.sort_spill_reservation_bytes;
*reservation = (*reservation).min(cap);
}

/// Configure DataFusion session context.
#[allow(clippy::too_many_arguments)]
fn prepare_datafusion_session_context(
batch_size: usize,
memory_pool: Arc<dyn MemoryPool>,
Expand All @@ -828,6 +846,7 @@ fn prepare_datafusion_session_context(
task_cpus: usize,
spark_config: &HashMap<String, String>,
spark_plan: &Operator,
off_heap_limit: Option<usize>,
) -> CometResult<SessionContext> {
let paths = local_dirs.into_iter().map(PathBuf::from).collect();
let disk_manager = DiskManagerBuilder::default()
Expand All @@ -845,6 +864,11 @@ fn prepare_datafusion_session_context(
// modified by changing spark.task.cpus in the Spark config.
.with_batch_size(batch_size);

if let Some(limit) = off_heap_limit {
let executor_cores = spark_config.get_usize(SPARK_EXECUTOR_CORES, 1);
configure_sort_spill_reservation(&mut session_config, limit, executor_cores, task_cpus);
}

// Translate the Comet-namespaced row-level pushdown flag into the equivalent
// DataFusion session options. `pushdown_filters` enables the parquet reader's
// RowFilter evaluation during decode (late materialization); `reorder_filters`
Expand Down Expand Up @@ -1837,6 +1861,23 @@ mod tests {
use datafusion_comet_proto::spark_expression::{AggExpr, Count, Expr, Sum};
use datafusion_comet_proto::spark_operator::{HashAggregate, ShuffleWriter};

#[test]
fn sort_merge_reserve_scales_with_the_task_budget() {
let mut config = SessionConfig::new();
let default = config.options().execution.sort_spill_reservation_bytes;
configure_sort_spill_reservation(&mut config, 64 * 1024 * 1024, 2, 1);
assert_eq!(
config.options().execution.sort_spill_reservation_bytes,
1024 * 1024
);
let mut config = SessionConfig::new();
configure_sort_spill_reservation(&mut config, 8 * 1024 * 1024 * 1024, 4, 1);
assert_eq!(
config.options().execution.sort_spill_reservation_bytes,
default
);
}

#[test]
fn skip_partial_eligibility_is_fail_closed() {
let count = AggExpr {
Expand Down
2 changes: 2 additions & 0 deletions native/core/src/execution/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,9 @@ pub use iceberg_write::IcebergWriteExec;
mod parquet_writer;
pub use parquet_writer::{ParquetCompression, ParquetWriterExec};
mod csv_scan;
mod partition_aggregate_window;
pub mod projection;
pub use partition_aggregate_window::PartitionAggregateWindowExec;
mod sample;
pub use sample::SampleExec;
mod rank_limit;
Expand Down
Loading