mz_row_spine/lib.rs
1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Types and traits in support of containers for row-encoded byte slices.
11//!
12//! This includes the vanilla `bytes_container` that holds byte slices in contiguous
13//! allocations, as well as a `dictionary` encoding wrapper that is able to rewrite
14//! the byte slices to use spare tags in each column to reference common values.
15
16pub use self::arc_batch::{ArcBatch, ArcBuilder};
17pub use self::dictionary::DatumContainer;
18pub use self::dictionary::DatumSeq;
19pub use self::dictionary::builders::RowRowColPagedState;
20pub use self::offset_opt::OffsetOptimized;
21pub use self::spines::{
22 ArcOrdKeyBuilder, ArcOrdKeySpine, ArcOrdValBuilder, ArcOrdValSpine, FundedValRowSpine,
23 RowBatcher, RowBuilder, RowRowBatcher, RowRowBuilder, RowRowColPagedBuilder, RowRowSpine,
24 RowSpine, RowValBatcher, RowValBuilder, RowValSpine, ValRowBatcher, ValRowBuilder,
25 ValRowColPagedBuilder, ValRowSpine,
26};
27
28mod arc_batch;
29
30use differential_dataflow::trace::implementations::OffsetList;
31
32/// Enable per-column dictionary compression in row containers.
33pub static DICTIONARY_COMPRESSION: std::sync::atomic::AtomicBool =
34 std::sync::atomic::AtomicBool::new(false);
35
36/// Spines specialized to contain `Row` types in keys and values.
37mod spines {
38 use columnation::Columnation;
39 use differential_dataflow::trace::implementations::BatchContainer;
40 use differential_dataflow::trace::implementations::Layout;
41 use differential_dataflow::trace::implementations::Update;
42 use differential_dataflow::trace::implementations::Vector;
43 use differential_dataflow::trace::implementations::merge_batcher::MergeBatcher;
44 use differential_dataflow::trace::implementations::ord_neu::{
45 OrdKeyBatch, OrdKeyBuilder, OrdValBatch, OrdValBuilder,
46 };
47 use differential_dataflow::trace::implementations::spine_fueled::Spine;
48 use mz_repr::Row;
49 use mz_timely_util::columnation::{ColInternalMerger, ColumnationStack};
50
51 use crate::arc_batch::{ArcBatch, ArcBuilder};
52 use crate::{DatumContainer, OffsetOptimized};
53
54 /// Batcher matching `mz_compute::typedefs::KeyValBatcher`, redeclared
55 /// locally so this crate does not need to depend on `mz_compute`.
56 type KeyValBatcher<K, V, T, D> = MergeBatcher<ColInternalMerger<(K, V), T, D>>;
57 type KeyBatcher<K, T, D> = KeyValBatcher<K, (), T, D>;
58
59 pub type RowRowSpine<T, R> = Spine<ArcBatch<OrdValBatch<RowRowLayout<((Row, Row), T, R)>>>>;
60 pub type RowRowBatcher<T, R> = KeyValBatcher<Row, Row, T, R>;
61 pub type RowRowBuilder<T, R> = ArcBuilder<crate::dictionary::builders::RowRowBuilder<T, R>>;
62
63 /// `RowRowBuilder` variant that consumes [`ColumnBody`] chunks. Pairs with
64 /// any batcher whose chains are bodies, spillable ([`AccountedChunkBatcher`]
65 /// behind an [`UnchunkBuilder`]) or resident ([`Col2ValColBatcher`]) alike,
66 /// so the `Paged` in the name records where it started rather than a
67 /// restriction. Installs a dictionary codec at seal time, gathering
68 /// statistics from the sealed chain, so columnar arrangements compress on
69 /// the same footing as the columnation-fed [`RowRowBuilder`].
70 ///
71 /// [`AccountedChunkBatcher`]: mz_timely_util::columnar::chunk::AccountedChunkBatcher
72 /// [`Col2ValColBatcher`]: mz_timely_util::columnar::Col2ValColBatcher
73 /// [`ColumnBody`]: mz_timely_util::columnar::body::ColumnBody
74 /// [`UnchunkBuilder`]: mz_timely_util::columnar::chunk::UnchunkBuilder
75 pub type RowRowColPagedBuilder<T, R> =
76 ArcBuilder<crate::dictionary::builders::RowRowColPagedBuilder<T, R>>;
77
78 pub type RowValSpine<V, T, R> = Spine<ArcBatch<OrdValBatch<RowValLayout<((Row, V), T, R)>>>>;
79 pub type RowValBatcher<V, T, R> = KeyValBatcher<Row, V, T, R>;
80 pub type RowValBuilder<V, T, R> =
81 ArcBuilder<crate::dictionary::builders::RowValBuilder<V, T, R>>;
82
83 /// Key-only `Row` spine. `DC` is the diff container, see `RowLayout`.
84 pub type RowSpine<T, R, DC = ColumnationStack<R>> =
85 Spine<ArcBatch<OrdKeyBatch<RowLayout<((Row, ()), T, R), DC>>>>;
86 pub type RowBatcher<T, R> = KeyBatcher<Row, T, R>;
87 pub type RowBuilder<T, R, DC = ColumnationStack<R>> =
88 ArcBuilder<crate::dictionary::builders::RowBuilder<T, R, DC>>;
89
90 pub type ValRowSpine<K, T, R> = Spine<ArcBatch<OrdValBatch<ValRowLayout<((K, Row), T, R)>>>>;
91 /// A `ValRowSpine` whose optional exertion is funded by its input.
92 pub type FundedValRowSpine<K, T, R> =
93 mz_timely_util::funded_spine::Spine<ArcBatch<OrdValBatch<ValRowLayout<((K, Row), T, R)>>>>;
94 pub type ValRowBatcher<K, T, R> = KeyValBatcher<K, Row, T, R>;
95 pub type ValRowBuilder<K, T, R> =
96 ArcBuilder<crate::dictionary::builders::ValRowBuilder<K, T, R>>;
97
98 /// `ValRowBuilder` variant that consumes [`Column`] chunks. Pairs with
99 /// `Col2ValPagedBatcher<K, Row, T, R>` for the spillable arrange path where
100 /// keys are arbitrary `Columnar` values (e.g. `UpsertKey`) and values are
101 /// packed `Row` bytes. Installs a dictionary codec on the value container at
102 /// seal time, gathering statistics from the sealed `Column` chain; keys are
103 /// not `Row`-shaped and so are left uncompressed.
104 ///
105 /// [`Column`]: mz_timely_util::columnar::Column
106 pub type ValRowColPagedBuilder<K, T, R> =
107 ArcBuilder<crate::dictionary::builders::ValRowColPagedBuilder<K, T, R>>;
108
109 /// A generic `Arc`-backed key/value spine, for callers outside `mz_compute` that need an
110 /// arrangement over non-`Row`-specialized types. The `Arc` handle rides on the local
111 /// [`ArcBatch`] newtype, so no differential-side `Arc` batch impls are required.
112 pub type ArcOrdValSpine<K, V, T, R> = Spine<ArcBatch<OrdValBatch<Vector<((K, V), T, R)>>>>;
113 /// Generic `Arc`-backed key-only spine. See [`ArcOrdValSpine`].
114 pub type ArcOrdKeySpine<K, T, R> = Spine<ArcBatch<OrdKeyBatch<Vector<((K, ()), T, R)>>>>;
115 /// Builder pairing with [`ArcOrdValSpine`].
116 pub type ArcOrdValBuilder<K, V, T, R> =
117 ArcBuilder<OrdValBuilder<Vector<((K, V), T, R)>, Vec<((K, V), T, R)>>>;
118 /// Builder pairing with [`ArcOrdKeySpine`].
119 pub type ArcOrdKeyBuilder<K, T, R> =
120 ArcBuilder<OrdKeyBuilder<Vector<((K, ()), T, R)>, Vec<((K, ()), T, R)>>>;
121
122 /// A layout based on timely stacks
123 pub struct RowRowLayout<U: Update<Key = Row, Val = Row>> {
124 phantom: std::marker::PhantomData<U>,
125 }
126 pub struct RowValLayout<U: Update<Key = Row>> {
127 phantom: std::marker::PhantomData<U>,
128 }
129 /// Layout for key-only `Row` updates. `DC` is the diff `BatchContainer`, a
130 /// columnation stack by default.
131 pub struct RowLayout<U, DC = ColumnationStack<<U as Update>::Diff>>
132 where
133 U: Update<Key = Row, Val = ()>,
134 {
135 phantom: std::marker::PhantomData<(U, DC)>,
136 }
137 /// Mirror of [`RowValLayout`] with the roles swapped: arbitrary `Columnation`
138 /// keys with `Row` values stored as packed bytes in a [`DatumContainer`].
139 pub struct ValRowLayout<U: Update<Val = Row>> {
140 phantom: std::marker::PhantomData<U>,
141 }
142
143 impl<U: Update<Key = Row, Val = Row>> Layout for RowRowLayout<U>
144 where
145 U::Time: Columnation,
146 U::Diff: Columnation,
147 {
148 type KeyContainer = DatumContainer;
149 type ValContainer = DatumContainer;
150 type TimeContainer = ColumnationStack<U::Time>;
151 type DiffContainer = ColumnationStack<U::Diff>;
152 type OffsetContainer = OffsetOptimized;
153 }
154 impl<U: Update<Key = Row>> Layout for RowValLayout<U>
155 where
156 U::Val: Columnation,
157 U::Time: Columnation,
158 U::Diff: Columnation,
159 {
160 type KeyContainer = DatumContainer;
161 type ValContainer = ColumnationStack<U::Val>;
162 type TimeContainer = ColumnationStack<U::Time>;
163 type DiffContainer = ColumnationStack<U::Diff>;
164 type OffsetContainer = OffsetOptimized;
165 }
166 impl<U: Update<Key = Row, Val = ()>, DC> Layout for RowLayout<U, DC>
167 where
168 U::Time: Columnation,
169 DC: BatchContainer<Owned = U::Diff>,
170 {
171 type KeyContainer = DatumContainer;
172 type ValContainer = ColumnationStack<()>;
173 type TimeContainer = ColumnationStack<U::Time>;
174 type DiffContainer = DC;
175 type OffsetContainer = OffsetOptimized;
176 }
177 impl<U: Update<Val = Row>> Layout for ValRowLayout<U>
178 where
179 U::Key: Columnation,
180 U::Time: Columnation,
181 U::Diff: Columnation,
182 {
183 type KeyContainer = ColumnationStack<U::Key>;
184 type ValContainer = DatumContainer;
185 type TimeContainer = ColumnationStack<U::Time>;
186 type DiffContainer = ColumnationStack<U::Diff>;
187 type OffsetContainer = OffsetOptimized;
188 }
189}
190
191#[cfg(test)]
192mod tests {
193 use crate::DatumContainer;
194 use crate::spines::{RowLayout, RowRowLayout, RowValLayout};
195 use differential_dataflow::trace::implementations::BatchContainer;
196 use differential_dataflow::trace::implementations::ord_neu::{OrdKeyBatch, OrdValBatch};
197 use mz_repr::adt::date::Date;
198 use mz_repr::adt::interval::Interval;
199 use mz_repr::{Datum, Diff, Row, SqlScalarType, Timestamp};
200 use mz_timely_util::columnation::ColumnationStack;
201
202 fn assert_send_sync<T: Send + Sync>() {}
203
204 /// The batch types backing our spines must stay `Send + Sync`, so that batches
205 /// can be shared across threads (for example behind an `Arc`) to serve reads
206 /// from outside the worker that maintains the trace. This holds because the
207 /// backing containers bottom out in `Vec`s, lgalloc regions, and `CompactBytes`,
208 /// all of which are thread-safe.
209 #[mz_ore::test]
210 fn batches_are_send_sync() {
211 assert_send_sync::<OrdValBatch<RowRowLayout<((Row, Row), Timestamp, Diff)>>>();
212 assert_send_sync::<OrdValBatch<RowValLayout<((Row, Row), Timestamp, Diff)>>>();
213 assert_send_sync::<OrdKeyBatch<RowLayout<((Row, ()), Timestamp, Diff)>>>();
214 assert_send_sync::<ColumnationStack<((Row, Row), Timestamp, Diff)>>();
215 }
216
217 /// A chain of chunks seals with the chain's dictionary codecs installed,
218 /// so its batch compresses like one sealed from columns directly. A seal
219 /// that skips the install stores raw bytes and, per
220 /// `DatumContainer::promote_stats_to_codec`, drags every merge it joins
221 /// down with it.
222 #[mz_ore::test]
223 #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
224 fn chunked_seal_installs_codecs() {
225 use std::sync::atomic::Ordering;
226
227 use differential_dataflow::trace::implementations::ord_neu::OrdValBatch;
228 use differential_dataflow::trace::{Builder, Description};
229 use mz_timely_util::columnar::body::ColumnBody;
230 use mz_timely_util::columnar::chunk::{ColumnChunk, UnchunkBuilder};
231 use timely::container::PushInto;
232 use timely::progress::{Antichain, Timestamp as _};
233
234 use crate::ArcBatch;
235
236 // Gate the dictionary path on. Safe for other tests: the flag only
237 // controls whether codecs are built, never decode results.
238 crate::DICTIONARY_COMPRESSION.store(true, Ordering::Relaxed);
239
240 // Low-cardinality rows, well under `STATS_THRESHOLD` (64Ki pushes), so
241 // a codec can only come from the seal-time install.
242 let time = Timestamp::minimum();
243 let updates: Vec<((Row, Row), Timestamp, i64)> = (0..2_000i64)
244 .map(|i| {
245 (
246 (
247 Row::pack_slice(&[Datum::Int64(i)]),
248 Row::pack_slice(&[Datum::String("a repeated string value")]),
249 ),
250 time,
251 1i64,
252 )
253 })
254 .collect();
255
256 // Cut the sorted run the way a merge batcher's chain is cut.
257 let columns = || -> Vec<ColumnBody<((Row, Row), Timestamp, i64)>> {
258 updates
259 .chunks(250)
260 .map(|part| {
261 let mut column: ColumnBody<((Row, Row), Timestamp, i64)> = Default::default();
262 for update in part {
263 column.push_into(update);
264 }
265 column
266 })
267 .collect()
268 };
269 let description = || {
270 Description::new(
271 Antichain::from_elem(time),
272 Antichain::new(),
273 Antichain::from_elem(time),
274 )
275 };
276
277 type Paged = crate::RowRowColPagedBuilder<Timestamp, i64>;
278 type Chunked = UnchunkBuilder<Paged, (Row, Row), Timestamp, i64>;
279
280 let mut chain = columns();
281 let from_columns = <Paged as Builder>::seal(&mut chain, description());
282
283 let mut chain: Vec<ColumnChunk<(Row, Row), Timestamp, i64>> =
284 columns().into_iter().map(ColumnChunk::from_body).collect();
285 let from_chunks = <Chunked as Builder>::seal(&mut chain, description());
286
287 assert!(
288 from_chunks.0.storage.keys.has_codec() && from_chunks.0.storage.vals.vals.has_codec(),
289 "the sealed batch carries the chain's codecs"
290 );
291
292 // Same updates and same cuts, so the chunked seal has to reach the
293 // encoded size the column seal does.
294 let heap = |c: &DatumContainer| {
295 let mut size = 0;
296 c.heap_size(|_, cap| size += cap);
297 size
298 };
299 let encoded = |b: &ArcBatch<OrdValBatch<RowRowLayout<((Row, Row), Timestamp, i64)>>>| {
300 heap(&b.0.storage.keys) + heap(&b.0.storage.vals.vals)
301 };
302
303 assert!(
304 encoded(&from_chunks) <= encoded(&from_columns),
305 "chunked seal should compress like the column seal: chunks={} columns={}",
306 encoded(&from_chunks),
307 encoded(&from_columns),
308 );
309 }
310
311 #[mz_ore::test]
312 #[cfg_attr(miri, ignore)] // unsupported operation: integer-to-pointer casts and `ptr::with_exposed_provenance` are not supported
313 fn test_round_trip() {
314 fn round_trip(datums: Vec<Datum>) {
315 let row = Row::pack(datums.clone());
316
317 let mut container = DatumContainer::with_capacity(row.byte_len());
318 container.push_own(&row);
319
320 // When run under miri this catches undefined bytes written to data
321 // eg by calling push_copy! on a type which contains undefined padding values
322 println!("{:?}", container.index(0).iter.data);
323
324 let datums2 = container.index(0).collect::<Vec<_>>();
325 assert_eq!(datums, datums2);
326 }
327
328 round_trip(vec![]);
329 round_trip(
330 SqlScalarType::enumerate()
331 .iter()
332 .flat_map(|r#type| r#type.interesting_datums())
333 .collect(),
334 );
335 round_trip(vec![
336 Datum::Null,
337 Datum::Null,
338 Datum::False,
339 Datum::True,
340 Datum::Int16(-21),
341 Datum::Int32(-42),
342 Datum::Int64(-2_147_483_648 - 42),
343 Datum::UInt8(0),
344 Datum::UInt8(1),
345 Datum::UInt16(0),
346 Datum::UInt16(1),
347 Datum::UInt16(1 << 8),
348 Datum::UInt32(0),
349 Datum::UInt32(1),
350 Datum::UInt32(1 << 8),
351 Datum::UInt32(1 << 16),
352 Datum::UInt32(1 << 24),
353 Datum::UInt64(0),
354 Datum::UInt64(1),
355 Datum::UInt64(1 << 8),
356 Datum::UInt64(1 << 16),
357 Datum::UInt64(1 << 24),
358 Datum::UInt64(1 << 32),
359 Datum::UInt64(1 << 40),
360 Datum::UInt64(1 << 48),
361 Datum::UInt64(1 << 56),
362 Datum::Date(Date::from_pg_epoch(365 * 45 + 21).unwrap()),
363 Datum::Interval(Interval {
364 months: 312,
365 ..Default::default()
366 }),
367 Datum::Interval(Interval::new(0, 0, 1_012_312)),
368 Datum::Bytes(&[]),
369 Datum::Bytes(&[0, 2, 1, 255]),
370 Datum::String(""),
371 Datum::String("العَرَبِيَّة"),
372 ]);
373 }
374
375 /// Exercises the *compressed* encode→decode paths, which the dyncfg-gated
376 /// `test_round_trip` never reaches (it installs no codec). We drive the codec
377 /// directly: observe a sample, build a codec via both `new_from([c1, c2])`
378 /// (the merge path) and `new_safe` (the safe-tag path), then round-trip every
379 /// row through it. We additionally assert the dictionary actually engaged, so
380 /// the test keeps covering the compressed branch rather than silently
381 /// degrading to raw fall-through.
382 #[mz_ore::test]
383 #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
384 fn test_codec_round_trip() {
385 use crate::row_codec::ColumnsCodec;
386
387 // Rows with a small set of repeated, multi-byte string values, so the
388 // dictionary installs entries (MisraGries keeps values with len > 1 and
389 // count > 1). Mixing in an integer column exercises the raw fall-through
390 // (and thus the new soundness `debug_assert`) alongside dictionary hits.
391 let values = ["apple", "banana", "cherry"];
392 let rows: Vec<Row> = (0..3_000)
393 .map(|i| {
394 Row::pack_slice(&[
395 Datum::String(values[i % values.len()]),
396 Datum::Int64(i64::try_from(i).unwrap()),
397 Datum::String(values[(i / 7) % values.len()]),
398 ])
399 })
400 .collect();
401
402 // Accumulate statistics in two independent observers, so the merge in
403 // `new_from([&stats1, &stats2])` is actually exercised.
404 let mut stats1 = ColumnsCodec::default();
405 let mut stats2 = ColumnsCodec::default();
406 let mut scratch = Vec::new();
407 for (i, row) in rows.iter().enumerate() {
408 scratch.clear();
409 let stats = if i % 2 == 0 { &mut stats1 } else { &mut stats2 };
410 stats.encode(ColumnsCodec::borrow_row(row), &mut scratch);
411 }
412
413 let merged = ColumnsCodec::new_from([&stats1, &stats2]);
414 let safe = stats1.new_safe();
415 for mut codec in [merged, safe] {
416 let mut compressed_any = false;
417 for row in &rows {
418 let mut buf = Vec::new();
419 codec.encode(ColumnsCodec::borrow_row(row), &mut buf);
420
421 let decoded = codec.decode(&buf).collect::<Vec<_>>();
422 let expected = ColumnsCodec::borrow_row(row).collect::<Vec<_>>();
423 assert_eq!(decoded, expected, "round-trip mismatch for {row:?}");
424
425 compressed_any |= buf.len() < row.data().len();
426 }
427 assert!(
428 compressed_any,
429 "dictionary never engaged; test no longer covers the compressed path",
430 );
431 }
432 }
433
434 /// Regression test for a dictionary-codec soundness bug in the safe-install
435 /// path (`new_safe`), reachable with the paged batcher enabled.
436 ///
437 /// A from-scratch container stores its pre-install rows *raw* while gathering
438 /// statistics, then installs a *safe* codec via `new_safe`. `new_safe` used to
439 /// discard the first-byte bitmap gathered over those raw rows. That bitmap is
440 /// soundness-critical: a later `new_from` merge consults it to decide which
441 /// one-byte tags are free to hand out as dictionary keys. With the bitmap
442 /// dropped, the merge could assign a dictionary tag equal to a raw datum's
443 /// first byte, after which `decode` resolves that literal datum to the
444 /// dictionary entry — returning the wrong value.
445 ///
446 /// We drive the lifecycle directly: observe short strings (first byte
447 /// `StringTiny`) into the pre-install statistics, install a safe codec, then
448 /// feed it many distinct *long* strings (first byte `StringShort`)
449 /// post-install so the merge has heavy hitters to compress. Merging via
450 /// `new_from` and re-encoding the short strings then exercises the raw
451 /// fall-through whose first byte the merge must not have claimed as a tag.
452 /// Before the fix the `StringTiny` tag was handed out and the round-trip
453 /// produced a long string (and tripped `encode`'s soundness `debug_assert`).
454 #[mz_ore::test]
455 #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
456 fn test_safe_codec_merge_bitmap_carryover() {
457 use crate::row_codec::ColumnsCodec;
458
459 // Short strings: length < 256, so they encode with the `StringTiny` tag.
460 // Unique, so MisraGries never makes them dictionary entries; they always
461 // fall through raw, exposing their first byte.
462 let short_rows: Vec<Row> = (0..256)
463 .map(|i| Row::pack_slice(&[Datum::String(&format!("s{i}"))]))
464 .collect();
465 // Long strings: length >= 256, so they encode with the `StringShort` tag —
466 // a *different* first byte than the short strings. Distinct values, each
467 // repeated, so the post-install codec accrues many heavy hitters and the
468 // merge assigns dictionary tags across the low byte range, reaching the
469 // short strings' `StringTiny` tag unless the bitmap reserves it.
470 let long_values: Vec<String> = (0..64).map(|i| format!("{i:0>300}")).collect();
471
472 // Pre-install statistics observe only the short strings' first bytes.
473 let mut stats = ColumnsCodec::default();
474 let mut scratch = Vec::new();
475 for row in &short_rows {
476 scratch.clear();
477 stats.encode(ColumnsCodec::borrow_row(row), &mut scratch);
478 }
479
480 // Install a safe codec, then feed it the long strings post-install so it
481 // accrues heavy hitters (and observes only the `StringShort` first byte).
482 let mut safe = stats.new_safe();
483 for _ in 0..8 {
484 for v in &long_values {
485 let row = Row::pack_slice(&[Datum::String(v)]);
486 scratch.clear();
487 safe.encode(ColumnsCodec::borrow_row(&row), &mut scratch);
488 }
489 }
490
491 // Merge, then round-trip the short strings. With the bitmap carried over,
492 // no dictionary tag collides with the short strings' first byte; without
493 // it, one does.
494 let mut merged = ColumnsCodec::new_from([&safe]);
495 for row in &short_rows {
496 let mut buf = Vec::new();
497 merged.encode(ColumnsCodec::borrow_row(row), &mut buf);
498 let decoded = merged.decode(&buf).collect::<Vec<_>>();
499 let expected = ColumnsCodec::borrow_row(row).collect::<Vec<_>>();
500 assert_eq!(decoded, expected, "round-trip mismatch for {row:?}");
501 }
502 }
503
504 /// Confirms the structural assumption underpinning `SAFE_TAG_BASE`: every
505 /// datum the row format produces encodes with a first byte strictly less
506 /// than `SAFE_TAG_BASE`. If `mz_repr` ever introduces a tag that crosses
507 /// the boundary, `DictionaryCodec::new_safe` would assign a dictionary tag
508 /// that collides with a literal datum first-byte, breaking decoding.
509 #[mz_ore::test]
510 fn test_safe_tag_base() {
511 use crate::row_codec::SAFE_TAG_BASE;
512 let check = |datum: Datum| {
513 let row = Row::pack_slice(&[datum]);
514 let data = row.data();
515 assert!(!data.is_empty(), "empty encoding for {datum:?}");
516 assert!(
517 data[0] < SAFE_TAG_BASE,
518 "datum {datum:?} encodes with first byte {} >= SAFE_TAG_BASE ({}); \
519 a new row tag has crossed the safe boundary",
520 data[0],
521 SAFE_TAG_BASE,
522 );
523 };
524 for ty in SqlScalarType::enumerate().iter() {
525 for datum in ty.interesting_datums() {
526 check(datum);
527 }
528 }
529 }
530
531 /// A batch built via the builder's `push`/`done` path (as the `reduce` operator
532 /// does) that stays under `STATS_THRESHOLD` never installs a codec at build time.
533 /// `done` now promotes the gathered statistics into the codec slot, so the batch
534 /// carries a codec + heavy-hitter summary and does not poison a later merge.
535 ///
536 /// This drives that container lifecycle directly: gather raw (well under the
537 /// threshold), promote at "done", then merge two such containers the way a spine
538 /// compaction does. With promotion the merge takes the `new_from` path and
539 /// compresses; without it both inputs are codec-less and the merge stays raw.
540 /// Every merged row must still round-trip.
541 #[mz_ore::test]
542 #[cfg_attr(miri, ignore)] // integer-to-pointer casts in row decoding are unsupported under miri
543 fn push_done_promotion_avoids_merge_poison() {
544 use std::sync::atomic::Ordering;
545 use timely::container::PushInto;
546
547 // Gate the dictionary path on. Safe for other tests: the flag only controls
548 // whether `DatumContainer` gathers stats; it never changes decode results.
549 crate::DICTIONARY_COMPRESSION.store(true, Ordering::Relaxed);
550
551 // Low-cardinality rows, well under `STATS_THRESHOLD` (64Ki): a repeated
552 // multi-byte string the dictionary compresses, plus an integer column that
553 // exercises raw fall-through.
554 let rows: Vec<Row> = (0..2_000i64)
555 .map(|i| {
556 Row::pack_slice(&[
557 Datum::Int64(i % 8),
558 Datum::String("a repeated string value"),
559 ])
560 })
561 .collect();
562
563 // Build a container the way the push/done path does: gather raw without ever
564 // crossing `STATS_THRESHOLD`, optionally promoting at "done".
565 let build = |promote: bool| {
566 let mut c = DatumContainer::with_capacity(rows.len());
567 for row in &rows {
568 c.push_into(row);
569 }
570 if promote {
571 c.promote_stats_to_codec();
572 }
573 c
574 };
575
576 // Merge two containers as a spine compaction does: allocate via
577 // `merge_capacity`, then copy every row through.
578 let merge = |a: &DatumContainer, b: &DatumContainer| {
579 let mut m = DatumContainer::merge_capacity(a, b);
580 for i in 0..a.len() {
581 m.push_into(a.index(i));
582 }
583 for i in 0..b.len() {
584 m.push_into(b.index(i));
585 }
586 m
587 };
588
589 let heap = |c: &DatumContainer| {
590 let mut size = 0;
591 c.heap_size(|_, cap| size += cap);
592 size
593 };
594
595 // Codec-less inputs (no promotion): the merge cannot `new_from` and stays raw.
596 let poisoned = merge(&build(false), &build(false));
597 // Promoted inputs carry a codec + summary: the merge `new_from`s and compresses.
598 let compressed = merge(&build(true), &build(true));
599
600 // Round-trip: every merged row decodes back to the corresponding input row
601 // (the merge here concatenates a's rows then b's rows, no consolidation).
602 assert_eq!(compressed.len(), rows.len() * 2);
603 for i in 0..compressed.len() {
604 let got = compressed.index(i).collect::<Vec<_>>();
605 let want = rows[i % rows.len()].iter().collect::<Vec<_>>();
606 assert_eq!(got, want, "merged row {i} round-trips");
607 }
608
609 // The promoted merge must actually compress relative to the poisoned one,
610 // confirming promotion carried a usable summary into `new_from`.
611 assert!(
612 heap(&compressed) < heap(&poisoned),
613 "promotion should let the merge compress: compressed={} poisoned={}",
614 heap(&compressed),
615 heap(&poisoned),
616 );
617 }
618}
619
620/// A `[u8]`-specialized container.
621mod bytes_container {
622
623 use differential_dataflow::trace::implementations::BatchContainer;
624 use timely::container::PushInto;
625
626 use mz_ore::region::Region;
627
628 /// A slice container with four bytes overhead per slice.
629 pub struct BytesContainer {
630 /// Total length of `batches`, maintained because recomputation is expensive.
631 length: usize,
632 batches: Vec<BytesBatch>,
633 }
634
635 impl BytesContainer {
636 /// Visit contained allocations to determine their size and capacity.
637 #[inline]
638 pub fn heap_size(&self, mut callback: impl FnMut(usize, usize)) {
639 // Calculate heap size for local, stash, and stash entries
640 callback(
641 self.batches.len() * std::mem::size_of::<BytesBatch>(),
642 self.batches.capacity() * std::mem::size_of::<BytesBatch>(),
643 );
644 for batch in self.batches.iter() {
645 batch.offsets.heap_size(&mut callback);
646 callback(batch.storage.len(), batch.storage.capacity());
647 }
648 }
649 }
650
651 impl BatchContainer for BytesContainer {
652 type Owned = Vec<u8>;
653 type ReadItem<'a> = &'a [u8];
654
655 #[inline]
656 fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
657 item.to_vec()
658 }
659
660 #[inline]
661 fn clone_onto<'a>(item: Self::ReadItem<'a>, other: &mut Self::Owned) {
662 other.clear();
663 other.extend_from_slice(item);
664 }
665
666 #[inline(always)]
667 fn push_ref(&mut self, item: Self::ReadItem<'_>) {
668 self.push_into(item);
669 }
670
671 #[inline(always)]
672 fn push_own(&mut self, item: &Self::Owned) {
673 self.push_into(item.as_slice())
674 }
675
676 fn clear(&mut self) {
677 self.batches.clear();
678 self.batches.push(BytesBatch::with_capacities(0, 0));
679 self.length = 0;
680 }
681
682 fn with_capacity(size: usize) -> Self {
683 Self {
684 length: 0,
685 batches: vec![BytesBatch::with_capacities(size, size)],
686 }
687 }
688
689 fn merge_capacity(cont1: &Self, cont2: &Self) -> Self {
690 let mut item_cap = 1;
691 let mut byte_cap = 0;
692 for batch in cont1.batches.iter() {
693 item_cap += batch.offsets.len() - 1;
694 byte_cap += batch.storage.len();
695 }
696 for batch in cont2.batches.iter() {
697 item_cap += batch.offsets.len() - 1;
698 byte_cap += batch.storage.len();
699 }
700 Self {
701 length: 0,
702 batches: vec![BytesBatch::with_capacities(item_cap, byte_cap)],
703 }
704 }
705
706 #[inline(always)]
707 fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
708 item
709 }
710
711 #[inline]
712 fn index(&self, mut index: usize) -> Self::ReadItem<'_> {
713 for batch in self.batches.iter() {
714 if index < batch.len() {
715 return batch.index(index);
716 }
717 index -= batch.len();
718 }
719 panic!("Index out of bounds");
720 }
721
722 #[inline(always)]
723 fn len(&self) -> usize {
724 self.length
725 }
726 }
727
728 impl PushInto<&[u8]> for BytesContainer {
729 #[inline]
730 fn push_into(&mut self, item: &[u8]) {
731 self.length += 1;
732 if let Some(batch) = self.batches.last_mut() {
733 let success = batch.try_push(item);
734 if !success {
735 // double the lengths from `batch`.
736 let item_cap = 2 * batch.offsets.len();
737 let byte_cap = std::cmp::max(2 * batch.storage.capacity(), item.len());
738 let mut new_batch = BytesBatch::with_capacities(item_cap, byte_cap);
739 assert!(new_batch.try_push(item));
740 self.batches.push(new_batch);
741 }
742 }
743 }
744 }
745
746 /// A batch of slice storage.
747 ///
748 /// The backing storage for this batch will not be resized.
749 pub struct BytesBatch {
750 offsets: crate::OffsetOptimized,
751 storage: Region<u8>,
752 len: usize,
753 }
754
755 impl BytesBatch {
756 /// Either accepts the slice and returns true,
757 /// or does not and returns false.
758 fn try_push(&mut self, slice: &[u8]) -> bool {
759 if self.storage.len() + slice.len() <= self.storage.capacity() {
760 self.storage.extend_from_slice(slice);
761 self.offsets.push_into(self.storage.len());
762 self.len += 1;
763 true
764 } else {
765 false
766 }
767 }
768 #[inline]
769 fn index(&self, index: usize) -> &[u8] {
770 let lower = self.offsets.index(index);
771 let upper = self.offsets.index(index + 1);
772 &self.storage[lower..upper]
773 }
774 #[inline(always)]
775 fn len(&self) -> usize {
776 mz_ore::soft_assert_eq_no_log!(self.len, self.offsets.len() - 1);
777 self.len
778 }
779
780 fn with_capacities(item_cap: usize, byte_cap: usize) -> Self {
781 // TODO: be wary of `byte_cap` greater than 2^32.
782 let mut offsets = crate::OffsetOptimized::with_capacity(item_cap + 1);
783 offsets.push_into(0);
784 Self {
785 offsets,
786 storage: Region::new_auto(byte_cap.next_power_of_two()),
787 len: 0,
788 }
789 }
790 }
791}
792
793mod offset_opt {
794 use differential_dataflow::trace::implementations::BatchContainer;
795 use differential_dataflow::trace::implementations::OffsetList;
796 use timely::container::PushInto;
797
798 enum OffsetStride {
799 Empty,
800 Zero,
801 Striding(usize, usize),
802 Saturated(usize, usize, usize),
803 }
804
805 impl OffsetStride {
806 /// Accepts or rejects a newly pushed element.
807 #[inline]
808 fn push(&mut self, item: usize) -> bool {
809 match self {
810 OffsetStride::Empty => {
811 if item == 0 {
812 *self = OffsetStride::Zero;
813 true
814 } else {
815 false
816 }
817 }
818 OffsetStride::Zero => {
819 *self = OffsetStride::Striding(item, 2);
820 true
821 }
822 OffsetStride::Striding(stride, count) => {
823 if item == *stride * *count {
824 *count += 1;
825 true
826 } else if item == *stride * (*count - 1) {
827 *self = OffsetStride::Saturated(*stride, *count, 1);
828 true
829 } else {
830 false
831 }
832 }
833 OffsetStride::Saturated(stride, count, reps) => {
834 if item == *stride * (*count - 1) {
835 *reps += 1;
836 true
837 } else {
838 false
839 }
840 }
841 }
842 }
843
844 #[inline]
845 fn index(&self, index: usize) -> usize {
846 match self {
847 OffsetStride::Empty => {
848 panic!("Empty OffsetStride")
849 }
850 OffsetStride::Zero => 0,
851 OffsetStride::Striding(stride, _steps) => *stride * index,
852 OffsetStride::Saturated(stride, steps, _reps) => {
853 if index < *steps {
854 *stride * index
855 } else {
856 *stride * (*steps - 1)
857 }
858 }
859 }
860 }
861
862 #[inline]
863 fn len(&self) -> usize {
864 match self {
865 OffsetStride::Empty => 0,
866 OffsetStride::Zero => 1,
867 OffsetStride::Striding(_stride, steps) => *steps,
868 OffsetStride::Saturated(_stride, steps, reps) => *steps + *reps,
869 }
870 }
871 }
872
873 pub struct OffsetOptimized {
874 strided: OffsetStride,
875 spilled: OffsetList,
876 }
877
878 impl BatchContainer for OffsetOptimized {
879 type Owned = usize;
880 type ReadItem<'a> = usize;
881
882 #[inline]
883 fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
884 item
885 }
886
887 #[inline]
888 fn push_ref(&mut self, item: Self::ReadItem<'_>) {
889 self.push_into(item)
890 }
891
892 #[inline]
893 fn push_own(&mut self, item: &Self::Owned) {
894 self.push_into(*item)
895 }
896
897 fn clear(&mut self) {
898 self.strided = OffsetStride::Empty;
899 self.spilled.clear();
900 }
901
902 fn with_capacity(_size: usize) -> Self {
903 Self {
904 strided: OffsetStride::Empty,
905 spilled: OffsetList::with_capacity(0),
906 }
907 }
908
909 fn merge_capacity(_cont1: &Self, _cont2: &Self) -> Self {
910 Self {
911 strided: OffsetStride::Empty,
912 spilled: OffsetList::with_capacity(0),
913 }
914 }
915
916 #[inline]
917 fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
918 item
919 }
920
921 #[inline]
922 fn index(&self, index: usize) -> Self::ReadItem<'_> {
923 if index < self.strided.len() {
924 self.strided.index(index)
925 } else {
926 self.spilled.index(index - self.strided.len())
927 }
928 }
929
930 #[inline]
931 fn len(&self) -> usize {
932 self.strided.len() + self.spilled.len()
933 }
934 }
935
936 impl PushInto<usize> for OffsetOptimized {
937 #[inline]
938 fn push_into(&mut self, item: usize) {
939 if !self.spilled.is_empty() {
940 self.spilled.push(item);
941 } else {
942 let inserted = self.strided.push(item);
943 if !inserted {
944 self.spilled.push(item);
945 }
946 }
947 }
948 }
949
950 impl OffsetOptimized {
951 pub fn heap_size(&self, callback: impl FnMut(usize, usize)) {
952 crate::offset_list_size(&self.spilled, callback);
953 }
954 }
955}
956
957/// Helper to compute the size of an [`OffsetList`] in memory.
958#[inline]
959pub(crate) fn offset_list_size(data: &OffsetList, mut callback: impl FnMut(usize, usize)) {
960 // Private `vec_size` because we should only use it where data isn't region-allocated.
961 // `T: Copy` makes sure the implementation is correct even if types change!
962 #[inline(always)]
963 fn vec_size<T: Copy>(data: &Vec<T>, mut callback: impl FnMut(usize, usize)) {
964 let size_of_t = std::mem::size_of::<T>();
965 callback(data.len() * size_of_t, data.capacity() * size_of_t);
966 }
967
968 vec_size(&data.smol, &mut callback);
969 vec_size(&data.chonk, callback);
970}
971
972/// A `Row`-specialized container using dictionary compression.
973///
974/// The approach is to establish for each column lists of common values, and to use "unoccupied"
975/// tags in the row encoding (e.g. where we would indicate types) to replace these common values.
976/// This substitution is opt-in, in that we don't need to do it, and in particular do not do it
977/// while we are collecting preliminary information about common values, and then start to use it
978/// once we believe we have enough information. Once we have started to use the substitutions we
979/// cannot change the meaning of a reserved byte pattern, for the container we are populating.
980///
981/// Each from-scratch container observes `STATS_THRESHOLD` records before establishing a mapping
982/// from spare tags to common values. Containers that are formed from merging other containers
983/// use those input containers' common values to populate a codec and use it immediately.
984///
985/// The dictionary behavior is controlled by the `DICTIONARY_COMPRESSION` flag, which if disabled
986/// prevents the construction of codecs, which when absent simply cause the wrapper to behave as
987/// a no-op that fails to use any spare tags for common values. The flag is set once, when a
988/// replica is created (from compute's `InstanceConfig::arrangement_dictionary_compression`, itself
989/// captured from the `enable_arrangement_dictionary_compression_alpha` dyncfg at that moment), and is
990/// not changed for the life of the process; flipping the dyncfg only affects replicas created
991/// afterwards. Even with the flag fixed, a single replica can hold a mix of compressed and
992/// uncompressed containers — e.g. containers that never observed enough records to install a
993/// codec, or that were merged from uncompressed inputs.
994mod dictionary {
995
996 use differential_dataflow::trace::implementations::BatchContainer;
997
998 use mz_repr::{Row, RowRef};
999
1000 use super::row_codec::{ColumnsCodec, ColumnsIter};
1001
1002 /// Wrapper types that exist to support the creation of dictionary codecs.
1003 ///
1004 /// These types interpose at the seal() call, to traverse the data that is being sealed and
1005 /// then construct codecs that are used to encode the row-shaped keys and values. There are
1006 /// several variants, corresponding to the RowRow, RowVal, and Row-only spine types.
1007 pub mod builders {
1008
1009 use columnar::{Columnar, Index};
1010 use columnation::Columnation;
1011 use differential_dataflow::difference::Semigroup;
1012 use differential_dataflow::lattice::Lattice;
1013 use differential_dataflow::trace::Builder;
1014 use differential_dataflow::trace::Description;
1015 use differential_dataflow::trace::implementations::BatchContainer;
1016 use differential_dataflow::trace::implementations::ord_neu::{OrdKeyBatch, OrdKeyBuilder};
1017 use differential_dataflow::trace::implementations::ord_neu::{OrdValBatch, OrdValBuilder};
1018 use mz_timely_util::columnar::Column;
1019 use mz_timely_util::columnar::body::ColumnBody;
1020 use mz_timely_util::columnar::chunk::ChainState;
1021 use mz_timely_util::columnation::ColumnationStack as TimelyStack;
1022 use timely::progress::Timestamp;
1023
1024 use mz_repr::{Row, RowRef};
1025
1026 use super::super::row_codec::ColumnsCodec;
1027 use super::{DatumContainer, DatumSeq};
1028 use crate::DICTIONARY_COMPRESSION;
1029 use crate::spines::{RowLayout, RowRowLayout, RowValLayout, ValRowLayout};
1030
1031 /// Gather encoding statistics across `rows` and produce a codec from them.
1032 ///
1033 /// Accepts anything that borrows as a [`RowRef`], so it serves both the
1034 /// columnation-fed builders (which yield `&Row`) and the paged builders
1035 /// (which yield `&RowRef` straight out of a [`Column`] chunk).
1036 ///
1037 /// Returns `None` when dictionary compression is disabled.
1038 fn build_codec<'a, B>(rows: impl IntoIterator<Item = &'a B>) -> Option<ColumnsCodec>
1039 where
1040 B: std::borrow::Borrow<RowRef> + ?Sized + 'a,
1041 {
1042 if !DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed) {
1043 return None;
1044 }
1045 let mut stats = ColumnsCodec::default();
1046 observe_rows(&mut stats, rows);
1047 Some(ColumnsCodec::new_from([&stats]))
1048 }
1049
1050 /// Fold `rows` into `stats`, the accumulator a codec is built from.
1051 fn observe_rows<'a, B>(stats: &mut ColumnsCodec, rows: impl IntoIterator<Item = &'a B>)
1052 where
1053 B: std::borrow::Borrow<RowRef> + ?Sized + 'a,
1054 {
1055 for row in rows {
1056 let row = row.borrow();
1057 if !row.is_empty() {
1058 // Gather stats only; the encoded output would be thrown away here, so
1059 // `observe` skips the per-value lookup and the throwaway-buffer memcpy
1060 // that `encode` would do (see `ColumnsCodec::observe`).
1061 stats.observe(DatumSeq::borrow_as(row).bytes_iter());
1062 }
1063 }
1064 }
1065
1066 pub struct RowRowBuilder<
1067 T: Lattice + Timestamp + Columnation,
1068 R: Ord + Semigroup + Columnation + 'static,
1069 > {
1070 inner: OrdValBuilder<RowRowLayout<((Row, Row), T, R)>, TimelyStack<((Row, Row), T, R)>>,
1071 }
1072
1073 impl<T: Lattice + Timestamp + Columnation, R: Ord + Semigroup + Columnation + 'static>
1074 Builder for RowRowBuilder<T, R>
1075 {
1076 type Input = TimelyStack<((Row, Row), T, R)>;
1077 type Time = T;
1078 type Output = OrdValBatch<RowRowLayout<((Row, Row), T, R)>>;
1079
1080 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1081 Self {
1082 inner: Builder::with_capacity(keys, vals, upds),
1083 }
1084 }
1085 fn push(&mut self, chunk: &mut Self::Input) {
1086 self.inner.push(chunk)
1087 }
1088 fn done(self, description: Description<Self::Time>) -> Self::Output {
1089 // The push/done build path (e.g. the `reduce` operator, which builds
1090 // batches with `Builder::new()` + `push` + `done` rather than `seal`)
1091 // never runs `seal`'s codec install. Install a codec here from the
1092 // statistics gathered during `push`, mirroring `seal` — but without
1093 // building a dictionary or re-encoding the rows; see
1094 // `DatumContainer::promote_stats_to_codec` for why a codec-less batch
1095 // must be avoided even though its rows stay raw.
1096 let mut inner = self.inner;
1097 inner.result.keys.promote_stats_to_codec();
1098 inner.result.vals.vals.promote_stats_to_codec();
1099 inner.done(description)
1100 }
1101 fn seal(
1102 chain: &mut Vec<Self::Input>,
1103 description: Description<Self::Time>,
1104 ) -> Self::Output {
1105 let key_codec = build_codec(
1106 chain
1107 .iter()
1108 .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1109 );
1110 let val_codec = build_codec(
1111 chain
1112 .iter()
1113 .flat_map(|link| link.iter().map(|((_, v), _, _)| v)),
1114 );
1115
1116 use differential_dataflow::trace::implementations::BuilderInput;
1117
1118 let (keys, vals, upds) = <Self::Input as BuilderInput<
1119 DatumContainer,
1120 DatumContainer,
1121 >>::key_val_upd_counts(&chain[..]);
1122 let mut builder = Self::with_capacity(keys, vals, upds);
1123 // The seal path installs a codec directly, so the per-container stats
1124 // gatherer (which `with_capacity` may have allocated) is dead weight and
1125 // would contradict the `stats: None once codec installed` invariant.
1126 builder.inner.result.keys.codec = key_codec;
1127 builder.inner.result.keys.stats = None;
1128 builder.inner.result.vals.vals.codec = val_codec;
1129 builder.inner.result.vals.vals.stats = None;
1130
1131 for mut chunk in chain.drain(..) {
1132 builder.push(&mut chunk);
1133 }
1134
1135 builder.done(description)
1136 }
1137 }
1138
1139 pub struct RowValBuilder<
1140 V: Ord + Clone + Columnation + 'static,
1141 T: Lattice + Timestamp + Columnation,
1142 R: Ord + Semigroup + Columnation + 'static,
1143 > {
1144 inner: OrdValBuilder<RowValLayout<((Row, V), T, R)>, TimelyStack<((Row, V), T, R)>>,
1145 }
1146
1147 impl<
1148 V: Ord + Clone + Columnation,
1149 T: Lattice + Timestamp + Columnation,
1150 R: Ord + Semigroup + Columnation + 'static,
1151 > Builder for RowValBuilder<V, T, R>
1152 {
1153 type Input = TimelyStack<((Row, V), T, R)>;
1154 type Time = T;
1155 type Output = OrdValBatch<RowValLayout<((Row, V), T, R)>>;
1156
1157 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1158 Self {
1159 inner: Builder::with_capacity(keys, vals, upds),
1160 }
1161 }
1162 fn push(&mut self, chunk: &mut Self::Input) {
1163 self.inner.push(chunk)
1164 }
1165 fn done(self, description: Description<Self::Time>) -> Self::Output {
1166 // See `RowRowBuilder::done`: install a codec on the `Row`-shaped key
1167 // container for the push/done (e.g. `reduce`) path that skips `seal`.
1168 let mut inner = self.inner;
1169 inner.result.keys.promote_stats_to_codec();
1170 inner.done(description)
1171 }
1172 fn seal(
1173 chain: &mut Vec<Self::Input>,
1174 description: Description<Self::Time>,
1175 ) -> Self::Output {
1176 let key_codec = build_codec(
1177 chain
1178 .iter()
1179 .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1180 );
1181
1182 use differential_dataflow::trace::implementations::BuilderInput;
1183
1184 let (keys, vals, upds) = <Self::Input as BuilderInput<
1185 DatumContainer,
1186 TimelyStack<V>,
1187 >>::key_val_upd_counts(&chain[..]);
1188 let mut builder = Self::with_capacity(keys, vals, upds);
1189 // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1190 builder.inner.result.keys.codec = key_codec;
1191 builder.inner.result.keys.stats = None;
1192
1193 for mut chunk in chain.drain(..) {
1194 builder.push(&mut chunk);
1195 }
1196
1197 builder.done(description)
1198 }
1199 }
1200
1201 pub struct RowBuilder<
1202 T: Lattice + Timestamp + Columnation,
1203 R: Ord + Semigroup + Columnation + 'static,
1204 DC: BatchContainer<Owned = R> = TimelyStack<R>,
1205 > {
1206 inner: OrdKeyBuilder<RowLayout<((Row, ()), T, R), DC>, TimelyStack<((Row, ()), T, R)>>,
1207 }
1208
1209 impl<T, R, DC> Builder for RowBuilder<T, R, DC>
1210 where
1211 T: Lattice + Timestamp + Columnation,
1212 R: Ord + Semigroup + Columnation + 'static,
1213 DC: BatchContainer<Owned = R>,
1214 {
1215 type Input = TimelyStack<((Row, ()), T, R)>;
1216 type Time = T;
1217 type Output = OrdKeyBatch<RowLayout<((Row, ()), T, R), DC>>;
1218
1219 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1220 Self {
1221 inner: Builder::with_capacity(keys, vals, upds),
1222 }
1223 }
1224 fn push(&mut self, chunk: &mut Self::Input) {
1225 self.inner.push(chunk)
1226 }
1227 fn done(self, description: Description<Self::Time>) -> Self::Output {
1228 // See `RowRowBuilder::done`: install a codec on the `Row`-shaped key
1229 // container for the push/done (e.g. `reduce`) path that skips `seal`.
1230 let mut inner = self.inner;
1231 inner.result.keys.promote_stats_to_codec();
1232 inner.done(description)
1233 }
1234 fn seal(
1235 chain: &mut Vec<Self::Input>,
1236 description: Description<Self::Time>,
1237 ) -> Self::Output {
1238 let key_codec = build_codec(
1239 chain
1240 .iter()
1241 .flat_map(|link| link.iter().map(|((k, _), _, _)| k)),
1242 );
1243
1244 use differential_dataflow::trace::implementations::BuilderInput;
1245
1246 let (keys, vals, upds) = <Self::Input as BuilderInput<
1247 DatumContainer,
1248 TimelyStack<()>,
1249 >>::key_val_upd_counts(&chain[..]);
1250 let mut builder = Self::with_capacity(keys, vals, upds);
1251 // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1252 builder.inner.result.keys.codec = key_codec;
1253 builder.inner.result.keys.stats = None;
1254
1255 for mut chunk in chain.drain(..) {
1256 builder.push(&mut chunk);
1257 }
1258
1259 builder.done(description)
1260 }
1261 }
1262
1263 /// Mirror of [`RowValBuilder`] with the roles swapped: arbitrary keys and
1264 /// `Row` *values*, so the dictionary codec is built for and installed on the
1265 /// value container.
1266 pub struct ValRowBuilder<
1267 K: Ord + Clone + Columnation + 'static,
1268 T: Lattice + Timestamp + Columnation,
1269 R: Ord + Semigroup + Columnation + 'static,
1270 > {
1271 inner: OrdValBuilder<ValRowLayout<((K, Row), T, R)>, TimelyStack<((K, Row), T, R)>>,
1272 }
1273
1274 impl<
1275 K: Ord + Clone + Columnation,
1276 T: Lattice + Timestamp + Columnation,
1277 R: Ord + Semigroup + Columnation + 'static,
1278 > Builder for ValRowBuilder<K, T, R>
1279 {
1280 type Input = TimelyStack<((K, Row), T, R)>;
1281 type Time = T;
1282 type Output = OrdValBatch<ValRowLayout<((K, Row), T, R)>>;
1283
1284 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1285 Self {
1286 inner: Builder::with_capacity(keys, vals, upds),
1287 }
1288 }
1289 fn push(&mut self, chunk: &mut Self::Input) {
1290 self.inner.push(chunk)
1291 }
1292 fn done(self, description: Description<Self::Time>) -> Self::Output {
1293 // See `RowRowBuilder::done`: install a codec on the `Row`-shaped value
1294 // container for the push/done (e.g. `reduce`) path that skips `seal`.
1295 let mut inner = self.inner;
1296 inner.result.vals.vals.promote_stats_to_codec();
1297 inner.done(description)
1298 }
1299 fn seal(
1300 chain: &mut Vec<Self::Input>,
1301 description: Description<Self::Time>,
1302 ) -> Self::Output {
1303 let val_codec = build_codec(
1304 chain
1305 .iter()
1306 .flat_map(|link| link.iter().map(|((_, v), _, _)| v)),
1307 );
1308
1309 use differential_dataflow::trace::implementations::BuilderInput;
1310
1311 let (keys, vals, upds) = <Self::Input as BuilderInput<
1312 TimelyStack<K>,
1313 DatumContainer,
1314 >>::key_val_upd_counts(&chain[..]);
1315 let mut builder = Self::with_capacity(keys, vals, upds);
1316 // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1317 builder.inner.result.vals.vals.codec = val_codec;
1318 builder.inner.result.vals.vals.stats = None;
1319
1320 for mut chunk in chain.drain(..) {
1321 builder.push(&mut chunk);
1322 }
1323
1324 builder.done(description)
1325 }
1326 }
1327
1328 /// Counterpart of [`RowRowBuilder`] that consumes [`ColumnBody`] chunks
1329 /// instead of columnation stacks, whether or not the batcher that
1330 /// produced them spills. Mirrors `RowRowBuilder::seal`:
1331 /// it gathers key and value statistics from the sealed chain and
1332 /// installs codecs directly, then drops the per-container stats gatherer.
1333 pub struct RowRowColPagedBuilder<
1334 T: Lattice + Timestamp + Columnation + Columnar,
1335 R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1336 > {
1337 inner: OrdValBuilder<RowRowLayout<((Row, Row), T, R)>, ColumnBody<((Row, Row), T, R)>>,
1338 }
1339
1340 impl<
1341 T: Lattice + Timestamp + Columnation + Columnar,
1342 R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1343 > Builder for RowRowColPagedBuilder<T, R>
1344 {
1345 type Input = ColumnBody<((Row, Row), T, R)>;
1346 type Time = T;
1347 type Output = OrdValBatch<RowRowLayout<((Row, Row), T, R)>>;
1348
1349 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1350 Self {
1351 inner: Builder::with_capacity(keys, vals, upds),
1352 }
1353 }
1354 fn push(&mut self, chunk: &mut Self::Input) {
1355 self.inner.push(chunk)
1356 }
1357 fn done(self, description: Description<Self::Time>) -> Self::Output {
1358 // See `RowRowBuilder::done`: a builder driven by `push` and `done`
1359 // alone gets its codec from the statistics gathered during the
1360 // pushes, since a codec-less batch poisons the merges it joins.
1361 // Installing a codec here is a no-op once one is installed.
1362 let mut inner = self.inner;
1363 inner.result.keys.promote_stats_to_codec();
1364 inner.result.vals.vals.promote_stats_to_codec();
1365 inner.done(description)
1366 }
1367 fn seal(
1368 chain: &mut Vec<Self::Input>,
1369 description: Description<Self::Time>,
1370 ) -> Self::Output {
1371 // `into_index_iter` yields the value column's `Row`s as `&RowRef`,
1372 // which `build_codec` consumes directly.
1373 let key_codec = build_codec(
1374 chain
1375 .iter()
1376 .flat_map(|c| c.borrow().into_index_iter().map(|((k, _), _, _)| k)),
1377 );
1378 let val_codec = build_codec(
1379 chain
1380 .iter()
1381 .flat_map(|c| c.borrow().into_index_iter().map(|((_, v), _, _)| v)),
1382 );
1383
1384 use differential_dataflow::trace::implementations::BuilderInput;
1385
1386 let (keys, vals, upds) = <Self::Input as BuilderInput<
1387 DatumContainer,
1388 DatumContainer,
1389 >>::key_val_upd_counts(&chain[..]);
1390 let mut builder = Self::with_capacity(keys, vals, upds);
1391 // See `RowRowBuilder::seal`: install the codecs and drop the
1392 // now-redundant per-container stats gatherer.
1393 builder.inner.result.keys.codec = key_codec;
1394 builder.inner.result.keys.stats = None;
1395 builder.inner.result.vals.vals.codec = val_codec;
1396 builder.inner.result.vals.vals.stats = None;
1397
1398 for mut chunk in chain.drain(..) {
1399 builder.push(&mut chunk);
1400 }
1401
1402 builder.done(description)
1403 }
1404 }
1405
1406 /// The chain-wide state of [`RowRowColPagedBuilder`]: the statistics its
1407 /// key and value codecs are built from, and the container sizes the
1408 /// chain implies.
1409 ///
1410 /// [`RowRowColPagedBuilder`]: crate::RowRowColPagedBuilder
1411 #[derive(Default)]
1412 pub struct RowRowColPagedState {
1413 keys: ColumnsCodec,
1414 vals: ColumnsCodec,
1415 /// Whether `keys` and `vals` hold statistics from any body.
1416 gathered: bool,
1417 counts: (usize, usize, usize),
1418 }
1419
1420 impl<
1421 T: Lattice + Timestamp + Columnation + Columnar,
1422 R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1423 > ChainState for RowRowColPagedBuilder<T, R>
1424 {
1425 type State = RowRowColPagedState;
1426
1427 fn wants_bodies() -> bool {
1428 DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1429 }
1430
1431 fn observe(state: &mut Self::State, input: &Self::Input) {
1432 if Self::wants_bodies() {
1433 // `into_index_iter` yields the key and value columns' `Row`s
1434 // as `&RowRef`, which `observe_rows` consumes directly.
1435 observe_rows(
1436 &mut state.keys,
1437 input.borrow().into_index_iter().map(|((k, _), _, _)| k),
1438 );
1439 observe_rows(
1440 &mut state.vals,
1441 input.borrow().into_index_iter().map(|((_, v), _, _)| v),
1442 );
1443 state.gathered = true;
1444 }
1445
1446 use differential_dataflow::trace::implementations::BuilderInput;
1447 let (keys, vals, upds) = <Self::Input as BuilderInput<
1448 DatumContainer,
1449 DatumContainer,
1450 >>::key_val_upd_counts(
1451 std::slice::from_ref(input)
1452 );
1453 state.counts.0 += keys;
1454 state.counts.1 += vals;
1455 state.counts.2 += upds;
1456 }
1457
1458 fn observe_records(state: &mut Self::State, records: usize) {
1459 // Key and value counts need the bodies. Sizing the update
1460 // containers alone is what a record count buys.
1461 state.counts.2 += records;
1462 }
1463
1464 fn from_state(state: Self::State) -> Self {
1465 let mut builder =
1466 Self::with_capacity(state.counts.0, state.counts.1, state.counts.2);
1467 // Install from the statistics that were gathered, not from the
1468 // compression flag: the flag is process-global and a worker can
1469 // flip it mid-seal, and a codec built from no observations would
1470 // carry an empty summary and displace the gatherer that `done`
1471 // would otherwise promote.
1472 if state.gathered {
1473 // See `RowRowBuilder::seal`: install the codecs the chain's
1474 // statistics imply and drop the now-redundant per-container
1475 // stats gatherer.
1476 builder.inner.result.keys.codec = Some(ColumnsCodec::new_from([&state.keys]));
1477 builder.inner.result.keys.stats = None;
1478 builder.inner.result.vals.vals.codec =
1479 Some(ColumnsCodec::new_from([&state.vals]));
1480 builder.inner.result.vals.vals.stats = None;
1481 }
1482 builder
1483 }
1484 }
1485
1486 /// Paged counterpart of [`ValRowBuilder`] that consumes [`Column`]
1487 /// chunks. Keys are arbitrary `Columnar` values (not `Row`-shaped) and
1488 /// stay uncompressed; only the value container receives a codec.
1489 pub struct ValRowColPagedBuilder<
1490 K: Ord + Clone + Columnation + Columnar + 'static,
1491 T: Lattice + Timestamp + Columnation + Columnar,
1492 R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1493 > {
1494 inner: OrdValBuilder<ValRowLayout<((K, Row), T, R)>, Column<((K, Row), T, R)>>,
1495 }
1496
1497 impl<
1498 K: Ord + Clone + Columnation + Columnar + 'static,
1499 T: Lattice + Timestamp + Columnation + Columnar,
1500 R: Ord + Semigroup + Columnation + Columnar + Clone + 'static,
1501 > Builder for ValRowColPagedBuilder<K, T, R>
1502 where
1503 for<'a> columnar::Ref<'a, K>: Copy + Ord,
1504 for<'a, 'b> &'a K: PartialEq<columnar::Ref<'b, K>>,
1505 for<'a> TimelyStack<K>: timely::container::PushInto<columnar::Ref<'a, K>>,
1506 {
1507 type Input = Column<((K, Row), T, R)>;
1508 type Time = T;
1509 type Output = OrdValBatch<ValRowLayout<((K, Row), T, R)>>;
1510
1511 fn with_capacity(keys: usize, vals: usize, upds: usize) -> Self {
1512 Self {
1513 inner: Builder::with_capacity(keys, vals, upds),
1514 }
1515 }
1516 fn push(&mut self, chunk: &mut Self::Input) {
1517 self.inner.push(chunk)
1518 }
1519 fn done(self, description: Description<Self::Time>) -> Self::Output {
1520 self.inner.done(description)
1521 }
1522 fn seal(
1523 chain: &mut Vec<Self::Input>,
1524 description: Description<Self::Time>,
1525 ) -> Self::Output {
1526 let val_codec = build_codec(
1527 chain
1528 .iter()
1529 .flat_map(|c| c.borrow().into_index_iter().map(|((_, v), _, _)| v)),
1530 );
1531
1532 use differential_dataflow::trace::implementations::BuilderInput;
1533
1534 let (keys, vals, upds) = <Self::Input as BuilderInput<
1535 TimelyStack<K>,
1536 DatumContainer,
1537 >>::key_val_upd_counts(&chain[..]);
1538 let mut builder = Self::with_capacity(keys, vals, upds);
1539 // See `RowRowBuilder::seal`: drop the now-redundant stats gatherer.
1540 builder.inner.result.vals.vals.codec = val_codec;
1541 builder.inner.result.vals.vals.stats = None;
1542
1543 for mut chunk in chain.drain(..) {
1544 builder.push(&mut chunk);
1545 }
1546
1547 builder.done(description)
1548 }
1549 }
1550 }
1551
1552 pub struct DatumContainer {
1553 /// Encoder/decoder used to translate between row bytes and the stored bytes.
1554 /// `None` until enough pushes have been observed (or if compression is disabled).
1555 codec: Option<ColumnsCodec>,
1556 /// The stored, possibly-encoded, row bytes.
1557 inner: super::bytes_container::BytesContainer,
1558 /// Staging buffer for ingested `Row` types.
1559 staging: Vec<u8>,
1560 /// Statistics gatherer, used to build a safe codec after enough pushes.
1561 /// `None` once the codec has been installed or if compression is disabled.
1562 stats: Option<ColumnsCodec>,
1563 }
1564
1565 impl BatchContainer for DatumContainer {
1566 type Owned = Row;
1567 type ReadItem<'a> = DatumSeq<'a>;
1568
1569 fn with_capacity(size: usize) -> Self {
1570 let stats = if crate::DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1571 {
1572 Some(Default::default())
1573 } else {
1574 None
1575 };
1576
1577 Self {
1578 codec: None,
1579 inner: BatchContainer::with_capacity(size),
1580 staging: Vec::new(),
1581 stats,
1582 }
1583 }
1584 fn merge_capacity(cont1: &Self, cont2: &Self) -> Self {
1585 // We only build a merged codec when *both* inputs carry one. A codec is
1586 // sound only for the data whose tag usage it observed, so we cannot reuse
1587 // one side's codec to decode the other side's rows. When exactly one side
1588 // is compressed we conservatively produce an uncompressed container rather
1589 // than risk a tag collision; the merged container re-gathers stats and may
1590 // install a fresh codec later via the `STATS_THRESHOLD` path.
1591 let codec = match (&cont1.codec, &cont2.codec) {
1592 (Some(c1), Some(c2)) => Some(ColumnsCodec::new_from([c1, c2])),
1593 _ => None,
1594 };
1595
1596 Self {
1597 codec,
1598 inner: BatchContainer::merge_capacity(&cont1.inner, &cont2.inner),
1599 staging: Vec::new(),
1600 stats: None,
1601 }
1602 }
1603 #[inline]
1604 fn index(&self, index: usize) -> Self::ReadItem<'_> {
1605 let data = self.inner.index(index);
1606 let iter = if let Some(codec) = &self.codec {
1607 codec.decode(data)
1608 } else {
1609 // Safety: without a codec we only push rows or datumseqs into `self.inner`.
1610 // Each retrieved byte slice should be row-encoded data, as long as we have
1611 // not unset the codec in the interim.
1612 unsafe { ColumnsIter::without_codec(data) }
1613 };
1614 DatumSeq { iter }
1615 }
1616 #[inline(always)]
1617 fn len(&self) -> usize {
1618 self.inner.len()
1619 }
1620
1621 #[inline(always)]
1622 fn reborrow<'b, 'a: 'b>(item: Self::ReadItem<'a>) -> Self::ReadItem<'b> {
1623 item
1624 }
1625
1626 #[inline(always)]
1627 fn into_owned<'a>(item: Self::ReadItem<'a>) -> Self::Owned {
1628 // Fast path: unencoded data is already row-formatted bytes.
1629 if item.iter.index.is_none() {
1630 // SAFETY: `iter.data` is raw row-encoded bytes when there is no codec.
1631 return unsafe { Row::from_bytes_unchecked(item.iter.data) };
1632 }
1633 Row::pack(item)
1634 }
1635
1636 #[inline(always)]
1637 fn clone_onto<'a>(item: Self::ReadItem<'a>, other: &mut Self::Owned) {
1638 // Fast path: unencoded data is already row-formatted bytes.
1639 if item.iter.index.is_none() {
1640 let mut packer = other.packer();
1641 // SAFETY: `iter.data` is raw row-encoded bytes when there is no codec.
1642 unsafe { packer.extend_by_slice_unchecked(item.iter.data) };
1643 return;
1644 }
1645 other.packer().extend(item);
1646 }
1647
1648 #[inline(always)]
1649 fn push_ref(&mut self, item: Self::ReadItem<'_>) {
1650 // Fast path: both sides unencoded — push raw bytes directly.
1651 if self.codec.is_none() && self.stats.is_none() && item.iter.index.is_none() {
1652 self.inner.push_ref(item.iter.data);
1653 return;
1654 }
1655 self.push_into(item);
1656 }
1657
1658 #[inline(always)]
1659 fn push_own(&mut self, item: &Self::Owned) {
1660 // Fast path: container is unencoded — push raw row bytes directly.
1661 if self.codec.is_none() && self.stats.is_none() {
1662 self.inner.push_ref(item.data());
1663 return;
1664 }
1665 self.push_into(item);
1666 }
1667
1668 #[inline(always)]
1669 fn clear(&mut self) {
1670 self.inner.clear();
1671 self.staging.clear();
1672 // Reset to the same state as a fresh `with_capacity`: drop any installed
1673 // codec and restore stats gathering (if compression is enabled). Keeping a
1674 // now-empty codec would leave `codec.is_some()`, which permanently routes
1675 // pushes down the encode path with an empty dictionary and prevents the
1676 // `STATS_THRESHOLD` install logic from ever re-engaging compression.
1677 self.codec = None;
1678 self.stats = if crate::DICTIONARY_COMPRESSION.load(std::sync::atomic::Ordering::Relaxed)
1679 {
1680 Some(Default::default())
1681 } else {
1682 None
1683 };
1684 }
1685 }
1686
1687 impl DatumContainer {
1688 /// Visit contained allocations to determine their size and capacity.
1689 #[inline]
1690 pub fn heap_size(&self, mut callback: impl FnMut(usize, usize)) {
1691 self.inner.heap_size(&mut callback);
1692 // The staging buffer and the (possibly absent) codec and stats gatherer all
1693 // hold heap allocations that the bare `inner` accounting misses.
1694 callback(self.staging.len(), self.staging.capacity());
1695 if let Some(codec) = &self.codec {
1696 codec.heap_size(&mut callback);
1697 }
1698 if let Some(stats) = &self.stats {
1699 stats.heap_size(&mut callback);
1700 }
1701 }
1702
1703 /// Promote a gathered-but-uninstalled statistics summary into the codec slot.
1704 ///
1705 /// A container filled via the builder's `push`/`done` path — as the `reduce`
1706 /// operator does, building batches with `Builder::new()` + `push` + `done`
1707 /// rather than `seal` — gathers statistics on every push but never reaches
1708 /// `seal`'s codec install, and only crosses the mid-formation
1709 /// `STATS_THRESHOLD` install if it grows past it. A smaller such container
1710 /// would otherwise be finalized with no codec at all, even with the flag on.
1711 ///
1712 /// That is a problem not because this batch needs compressing — its rows are
1713 /// already stored raw and we deliberately do *not* re-encode them here — but
1714 /// because a codec-less batch poisons future merges: [`Self::merge_capacity`]
1715 /// keys off the presence of a codec, so a codec-less input forces the merged
1716 /// container onto the uncompressed path. Moving the gathered statistics into
1717 /// the codec slot leaves the batch carrying a codec whose retained heavy-hitter
1718 /// summary a later merge can rebuild from via `ColumnsCodec::new_from`, while
1719 /// installing no dictionary: the empty `decode` map resolves every stored
1720 /// (raw) column through the literal-datum fall-through, so reads stay correct.
1721 ///
1722 /// We move the summary as-is rather than building a dictionary via `new_safe`
1723 /// / `new_from` (which reset the summary): unlike `seal` and the mid-formation
1724 /// install, `done` has no further rows to re-observe, so a reset summary would
1725 /// leave the eventual merge nothing to rebuild from.
1726 pub(crate) fn promote_stats_to_codec(&mut self) {
1727 if self.codec.is_none() {
1728 self.codec = self.stats.take();
1729 }
1730 }
1731
1732 /// Whether a codec is installed. A container without one stores raw
1733 /// bytes and blocks compression in every merge it takes part in.
1734 #[cfg(test)]
1735 pub(crate) fn has_codec(&self) -> bool {
1736 self.codec.is_some()
1737 }
1738 }
1739
1740 use timely::container::PushInto;
1741 impl PushInto<Row> for DatumContainer {
1742 #[inline(always)]
1743 fn push_into(&mut self, item: Row) {
1744 self.push_into(&item);
1745 }
1746 }
1747
1748 impl PushInto<&Row> for DatumContainer {
1749 #[inline(always)]
1750 fn push_into(&mut self, item: &Row) {
1751 self.push_into(DatumSeq::borrow_as(item));
1752 }
1753 }
1754
1755 impl PushInto<&RowRef> for DatumContainer {
1756 #[inline(always)]
1757 fn push_into(&mut self, item: &RowRef) {
1758 self.push_into(DatumSeq::borrow_as(item));
1759 }
1760 }
1761
1762 /// Number of pushes a from-scratch container observes before it turns its
1763 /// gathered stats into a safe codec.
1764 ///
1765 /// A safe codec has at most `256 - SAFE_TAG_BASE` (= 134) dictionary slots per
1766 /// column, so we only need to identify ~134 genuinely-popular values. The
1767 /// `MisraGries` summary retains up to `2 * k` (= 1024) distinct candidates
1768 /// between tidies and reduces to `k` (= 512), comfortably more than 134, so the
1769 /// threshold just needs to be large enough that heavy hitters accumulate counts
1770 /// well above 1 before we freeze the codec. 64Ki pushes gives that headroom while
1771 /// keeping the pre-codec (uncompressed) window short.
1772 const STATS_THRESHOLD: usize = 64 * 1024;
1773
1774 impl PushInto<DatumSeq<'_>> for DatumContainer {
1775 #[inline]
1776 fn push_into(&mut self, item: DatumSeq<'_>) {
1777 // Fast path: container and item are both unencoded.
1778 // This is the hot path when dictionary compression is disabled.
1779 if self.codec.is_none() && self.stats.is_none() && item.iter.index.is_none() {
1780 self.inner.push_ref(item.iter.data);
1781 return;
1782 }
1783
1784 // Check if we've gathered enough stats to install a safe codec.
1785 if self.codec.is_none() && self.stats.is_some() && self.inner.len() >= STATS_THRESHOLD {
1786 let stats = self.stats.take().unwrap();
1787 self.codec = Some(stats.new_safe());
1788 }
1789
1790 if let Some(codec) = &mut self.codec {
1791 // Encode using the installed codec.
1792 codec.encode(item.bytes_iter(), &mut self.staging);
1793 } else if let Some(stats) = &mut self.stats {
1794 // Stats-gathering phase: feed the statistics but store raw bytes.
1795 // `observe` updates the heavy-hitter/tag summaries without encoding, so
1796 // we copy each row exactly once (below) instead of also encoding it into
1797 // a buffer we would immediately discard.
1798 stats.observe(item.bytes_iter());
1799 for slice in item.bytes_iter() {
1800 self.staging.extend_from_slice(slice);
1801 }
1802 } else {
1803 // No codec, no stats: raw copy.
1804 for slice in item.bytes_iter() {
1805 self.staging.extend_from_slice(slice);
1806 }
1807 }
1808 self.inner.push_ref(&self.staging[..]);
1809 self.staging.clear();
1810 }
1811 }
1812
1813 use mz_repr::{Datum, read_datum};
1814
1815 /// A reference that can be resolved to a sequence of `Datum`s.
1816 ///
1817 /// This type must "compare" as if decoded to a `Row`, which means it needs to track
1818 /// various nuances of `Row::cmp`, which at the moment is first by length, and then by
1819 /// the raw binary slice backing the row. Neither of those are explicit in this struct.
1820 /// We will need to produce them in order to perform comparisons.
1821 #[derive(Debug)]
1822 pub struct DatumSeq<'a> {
1823 pub iter: ColumnsIter<'a>,
1824 }
1825
1826 impl<'a> DatumSeq<'a> {
1827 #[inline(always)]
1828 fn borrow_as(other: &'a RowRef) -> Self {
1829 Self {
1830 iter: ColumnsCodec::borrow_row(other),
1831 }
1832 }
1833
1834 /// Borrow a `Row` as a `DatumSeq` so that it can be used to seek into a
1835 /// trace whose key/value container is a [`DatumContainer`].
1836 #[inline]
1837 pub fn from_row(row: &'a Row) -> Self {
1838 Self::borrow_as(row)
1839 }
1840
1841 #[inline]
1842 pub fn to_row(&self) -> Row {
1843 // Fast path: unencoded data is already row-formatted bytes.
1844 if self.iter.index.is_none() {
1845 return unsafe { Row::from_bytes_unchecked(self.iter.data) };
1846 }
1847 Row::pack(*self)
1848 }
1849 }
1850
1851 impl<'a> Copy for DatumSeq<'a> {}
1852 impl<'a> Clone for DatumSeq<'a> {
1853 #[inline(always)]
1854 fn clone(&self) -> Self {
1855 *self
1856 }
1857 }
1858
1859 use std::cmp::Ordering;
1860 impl<'a, 'b> PartialEq<DatumSeq<'a>> for DatumSeq<'b> {
1861 #[inline(always)]
1862 fn eq(&self, other: &DatumSeq<'a>) -> bool {
1863 // Fast path: both sides are unencoded raw row bytes.
1864 if self.iter.index.is_none() && other.iter.index.is_none() {
1865 return self.iter.data == other.iter.data;
1866 }
1867 Iterator::eq(self.iter, other.iter)
1868 }
1869 }
1870 impl<'a> Eq for DatumSeq<'a> {}
1871 impl<'a, 'b> PartialOrd<DatumSeq<'a>> for DatumSeq<'b> {
1872 #[inline(always)]
1873 fn partial_cmp(&self, other: &DatumSeq<'a>) -> Option<Ordering> {
1874 // Fast path: both sides are unencoded raw row bytes.
1875 if self.iter.index.is_none() && other.iter.index.is_none() {
1876 let left = self.iter.data;
1877 let right = other.iter.data;
1878 return Some(match left.len().cmp(&right.len()) {
1879 Ordering::Equal => left.cmp(right),
1880 other => other,
1881 });
1882 }
1883 // Slow path: at least one side is dictionary-encoded.
1884 // Fused length + lexicographic comparison in a single pass per side.
1885 // Row ordering is: shorter < longer; equal lengths compared lexicographically.
1886 //
1887 // We compare byte-by-byte (via `flatten`) rather than slice-by-slice on
1888 // purpose: a dictionary tag expands to a multi-byte value on one side while
1889 // the other side may store those same bytes raw, so the per-column slice
1890 // boundaries do not line up between the two iterators. Decoding to a flat
1891 // byte stream is the only representation in which both sides are directly
1892 // comparable. This path is cold — it only runs when at least one operand is
1893 // dictionary-encoded; the common unencoded case is handled by the fast path
1894 // above with a single slice comparison.
1895 let mut left = self.iter.flatten();
1896 let mut right = other.iter.flatten();
1897 let mut first_diff = Ordering::Equal;
1898 loop {
1899 match (left.next(), right.next()) {
1900 (Some(l), Some(r)) => {
1901 if first_diff == Ordering::Equal {
1902 first_diff = l.cmp(r);
1903 }
1904 }
1905 // Left exhausted first: left is shorter, so Less.
1906 (None, Some(_)) => return Some(Ordering::Less),
1907 // Right exhausted first: right is shorter, so Greater.
1908 (Some(_), None) => return Some(Ordering::Greater),
1909 // Same length: use first lexicographic difference.
1910 (None, None) => return Some(first_diff),
1911 }
1912 }
1913 }
1914 }
1915 impl<'a> Ord for DatumSeq<'a> {
1916 #[inline(always)]
1917 fn cmp(&self, other: &Self) -> Ordering {
1918 self.partial_cmp(other).unwrap()
1919 }
1920 }
1921
1922 impl<'a> PartialEq<&'a Row> for DatumSeq<'a> {
1923 #[inline(always)]
1924 fn eq(&self, other: &&'a Row) -> bool {
1925 self.eq(&Self::borrow_as(*other))
1926 }
1927 }
1928
1929 // Lifetimes decoupled (`'b` independent of `'a`): the arrange machinery
1930 // requires `for<'b> DatumSeq<'a>: PartialEq<&'b RowRef>`, i.e. a fixed
1931 // `DatumSeq` must compare against a `&RowRef` of any lifetime.
1932 impl<'a, 'b> PartialEq<&'b RowRef> for DatumSeq<'a> {
1933 #[inline(always)]
1934 fn eq(&self, other: &&'b RowRef) -> bool {
1935 self.eq(&DatumSeq::borrow_as(*other))
1936 }
1937 }
1938
1939 impl<'a> DatumSeq<'a> {
1940 #[inline(always)]
1941 pub fn bytes_iter(self) -> ColumnsIter<'a> {
1942 self.iter
1943 }
1944 }
1945
1946 impl<'a> Iterator for DatumSeq<'a> {
1947 type Item = Datum<'a>;
1948 #[inline(always)]
1949 fn next(&mut self) -> Option<Self::Item> {
1950 // Delegate to `ColumnsIter`, which handles both the codec and no-codec
1951 // cases. The no-codec scan hot path is served directly by `extend_datums`
1952 // (which decodes without going through this iterator), so the only callers
1953 // left here are the codec-encoded `extend_datums`/`to_row` paths and tests;
1954 // none warrant a dedicated no-codec fast path.
1955 self.iter
1956 .next()
1957 .map(|mut bytes| unsafe { read_datum(&mut bytes) })
1958 }
1959 }
1960
1961 use mz_repr::RowArena;
1962 use mz_repr::fixed_length::ExtendDatums;
1963 impl<'long> ExtendDatums for DatumSeq<'long> {
1964 #[inline]
1965 fn extend_datums<'a>(
1966 &'a self,
1967 _arena: &'a RowArena,
1968 target: &mut Vec<Datum<'a>>,
1969 max: Option<usize>,
1970 ) {
1971 // Branch on codec presence ONCE per row rather than once per datum.
1972 // With no codec (the common, feature-off case) push raw datums in a
1973 // tight loop, matching the pre-dictionary path; with a codec, fall
1974 // back to the per-column iterator. This keeps the codec check out of
1975 // the per-datum loop — the source of the feature-off scan overhead.
1976 if self.iter.index.is_none() {
1977 let mut data = self.iter.data;
1978 match max {
1979 Some(max) => {
1980 let mut n = 0;
1981 while n < max && !data.is_empty() {
1982 target.push(unsafe { read_datum(&mut data) });
1983 n += 1;
1984 }
1985 }
1986 None => {
1987 while !data.is_empty() {
1988 target.push(unsafe { read_datum(&mut data) });
1989 }
1990 }
1991 }
1992 } else {
1993 match max {
1994 Some(max) => target.extend((*self).take(max)),
1995 None => target.extend(*self),
1996 }
1997 }
1998 }
1999 }
2000}
2001
2002/// Traits abstracting the processes of encoding and decoding row-encoded byte sequences.
2003///
2004/// It is unsafe to use these types to encode byte sequences that are not row-encoded,
2005/// as they are parsed out of contiguous `[u8]` slices using `mz_repr::read_datum`.
2006mod row_codec {
2007
2008 pub use self::misra_gries::MisraGries;
2009 pub use columns::{ColumnsCodec, ColumnsIter};
2010 pub use dictionary::DictionaryCodec;
2011 #[cfg(test)]
2012 pub use dictionary::SAFE_TAG_BASE;
2013
2014 // Deterministic hasher state for the codecs' hash maps: a fixed-seed
2015 // `ahash::RandomState` shared with `mz_timely_util`'s consolidation hasher, so
2016 // the heavy-hitter summaries — and therefore which values each codec compresses
2017 // — are identical across runs and replicas, as the old `BTreeMap` backing was.
2018 use mz_timely_util::hash::fixed_state;
2019
2020 // The codecs encode and decode `[u8]` data specific to the `[Row]` encoding. They
2021 // soundly decode data they themselves encoded from valid `[Row]` data, but may be
2022 // unsound if asked to decode data that was not row-encoded, or was encoded with a
2023 // different codec. `ColumnsCodec` (a per-column wrapper around `DictionaryCodec`) is
2024 // the only codec the spine instantiates; the methods are inherent rather than behind
2025 // a `Codec` trait because nothing ever dispatches over codecs generically.
2026
2027 mod columns {
2028
2029 use mz_repr::{RowRef, read_datum};
2030
2031 use super::DictionaryCodec;
2032
2033 /// Independently encodes each column.
2034 #[derive(Default, Debug)]
2035 pub struct ColumnsCodec {
2036 columns: Vec<DictionaryCodec>,
2037 }
2038
2039 impl ColumnsCodec {
2040 /// Decode a row-encoded byte slice into per-column byte slices.
2041 pub(crate) fn decode<'a>(&'a self, bytes: &'a [u8]) -> ColumnsIter<'a> {
2042 ColumnsIter {
2043 index: Some(self),
2044 column: 0,
2045 data: bytes,
2046 }
2047 }
2048 /// Encode a sequence of column byte slices, updating per-column statistics.
2049 pub(crate) fn encode<'a, I>(&mut self, iter: I, output: &mut Vec<u8>)
2050 where
2051 I: IntoIterator<Item = &'a [u8]>,
2052 {
2053 for (index, bytes) in iter.into_iter().enumerate() {
2054 if self.columns.len() <= index {
2055 self.columns.push(Default::default());
2056 }
2057 self.columns[index].encode(std::iter::once(bytes), output);
2058 }
2059 }
2060
2061 /// Construct a codec valid for the union of the supplied codecs' data.
2062 pub(crate) fn new_from<'a>(stats: impl IntoIterator<Item = &'a Self>) -> Self {
2063 // An empty `stats` iterator yields a zero-column codec, which encodes and
2064 // decodes nothing; callers merging no inputs get an inert (but sound) codec.
2065 let stats = stats.into_iter().collect::<Vec<_>>();
2066 let cols = stats.iter().map(|s| s.columns.len()).max().unwrap_or(0);
2067 let mut columns = Vec::with_capacity(cols);
2068 let default: DictionaryCodec = Default::default();
2069 for index in 0..cols {
2070 columns.push(DictionaryCodec::new_from(
2071 stats
2072 .iter()
2073 .map(|s| s.columns.get(index).unwrap_or(&default)),
2074 ));
2075 }
2076 Self { columns }
2077 }
2078
2079 /// Reveal a row's bytes for fast-path comparison, with no codec to consult.
2080 #[inline(always)]
2081 pub(crate) fn borrow_row(row: &RowRef) -> ColumnsIter<'_> {
2082 ColumnsIter {
2083 index: None,
2084 column: 0,
2085 data: row.data(),
2086 }
2087 }
2088 }
2089
2090 impl ColumnsCodec {
2091 /// Visit contained allocations to determine their size and capacity.
2092 pub(crate) fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2093 let elem = std::mem::size_of::<DictionaryCodec>();
2094 callback(self.columns.len() * elem, self.columns.capacity() * elem);
2095 for column in &self.columns {
2096 column.heap_size(callback);
2097 }
2098 }
2099 }
2100
2101 impl ColumnsCodec {
2102 /// Record a row's column values in the statistics without encoding.
2103 ///
2104 /// Used during the stats-gathering phase, where we want the heavy-hitter
2105 /// and tag-usage information but store the row raw, so encoding into a
2106 /// throwaway buffer would be pure waste.
2107 #[inline]
2108 pub(crate) fn observe<'a, I>(&mut self, iter: I)
2109 where
2110 I: IntoIterator<Item = &'a [u8]>,
2111 {
2112 for (index, bytes) in iter.into_iter().enumerate() {
2113 if self.columns.len() <= index {
2114 self.columns.push(Default::default());
2115 }
2116 self.columns[index].observe(bytes);
2117 }
2118 }
2119 }
2120
2121 impl ColumnsCodec {
2122 /// Construct a codec using only structurally safe tags.
2123 ///
2124 /// Consumes `self`: this is only ever called on stats that have just been
2125 /// `take`n out of a container and are about to be discarded, so we move the
2126 /// per-column `MisraGries` summaries through rather than cloning them.
2127 pub(crate) fn new_safe(self) -> Self {
2128 let columns = self
2129 .columns
2130 .into_iter()
2131 .map(DictionaryCodec::new_safe)
2132 .collect();
2133 Self { columns }
2134 }
2135 }
2136
2137 #[derive(Debug, Copy, Clone)]
2138 pub struct ColumnsIter<'a> {
2139 // `None` when iterating an owned row directly, with no codec to consult.
2140 pub index: Option<&'a ColumnsCodec>,
2141 pub column: usize,
2142 pub data: &'a [u8],
2143 }
2144
2145 impl<'a> Iterator for ColumnsIter<'a> {
2146 type Item = &'a [u8];
2147 #[inline(always)]
2148 fn next(&mut self) -> Option<Self::Item> {
2149 if self.data.is_empty() {
2150 None
2151 } else if let Some(bytes) = self
2152 .index
2153 .as_ref()
2154 .and_then(|i| i.columns.get(self.column))
2155 .and_then(|i| i.decode.get(self.data[0].into()))
2156 {
2157 self.data = &self.data[1..];
2158 self.column += 1;
2159 Some(bytes)
2160 } else {
2161 let mut data = self.data;
2162 let data_len = data.len();
2163 unsafe {
2164 read_datum(&mut data);
2165 }
2166 let (prev, next) = self.data.split_at(data_len - data.len());
2167 self.data = next;
2168 self.column += 1;
2169 Some(prev)
2170 }
2171 }
2172 }
2173
2174 impl<'a> ColumnsIter<'a> {
2175 /// Create a column iterator without a codec.
2176 ///
2177 /// This requires the data to be row-formatted, and it will be erroneous otherwise.
2178 #[inline(always)]
2179 pub unsafe fn without_codec(data: &'a [u8]) -> Self {
2180 Self {
2181 index: None,
2182 column: 0,
2183 data,
2184 }
2185 }
2186 }
2187 }
2188
2189 /// A dictionary encoding codec for `[Row]` data.
2190 ///
2191 /// The dictionary harvests unused tags within each column and uses them to
2192 /// represent popular values within that column. There are two mechanisms it
2193 /// uses to accomplish this:
2194 ///
2195 /// 1. Statically free tags: `SAFE_TAG_BASE` is taken as an exclusive upper bound
2196 /// on the tags that will be used by `[Row]`, and tags greater or equal to this
2197 /// value are always safe to use.
2198 /// 2. Dynamically free tags: having seen an entire collection, we can use any
2199 /// tag not otherwise used by the collection, as it would not be ambiguous.
2200 ///
2201 /// It goes without saying that if either of these approaches are incorrect,
2202 /// there are calamitous unsoundness implications.
2203 mod dictionary {
2204 // The `encode` map is a pure value->tag lookup table (never iterated for logic),
2205 // so `mz_ore::collections::HashMap`'s order-hiding would suffice — but it offers
2206 // no fixed-seed constructor, and we want the same deterministic hasher as the
2207 // summary above. `heap_size`'s `keys()` walk is an order-insensitive sum.
2208 #![allow(clippy::disallowed_types)]
2209
2210 use std::collections::HashMap;
2211
2212 use super::fixed_state;
2213 pub use super::{BytesMap, MisraGries};
2214
2215 /// First byte value that is structurally unused by the datum encoding.
2216 /// All byte values >= this are safe to use as dictionary tags without
2217 /// observing the data, since no datum's first byte can have this value.
2218 ///
2219 /// `mz_repr`'s `Row` `Tag` enum currently has 84 variants (discriminants
2220 /// 0..=83), so the truly tight bound is 84. We deliberately pick a larger,
2221 /// round-ish constant to leave headroom for new tags without having to also
2222 /// bump the safe set, and the `test_safe_tag_base` test pins the real
2223 /// invariant: every datum the row format produces must encode with a first
2224 /// byte strictly less than this value. If a future tag crosses the boundary
2225 /// that test fails loudly rather than silently corrupting decoding.
2226 pub const SAFE_TAG_BASE: u8 = 122;
2227
2228 /// Per-column dictionary codec. Encodes column byte slices, replacing popular
2229 /// values with spare tags; decoding is performed by `ColumnsIter` reading the
2230 /// `decode` map directly.
2231 #[derive(Default, Debug)]
2232 pub struct DictionaryCodec {
2233 // Looked up once per value on the encode path; mostly misses (only popular
2234 // values compress), so a hash map beats a `BTreeMap`'s byte-slice walk. The
2235 // map is only ever read via `get` — never iterated — so its hasher seed has
2236 // no observable effect; the populated maps are built with `fixed_state` in
2237 // `new_from`/`new_safe` for consistency, while the derived-`Default` (stats
2238 // accumulator) variant stays empty and is never consulted.
2239 encode: HashMap<Vec<u8>, u8, ahash::RandomState>,
2240 pub decode: BytesMap,
2241 stats: (MisraGries<Vec<u8>>, [u64; 4]),
2242 }
2243
2244 impl DictionaryCodec {
2245 /// Encode a sequence of byte slices.
2246 ///
2247 /// Encoding also records statistics about the structure of the input.
2248 ///
2249 /// Decoding has no symmetric method here: a column's bytes are decoded by
2250 /// `ColumnsIter`, which consults the `decode` map directly.
2251 pub(super) fn encode<'a, I>(&mut self, iter: I, output: &mut Vec<u8>)
2252 where
2253 I: IntoIterator<Item = &'a [u8]>,
2254 {
2255 for bytes in iter.into_iter() {
2256 mz_ore::soft_assert_no_log!(
2257 !bytes.is_empty(),
2258 "row encoding never yields empty column slices",
2259 );
2260 // If we have an index referencing `bytes`, use the index key.
2261 if let Some(b) = self.encode.get(bytes) {
2262 output.push(*b);
2263 } else {
2264 // Raw fall-through. Soundness rests on `bytes[0]` never being a
2265 // tag we hand out as a dictionary key: `new_from`/`new_safe` only
2266 // assign dictionary tags from first-byte values that were never
2267 // observed (or are `>= SAFE_TAG_BASE`, which no datum first-byte
2268 // can equal). If a literal datum's first byte collided with a
2269 // dictionary tag, `decode` would resolve it to the dictionary
2270 // entry instead of reading the datum. This `debug_assert` makes
2271 // the load-bearing "no later first-byte outside the observed
2272 // union" invariant self-checking.
2273 mz_ore::soft_assert_no_log!(
2274 self.decode.get(bytes[0].into()).is_none(),
2275 "raw datum first-byte {} collides with a dictionary tag; \
2276 decode would be ambiguous",
2277 bytes[0],
2278 );
2279 output.extend(bytes);
2280 }
2281 self.observe(bytes);
2282 }
2283 }
2284
2285 /// Construct a new encoder from supplied statistics.
2286 pub(super) fn new_from<'a>(stats: impl IntoIterator<Item = &'a Self>) -> Self {
2287 // Collect most popular bytes from combined containers.
2288 let mut mg = MisraGries::default();
2289 let mut tags: [u64; 4] = [0; 4];
2290 for stat in stats.into_iter() {
2291 for (thing, count) in stat.stats.0.clone().done() {
2292 mg.update(thing, count);
2293 }
2294 tags[0] |= stat.stats.1[0];
2295 tags[1] |= stat.stats.1[1];
2296 tags[2] |= stat.stats.1[2];
2297 tags[3] |= stat.stats.1[3];
2298 }
2299 let mut mg = mg
2300 .done()
2301 .into_iter()
2302 .filter(|(next_bytes, count)| next_bytes.len() > 1 && count > &1);
2303 // Establish encoding and decoding rules.
2304 let mut encode = HashMap::with_hasher(fixed_state());
2305 let mut decode = BytesMap::default();
2306 for tag in 0..=255 {
2307 let tag_idx: usize = (tag % 4).into();
2308 let shift = tag >> 2;
2309 if (tags[tag_idx] >> shift) & 0x01 != 0 {
2310 // Tag is used by a literal datum first-byte; reserve the slot.
2311 decode.push(None);
2312 } else if let Some((next_bytes, _count)) = mg.next() {
2313 decode.push(Some(&next_bytes[..]));
2314 encode.insert(next_bytes, tag);
2315 } else {
2316 // Unused tag, but the heavy-hitter supply is exhausted. We must
2317 // still push a slot so that `decode`'s index stays aligned with
2318 // the tag value: every iteration pushes exactly once, keeping the
2319 // map length 256 and `decode.get(tag)` addressable by tag.
2320 decode.push(None);
2321 }
2322 }
2323
2324 Self {
2325 encode,
2326 decode,
2327 stats: (MisraGries::default(), [0u64; 4]),
2328 }
2329 }
2330 }
2331
2332 impl DictionaryCodec {
2333 /// Visit contained allocations to determine their size and capacity.
2334 ///
2335 /// The `encode` table is approximated as one logical entry's worth of bytes
2336 /// per element for size and its reserved `capacity()` for capacity; the
2337 /// dominant terms (the owned key bytes and the `decode` map's byte arena)
2338 /// are accounted exactly.
2339 pub fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2340 let entry = std::mem::size_of::<(Vec<u8>, u8)>();
2341 callback(self.encode.len() * entry, self.encode.capacity() * entry);
2342 for key in self.encode.keys() {
2343 callback(key.len(), key.capacity());
2344 }
2345 self.decode.heap_size(callback);
2346 self.stats.0.heap_size(callback);
2347 }
2348
2349 /// Record a single column value in this codec's statistics without
2350 /// producing any encoded output.
2351 ///
2352 /// Statistics come in two decoupled parts, with very different costs and
2353 /// purposes:
2354 ///
2355 /// 1. The tag bitmap (`stats.1`) records which first-byte values have been
2356 /// observed. It is cheap (four `u64` ORs) and *soundness critical*:
2357 /// `new_from`'s dynamic-tag path only hands out tags that this bitmap
2358 /// reports as unused, so it must stay accurate for the entire life of the
2359 /// codec, including on the hot encode path.
2360 /// 2. The MisraGries summary (`stats.0`) tracks heavy hitters and only
2361 /// affects *which* values a future codec compresses, never correctness.
2362 /// It is the expensive part (a `BTreeMap` insert per column per row). We
2363 /// keep feeding it after install, on the hot encode path, on purpose: a
2364 /// later merge rebuilds the merged codec from these summaries via
2365 /// `new_from`. If we froze the summary at install time, then as the
2366 /// collection evolves — records cancel under consolidation, the popular
2367 /// set drifts — the codec could never reclaim slots for newly-popular
2368 /// values and would eventually be left compressing values that no longer
2369 /// occur, ceasing to compress the ones that do.
2370 #[inline]
2371 pub fn observe(&mut self, bytes: &[u8]) {
2372 mz_ore::soft_assert_no_log!(
2373 !bytes.is_empty(),
2374 "row encoding never yields empty column slices",
2375 );
2376 let tag = bytes[0];
2377 let tag_idx: usize = (tag % 4).into();
2378 self.stats.1[tag_idx] |= 1 << (tag >> 2);
2379 self.stats.0.insert_ref(bytes);
2380 }
2381
2382 /// Construct a codec using only structurally safe tags (>= SAFE_TAG_BASE).
2383 /// These tags never collide with datum first-bytes, so the codec can be
2384 /// installed without observing all data first.
2385 pub(super) fn new_safe(stats: Self) -> Self {
2386 // The container stores its pre-install rows raw, so the first-byte
2387 // bitmap (`stats.1`) gathered while observing them must carry over to
2388 // the installed codec. The bitmap is soundness-critical: a later
2389 // `new_from` merge consults it to decide which one-byte tags are free
2390 // to hand out as dictionary keys. If we dropped it here, the merge
2391 // could assign a dictionary tag equal to a pre-install datum's first
2392 // byte, after which `decode` would resolve that literal datum to the
2393 // dictionary entry. The MisraGries summary (`stats.0`), by contrast,
2394 // is consumed below to seed the dictionary and is reset, since the
2395 // installed codec re-accumulates it from rows it sees post-install.
2396 let (mg, observed_tags) = stats.stats;
2397 let mut mg = mg
2398 .done()
2399 .into_iter()
2400 .filter(|(next_bytes, count)| next_bytes.len() > 1 && count > &1);
2401 let mut encode = HashMap::with_hasher(fixed_state());
2402 let mut decode = BytesMap::default();
2403 // Fill slots 0..SAFE_TAG_BASE with None (reserved for datum tags).
2404 for _ in 0..SAFE_TAG_BASE {
2405 decode.push(None);
2406 }
2407 // Assign dictionary entries to safe tags.
2408 for tag in SAFE_TAG_BASE..=255 {
2409 if let Some((next_bytes, _count)) = mg.next() {
2410 decode.push(Some(&next_bytes[..]));
2411 encode.insert(next_bytes, tag);
2412 }
2413 }
2414 Self {
2415 encode,
2416 decode,
2417 stats: (MisraGries::default(), observed_tags),
2418 }
2419 }
2420 }
2421 }
2422
2423 /// A map from `0 .. something` to `Option<&[u8]>`.
2424 ///
2425 /// Non-empty slices are pushed in order, and can be retrieved by index.
2426 /// Pushing an empty slice is equivalent to pushing `None`.
2427 #[derive(Debug)]
2428 pub struct BytesMap {
2429 offsets: Vec<usize>,
2430 bytes: Vec<u8>,
2431 }
2432 impl Default for BytesMap {
2433 #[inline(always)]
2434 fn default() -> Self {
2435 Self {
2436 offsets: vec![0],
2437 bytes: Vec::new(),
2438 }
2439 }
2440 }
2441 impl BytesMap {
2442 #[inline]
2443 fn push(&mut self, input: Option<&[u8]>) {
2444 if let Some(bytes) = input {
2445 self.bytes.extend(bytes);
2446 }
2447 self.offsets.push(self.bytes.len());
2448 }
2449 /// Visit contained allocations to determine their size and capacity.
2450 fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2451 let off = std::mem::size_of::<usize>();
2452 callback(self.offsets.len() * off, self.offsets.capacity() * off);
2453 callback(self.bytes.len(), self.bytes.capacity());
2454 }
2455 #[inline]
2456 fn get(&self, index: usize) -> Option<&[u8]> {
2457 if index < self.offsets.len() - 1 {
2458 let lower = self.offsets[index];
2459 let upper = self.offsets[index + 1];
2460 if lower < upper {
2461 Some(&self.bytes[lower..upper])
2462 } else {
2463 None
2464 }
2465 } else {
2466 None
2467 }
2468 }
2469 }
2470
2471 mod misra_gries {
2472 // The summary must iterate its entries (to extract heavy hitters in `done`, to
2473 // `tidy`, and to size itself), which `mz_ore::collections::HashMap` deliberately
2474 // forbids. We instead get determinism from the fixed-seed hasher (`fixed_state`)
2475 // plus the total-order sort in `done`; `tidy`/`heap_size` are order-insensitive.
2476 #![allow(clippy::disallowed_types)]
2477
2478 use std::collections::HashMap;
2479 use std::hash::Hash;
2480
2481 use super::fixed_state;
2482
2483 /// Maintains a summary of "heavy hitters" in a presented collection of items.
2484 ///
2485 /// Uses a hash map internally so that repeated observations of the same
2486 /// element only allocate once (on first sighting), and so the per-element
2487 /// `insert_ref` is an O(1) hash rather than an O(log n) walk of byte-slice
2488 /// comparisons. This is the hot path: one lookup per column per row, fed both
2489 /// while gathering stats and on the steady-state encode path. The hasher is
2490 /// fixed-seed (see [`fixed_state`]) so the summary — and thus which values a
2491 /// codec compresses — stays deterministic across runs and replicas.
2492 ///
2493 /// Tidy is performed when the number of *distinct* elements exceeds `2 * k`,
2494 /// reducing to at most `k` entries.
2495 #[derive(Clone, Debug)]
2496 pub struct MisraGries<T: Ord + Hash> {
2497 inner: HashMap<T, usize, ahash::RandomState>,
2498 k: usize,
2499 }
2500
2501 impl<T: Ord + Hash> Default for MisraGries<T> {
2502 #[inline(always)]
2503 fn default() -> Self {
2504 Self {
2505 inner: HashMap::with_hasher(fixed_state()),
2506 k: 512,
2507 }
2508 }
2509 }
2510
2511 impl<T: Ord + Hash> MisraGries<T> {
2512 /// Inserts an additional element to the summary.
2513 #[inline(always)]
2514 pub fn insert(&mut self, element: T) {
2515 self.update(element, 1);
2516 }
2517 /// Inserts multiple copies of an element to the summary.
2518 #[inline]
2519 pub fn update(&mut self, element: T, count: usize) {
2520 *self.inner.entry(element).or_insert(0) += count;
2521 if self.inner.len() > 2 * self.k {
2522 self.tidy();
2523 }
2524 }
2525
2526 /// Completes the summary, and extracts the items and their counts.
2527 pub fn done(self) -> Vec<(T, usize)> {
2528 let mut result: Vec<_> = self.inner.into_iter().collect();
2529 // Descending count, ties broken by key, so the values a codec selects
2530 // are deterministic regardless of hash-map iteration order.
2531 result.sort_by(|x, y| y.1.cmp(&x.1).then_with(|| x.0.cmp(&y.0)));
2532 result
2533 }
2534
2535 /// Reduces the summary down to at most `k` distinct items by
2536 /// subtracting the (k+1)-th largest count from all entries and
2537 /// discarding those that drop to zero or below.
2538 fn tidy(&mut self) {
2539 let mut counts: Vec<usize> = self.inner.values().copied().collect();
2540 counts.sort_unstable_by(|a, b| b.cmp(a));
2541 // The (k+1)-th largest count, or 0 if fewer than k+1 entries.
2542 let sub_weight = counts.get(self.k).copied().unwrap_or(0);
2543 if sub_weight > 0 {
2544 self.inner.retain(|_, count| {
2545 *count = count.saturating_sub(sub_weight);
2546 *count > 0
2547 });
2548 }
2549 }
2550 }
2551
2552 impl MisraGries<Vec<u8>> {
2553 /// Visit contained allocations to determine their size and capacity.
2554 ///
2555 /// The hash table is approximated as one logical entry per element for
2556 /// size and its reserved `capacity()` for capacity; the owned key bytes
2557 /// are accounted exactly.
2558 pub fn heap_size(&self, callback: &mut impl FnMut(usize, usize)) {
2559 let entry = std::mem::size_of::<(Vec<u8>, usize)>();
2560 callback(self.inner.len() * entry, self.inner.capacity() * entry);
2561 for key in self.inner.keys() {
2562 callback(key.len(), key.capacity());
2563 }
2564 }
2565
2566 /// Insert a borrowed byte slice, only allocating if the key is new.
2567 #[inline]
2568 pub fn insert_ref(&mut self, element: &[u8]) {
2569 if let Some(count) = self.inner.get_mut(element) {
2570 *count += 1;
2571 } else {
2572 self.insert(element.to_owned());
2573 }
2574 }
2575 }
2576
2577 impl<T: Ord + Hash> std::ops::AddAssign for MisraGries<T> {
2578 fn add_assign(&mut self, rhs: Self) {
2579 for (element, count) in rhs.done() {
2580 self.update(element, count);
2581 }
2582 }
2583 }
2584 }
2585}