Skip to content

Commit aa40adb

Browse files
committed
timely-util: spill young chunk generations uncompressed
A chunk at generational depth d is rewritten with frequency proportional to 2^-d under geometric merging, so compressing a shallow generation buys pool bytes back for only a short stay at a guaranteed near-term codec round-trip: fresh (depth 0) chunks are consumed by their first merge with certainty. Spill bodies below a configurable compression depth floor (default 1) under a new identity ExtentCodec in mz_ore::pool instead: still inserted into the pool, so they stay budget-accounted and swap-backed like every extent, but encode and decode reduce to copies. Keeping young generations out of the pool entirely was rejected: exempted bytes are invisible to the pool budget, and coalesced carries at 2 MiB each across many operators could accumulate unbudgeted resident state exactly when the system is busiest. The floor is wired through the replica-scoped column_chunk_compress_min_depth dyncfg alongside the other chunk knobs.
1 parent 0a18050 commit aa40adb

6 files changed

Lines changed: 181 additions & 8 deletions

File tree

‎misc/python/materialize/mzcompose/__init__.py‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -592,6 +592,7 @@ def get_default_system_parameters(
592592
"linear_join_yielding",
593593
"enable_column_paged_batcher",
594594
"enable_column_paged_batcher_spill",
595+
"column_chunk_compress_min_depth",
595596
"column_paged_batcher_budget_fraction",
596597
"column_paged_batcher_lz4",
597598
"column_paged_batcher_swap_pageout",

‎misc/python/materialize/parallel_workload/action.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3052,6 +3052,11 @@ def __init__(
30523052
"0.02",
30533053
]
30543054
self.flags_with_values["enable_upsert_paged_spill"] = BOOLEAN_FLAG_VALUES
3055+
self.flags_with_values["column_chunk_compress_min_depth"] = [
3056+
"0", # compress every spilled body
3057+
"1", # the default: fresh chunks store uncompressed
3058+
"4", # exempt several young generations
3059+
]
30553060
# 0 forces the estimated-size path for every table, the default forces
30563061
# the exact COUNT(*) path for workload-sized tables.
30573062
self.flags_with_values["mysql_source_snapshot_exact_count_max_rows"] = [

‎src/compute-types/src/dyncfgs.rs‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,24 @@ pub const ENABLE_COLUMN_PAGED_BATCHER_SPILL: Config<bool> = Config::new(
6565
)
6666
.scoped(ParameterScope::Replica);
6767

68+
/// The youngest chunk generation whose spilled bodies are compressed.
69+
///
70+
/// A chunk at generational depth `d` is rewritten with frequency
71+
/// proportional to `2^-d` under geometric merging, so compressing shallow
72+
/// generations buys pool bytes back for only a short stay at a guaranteed
73+
/// near-term codec round-trip. Generations below the floor spill under the
74+
/// identity codec: fully budgeted and swap-backed, with encode and decode
75+
/// reduced to copies. The default exempts only fresh (depth 0) chunks,
76+
/// which their first merge consumes with certainty. `0` compresses every
77+
/// spilled body.
78+
pub const COLUMN_CHUNK_COMPRESS_MIN_DEPTH: Config<u32> = Config::new(
79+
"column_chunk_compress_min_depth",
80+
1,
81+
"The youngest chunk generation whose spilled bodies are lz4-compressed in the buffer \
82+
pool; younger generations store uncompressed. 0 compresses every spilled body.",
83+
)
84+
.scoped(ParameterScope::Replica);
85+
6886
/// Resident-bytes budget fraction for chunk spilling. Two consumers read
6987
/// it: the column pager's tiered policy multiplies it against the
7088
/// announced memory limit, and the buffer pool (`mz_ore::pool`)
@@ -601,4 +619,5 @@ pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet {
601619
.add(&COLUMN_PAGED_BATCHER_SPILL_WORKER_COUNT)
602620
.add(&COLUMN_PAGED_BATCHER_EAGER_BACKING)
603621
.add(&COLUMN_PAGED_BATCHER_POOL_RSS_TARGET_FRACTION)
622+
.add(&COLUMN_CHUNK_COMPRESS_MIN_DEPTH)
604623
}

‎src/compute/src/compute_state.rs‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -428,6 +428,13 @@ impl ComputeState {
428428
warn!("chunk spill: buffer pool unavailable; chunks stay resident");
429429
}
430430
}
431+
432+
// The generational depth floor below which spilled bodies store
433+
// uncompressed. Subsystem-independent, so applied here alongside
434+
// the rest of the process-wide chunk configuration.
435+
let compress_min_depth =
436+
u8::try_from(COLUMN_CHUNK_COMPRESS_MIN_DEPTH.get(config)).unwrap_or(u8::MAX);
437+
mz_timely_util::columnar::chunk::set_compress_min_depth(compress_min_depth);
431438
}
432439

