apache / apache/datafusion

Investigate performance tradeoff in compressing spill files

Open
#16,367 3 comments 0 reactions 0 assignees View on GitHub
enhancement
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Is your feature request related to a problem or challenge?

Part of https://github.com/apache/datafusion/issues/16065 , Related to https://github.com/apache/datafusion/issues/14078

### Background
PR https://github.com/apache/datafusion/pull/16268 will introduce compression option when writing spill files in disk. These options are limited to compression options to what Arrow IPC Stream Writer supports : `zstd`, `lz4_frame` (compression level is not configurable)

To further enhance spilling execution, we need to investigate
- CPU / I/O tradeoff when `zstd` or `lz4_frame` compression is enabled i.e. compression ratio, extra latency spent for compression
- Current arrow ipc stream writer always write `batch` at a time in `append_batch`. In terms of compression, it is not sure yet how much single batch can benefit from compression.
- whether we need separate `Writer` or `Reader` implementation instead of IPC Stream Writer.
- how to introduce sort of `adaptiveness`.

### Describe the solution you'd like

First, we need to track (or update) how many bytes are written in spill files. Datafusion currently tracks `spilled_bytes` as part of `SpillMetrics`, but it is calculated based on in memory array size, which would be different from actual spill files size especially when we compress spill files.

Second, update the benchmarks or write a separate benchmarks to see the performance characteristics. One possible way is writing out spill-related metrics to output.json when running benches like tpch with `debug` option. Another idea is to generate some spill files for microbenchmark testing only spill writing - reading process.

### Describe alternatives you've considered

_No response_

### Additional context

_No response_

Contributor guide

Open the contributing guide

Research direction

Start by tracing SpillMetrics and its spilled_bytes calculation, then review the Arrow IPC Stream Writer spill path and the existing tpch benchmarks with the debug option. Compare compressed and uncompressed spill writing and reading, including compression ratio and latency. Done means actual spill-file bytes are tracked and benchmark output exposes enough metrics to evaluate the tradeoffs.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust
Domain
data-engineering, performance
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.