adriangb opened a new issue, #26141: URL: https://github.com/apache/datafusion/issues/26141
### Describe the bug `RepartitionExec` does not apply `datafusion.execution.spill_compression`. Its spill files are always uncompressed. `RepartitionExec::execute` builds its `SpillManager` without `.with_compression_type(context.session_config().spill_compression())`: https://github.com/apache/datafusion/blob/4978b3061c72753f73784f92d573eb0d903bcaca/datafusion/physical-plan/src/repartition/mod.rs#L1757-L1761 All other spilling operators pass the setting (sort, aggregate, sort-merge join, nested loop join). Thus a user who sets `spill_compression` gets compressed spill files from those operators, but full-size spill files from `RepartitionExec`. A real workload showed the effect. A hash repartition in front of an external sort wrote about 810 GB of uncompressed spill for a 32 GB Parquet input with `spill_compression = lz4_frame`. The lz4 spill of the sort for the same rows was about 100 GB. ### To Reproduce Add this test to the `tests` module in `datafusion/physical-plan/src/repartition/mod.rs` (it needs `use datafusion_common::config::SpillCompression;`), then run `cargo test -p datafusion-physical-plan --lib repartition_spill_honors_spill_compression`. The input is 20 batches of 8192 identical `UInt32` values, and a 1-byte memory limit makes each batch spill. ```rust /// Runs a spilling `RepartitionExec` with the given spill compression and /// returns `(spilled_rows, spilled_bytes)`. async fn repartition_spill_with_compression( spill_compression: SpillCompression, ) -> Result<(usize, usize)> { // Highly compressible input: 20 batches of 8192 identical values let schema = test_schema(false); let batch = RecordBatch::try_new( Arc::clone(&schema), vec![Arc::new(UInt32Array::from(vec![42; 8192]))], )?; let num_batches = 20; let input_partitions = vec![vec![batch; num_batches]]; // Tight memory limit to force every batch to spill let runtime = RuntimeEnvBuilder::default() .with_memory_limit(1, 1.0) .build_arc()?; let session_config = SessionConfig::new().with_spill_compression(spill_compression); let task_ctx = Arc::new( TaskContext::default() .with_runtime(runtime) .with_session_config(session_config), ); let exec = TestMemoryExec::try_new_exec(&input_partitions, Arc::clone(&schema), None)?; let exec = RepartitionExec::try_new(exec, Partitioning::RoundRobinBatch(4))?; let mut total_rows = 0; for i in 0..exec.partitioning().partition_count() { let mut stream = exec.execute(i, Arc::clone(&task_ctx))?; while let Some(result) = stream.next().await { let batch = result?; // Spilled data must read back intact regardless of the codec let values = as_uint32_array(batch.column(0))?; assert!(values.iter().all(|v| v == Some(42))); total_rows += batch.num_rows(); } } assert_eq!(total_rows, num_batches * 8192); let metrics = exec.metrics().unwrap(); Ok(( metrics.spilled_rows().unwrap(), metrics.spilled_bytes().unwrap(), )) } #[tokio::test] async fn repartition_spill_honors_spill_compression() -> Result<()> { let (uncompressed_rows, uncompressed_bytes) = repartition_spill_with_compression(SpillCompression::Uncompressed).await?; assert!(uncompressed_rows > 0, "Expected spilling to occur"); for spill_compression in [SpillCompression::Lz4Frame, SpillCompression::Zstd] { let (rows, bytes) = repartition_spill_with_compression(spill_compression).await?; assert_eq!( rows, uncompressed_rows, "Expected the same rows to spill with {spill_compression}" ); // The input is constant, so any codec shrinks it by far more than 2x assert!( bytes * 2 < uncompressed_bytes, "Expected {spill_compression} spill files to be much smaller than \ uncompressed ones: {bytes} vs {uncompressed_bytes} bytes" ); } Ok(()) } ``` On `main` at 4978b3061c72753f73784f92d573eb0d903bcaca the test fails: ``` Expected lz4_frame spill files to be much smaller than uncompressed ones: 679264 vs 679264 bytes ``` `spilled_bytes` of `RepartitionExec` for the 163,840 spilled rows: | `spill_compression` | `main` | with `.with_compression_type(...)` added | | --- | --- | --- | | `uncompressed` | 679,264 | 679,264 | | `lz4_frame` | 679,264 | 7,744 | | `zstd` | 679,264 | 5,024 | ### Expected behavior `RepartitionExec` writes its spill files with the codec from `datafusion.execution.spill_compression`, as the other spilling operators do. ### Additional context The read side needs no change. Spill files use the Arrow IPC Stream format, and the reader gets the codec from the stream. -- 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]