433440
// Remember the maintenance interval locally to avoid reading it from the config set on

‎src/ore/src/pool.rs‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,28 @@ pub trait ExtentCodec: std::fmt::Debug + Send + Sync {
126126
fn decode(&self, stored: &[u8], body: &mut [u8]);
127127
}
128128

129+
/// The identity [`ExtentCodec`]: the stored form is the body. Encode and
130+
/// decode are copies, and range reads copy the range directly, so a chunk
131+
/// stored under this codec pays no compression work in either direction
132+
/// while remaining fully budgeted and swap-backed like any other extent.
133+
#[derive(Debug)]
134+
pub struct IdentityCodec;
135+
136+
/// The [`IdentityCodec`] instance to pass to [`Pool::insert_with`].
137+
pub static IDENTITY_CODEC: IdentityCodec = IdentityCodec;
138+
139+
impl ExtentCodec for IdentityCodec {
140+
fn encode(&self, body: &[u8], out: &mut Vec<u8>) {
141+
out.clear();
142+
out.extend_from_slice(body);
143+
}
144+
145+
fn decode(&self, stored: &[u8], body: &mut [u8]) {
146+
assert_eq!(stored.len(), body.len(), "identity stored form is the body");
147+
body.copy_from_slice(stored);
148+
}
149+
}
150+
129151
/// The largest stored form [`ExtentCodec::encode`] may produce for a
130152
/// `body_len`-byte body: an incompressible-input expansion matching lz4's
131153
/// worst case plus a four-byte length prefix. The extent store's size-class
@@ -3703,6 +3725,24 @@ mod tests {
37033725
assert_eq!(pool.stats().resident_bytes, 0);
37043726
}
37053727

3728+
/// The identity codec stores the body verbatim: eviction and reads,
3729+
/// whole and by range, reconstruct it unchanged.
3730+
#[mz_ore::test]
3731+
fn identity_codec_round_trips() {
3732+
let pool = test_pool(usize::MAX);
3733+
let want = payload(SMALL, 601);
3734+
let h = pool.insert_with(SMALL, ChunkHints::default(), &IDENTITY_CODEC, |dst| {
3735+
dst.copy_from_slice(&want);
3736+
});
3737+
assert_eq!(read(&h), want);
3738+
pool.evict(&h);
3739+
assert_eq!(read(&h), want, "round-trips through the extent");
3740+
pool.evict(&h);
3741+
let mut range = Vec::new();
3742+
h.read_range_into(8..24, &mut range);
3743+
assert_eq!(range, want[8..24], "range reads copy the range directly");
3744+
}
3745+
37063746
#[mz_ore::test]
37073747
fn insert_with_fills_in_place() {
37083748
let pool = test_pool(usize::MAX);

‎src/timely-util/src/columnar/chunk.rs‎

Lines changed: 109 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -44,18 +44,18 @@
4444
//! resident fence metadata so a probe set faults only the chunk bodies it
4545
//! actually touches.
4646
47-
use std::cell::RefCell;
47+
use std::cell::{Cell, RefCell};
4848
use std::collections::VecDeque;
4949
use std::rc::Rc;
50-
use std::sync::atomic::{AtomicBool, Ordering};
50+
use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
5151

5252
use columnar::bytes::indexed;
5353
use columnar::{Borrow, BorrowedOf, Columnar, Container as _, FromBytes, Index, Len, Push as _};
5454
use differential_dataflow::difference::Semigroup;
5555
use differential_dataflow::lattice::Lattice;
5656
use differential_dataflow::trace::chunk::Chunk;
5757
use mz_ore::cast::CastFrom;
58-
use mz_ore::pool::{ChunkHandle, ChunkHints, ExtentCodec, Pool};
58+
use mz_ore::pool::{ChunkHandle, ChunkHints, ExtentCodec, IDENTITY_CODEC, Pool};
5959
use timely::Accountable;
6060
use timely::container::{ContainerBuilder, PushInto};
6161
use timely::dataflow::channels::ContainerBytes;
@@ -78,6 +78,11 @@ thread_local! {
7878
/// pool without touching process-global state.
7979
static SPILL_OVERRIDE: RefCell<Option<Pool>> = const { RefCell::new(None) };
8080

81+
/// A thread-scoped depth-floor override, taking precedence over the
82+
/// global value. Lets tests and benches pin the floor without racing
83+
/// concurrently running tests on the process-global state.
84+
static COMPRESS_MIN_DEPTH_OVERRIDE: Cell<Option<u8>> = const { Cell::new(None) };
85+
8186
/// Reusable staging for call-scoped reads of spilled bodies.
8287
static READ_SCRATCH: RefCell<Vec<u64>> = const { RefCell::new(Vec::new()) };
8388
}
@@ -113,6 +118,51 @@ pub fn set_spill_override(pool: Option<Pool>) {
113118
SPILL_OVERRIDE.with(|cell| *cell.borrow_mut() = pool);
114119
}
115120

121+
/// The youngest generational depth whose spilled bodies are compressed. See
122+
/// [`set_compress_min_depth`].
123+
static COMPRESS_MIN_DEPTH: AtomicU8 = AtomicU8::new(DEFAULT_COMPRESS_MIN_DEPTH);
124+
125+
/// Set the youngest generational depth whose spilled bodies are compressed.
126+
///
127+
/// A chunk at depth `d` is rewritten (merged, extracted, advanced) with
128+
/// frequency proportional to `2^-d` under geometric merging, so compressing
129+
/// a shallow chunk buys a short stay in the pool at the cost of a guaranteed
130+
/// near-term codec round-trip: the body is encoded only to be read back and
131+
/// decoded by the next rewrite. Generations below the floor spill under the
132+
/// identity codec instead: still budgeted and swap-backed like every extent,
133+
/// but encode and decode are copies. The floor never exempts a body from the
134+
/// pool, so it cannot grow unbudgeted resident state.
135+
///
136+
/// `0` compresses every spilled body. Consulted at every commit, so changes
137+
/// apply to running dataflows.
138+
pub fn set_compress_min_depth(depth: u8) {
139+
COMPRESS_MIN_DEPTH.store(depth, Ordering::Relaxed);
140+
}
141+
142+
/// Set or unset a thread-scoped depth-floor override, taking precedence over
143+
/// [`set_compress_min_depth`]. For tests and benches, which run concurrently
144+
/// and must not race on the process-global floor.
145+
pub fn set_compress_min_depth_override(depth: Option<u8>) {
146+
COMPRESS_MIN_DEPTH_OVERRIDE.with(|cell| cell.set(depth));
147+
}
148+
149+
/// The depth floor in effect for this thread's commits.
150+
fn compress_min_depth() -> u8 {
151+
COMPRESS_MIN_DEPTH_OVERRIDE
152+
.with(|cell| cell.get())
153+
.unwrap_or_else(|| COMPRESS_MIN_DEPTH.load(Ordering::Relaxed))
154+
}
155+
156+
/// The codec a body at `depth` stores under: identity below the compression
157+
/// floor, lz4 at and past it.
158+
fn codec_for_depth(depth: u8) -> &'static dyn ExtentCodec {
159+
if depth < compress_min_depth() {
160+
&IDENTITY_CODEC
161+
} else {
162+
&LZ4_CODEC
163+
}
164+
}
165+
116166
/// The pool committed chunks spill to, if any.
117167
fn spill_pool() -> Option<Pool> {
118168
if let Some(pool) = SPILL_OVERRIDE.with(|cell| cell.borrow().clone()) {
@@ -161,6 +211,16 @@ const COMMIT_BYTES: usize = 2 << 20;
161211
/// unbudgeted heap, and no accounting here would catch it.
162212
const SPILL_MIN_BYTES: usize = 64 << 10;
163213

214+
/// The default compression depth floor: fresh (depth 0) bodies spill
215+
/// uncompressed.
216+
///
217+
/// A fresh chunk is consumed by its first merge with certainty, so
218+
/// compressing it can never save pool bytes for longer than one merge
219+
/// cadence and always costs a full encode plus decode. Depth 1 and beyond
220+
/// have survived a merge and wait geometrically longer for the next, so
221+
/// their compression amortizes.
222+
const DEFAULT_COMPRESS_MIN_DEPTH: u8 = 1;
223+
164224
/// Whether a column is big enough to commit on its own. A monotone
165225
/// threshold, so settle's carry, which grows by whole chunks, cannot step
166226
/// over it.
@@ -309,14 +369,20 @@ impl<D: Columnar, T: Columnar, R: Columnar> ColumnChunk<D, T, R> {
309369

310370
/// Spill a non-empty column into `pool` unconditionally, capturing the
311371
/// resident fence metadata.
372+
///
373+
/// Generations below the compression depth floor store under the
374+
/// identity codec: rewritten too soon for compression to amortize, they
375+
/// stay budgeted and swap-backed while encode and decode reduce to
376+
/// copies.
312377
fn spill_body(column: Column<(D, T, R)>, pool: &Pool, depth: u8) -> Self {
378+
let codec = codec_for_depth(depth);
313379
let len_bytes = column.length_in_bytes();
314380
let view = column.borrow();
315381
let records = view.len();
316382
let mut fences = D::Container::default();
317383
fences.push(view.0.get(0));
318384
fences.push(view.0.get(records - 1));
319-
let handle = spill_column(column, pool, len_bytes, ChunkHints { depth });
385+
let handle = spill_column(column, pool, len_bytes, ChunkHints { depth }, codec);
320386
ColumnChunk::Spilled(Rc::new(SpilledBody {
321387
records,
322388
fences,
@@ -379,13 +445,14 @@ fn spill_column<C: Columnar>(
379445
pool: &Pool,
380446
len_bytes: usize,
381447
hints: ChunkHints,
448+
codec: &'static dyn ExtentCodec,
382449
) -> ChunkHandle {
383450
mz_ore::soft_assert_eq_no_log!(len_bytes % 8, 0);
384451
match column {
385-
Column::Align(words) => pool.insert_with(words.len(), hints, &LZ4_CODEC, |dst| {
386-
dst.copy_from_slice(&words)
387-
}),
388-
other => pool.insert_with(len_bytes / 8, hints, &LZ4_CODEC, |dst| {
452+
Column::Align(words) => {
453+
pool.insert_with(words.len(), hints, codec, |dst| dst.copy_from_slice(&words))
454+
}
455+
other => pool.insert_with(len_bytes / 8, hints, codec, |dst| {
389456
let bytes: &mut [u8] = bytemuck::cast_slice_mut(dst);
390457
let mut cursor = std::io::Cursor::new(bytes);
391458
other.into_bytes(&mut cursor);
@@ -1701,6 +1768,39 @@ mod tests {
17011768
set_spill_override(None);
17021769
}
17031770

1771+
/// The compression depth floor picks the codec, not whether a body
1772+
/// spills: shallow generations store at identity, the floor and deeper
1773+
/// at lz4, and every depth spills and round-trips.
1774+
#[mz_ore::test]
1775+
fn spill_codec_depth_floor() {
1776+
set_spill_override(Some(test_pool()));
1777+
set_compress_min_depth_override(Some(2));
1778+
// Codec identity via Debug: ZST statics and dyn vtables make
1779+
// pointer comparison unreliable.
1780+
let codec_name = |depth: u8| format!("{:?}", codec_for_depth(depth));
1781+
assert_eq!(codec_name(0), "IdentityCodec");
1782+
assert_eq!(codec_name(1), "IdentityCodec");
1783+
assert_eq!(codec_name(2), "Lz4Codec");
1784+
assert_eq!(codec_name(u8::MAX), "Lz4Codec");
1785+
1786+
let data: Vec<Tuple> = (0..20_000u64).map(|i| ((i, 0), 0, 1i64)).collect();
1787+
let data = consolidate(data);
1788+
let column = build_column(&data);
1789+
for depth in [0u8, 1, 2, 3] {
1790+
let chunk = TestChunk::commit(column.clone(), depth);
1791+
assert!(chunk.is_spilled(), "depth {depth} must spill");
1792+
assert_eq!(collect_column(&chunk.into_column()), data);
1793+
}
1794+
set_spill_override(None);
1795+
set_compress_min_depth_override(None);
1796+
1797+
// The default floor stores only fresh (depth 0) bodies at identity.
1798+
set_compress_min_depth_override(Some(DEFAULT_COMPRESS_MIN_DEPTH));
1799+
assert_eq!(codec_name(0), "IdentityCodec");
1800+
assert_eq!(codec_name(1), "Lz4Codec");
1801+
set_compress_min_depth_override(None);
1802+
}
1803+
17041804
/// The compute and storage spill gates compose as an OR: either gate
17051805
/// routes commits to the installed pool, and each setter writes only its
17061806
/// own gate.
@@ -1732,6 +1832,7 @@ mod tests {
17321832
assert!(commit(&col), "the compute gate alone spills");
17331833
set_compute_spill_enabled(false);
17341834
assert!(!commit(&col), "both gates off again");
1835+
set_compress_min_depth_override(None);
17351836
}
17361837

17371838
/// Re-spilling an already-serialized body exercises the `Column::Align`

0 commit comments

Comments
 (0)