Skip to main content

mz_compute/render/
columnar.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//! Columnar dataflow edge support.
11//!
12//! Defines [`ColCollection`], the columnar collection that dataflow edges between Plan
13//! nodes carry. Every producer emits this representation.
14//!
15//! Within a Plan node, operators may freely materialize `Vec` collections. Only
16//! the collection edge format is constrained. A node that produces a row-based
17//! collection re-encodes it to the columnar edge at its output leaf via
18//! [`vec_to_columnar`]. A node that must consume rows decodes at its input leaf
19//! via [`columnar_to_vec`]. Both are named operators (`VecToColumnar`,
20//! `ColumnarToVec`), so those leaf seams stay visible in dataflow
21//! introspection.
22
23use columnar::{Borrow, Columnar, Container, Index, Len, Push};
24use differential_dataflow::dynamic::pointstamp::{PointStamp, PointStampSummary};
25use differential_dataflow::{AsCollection, Collection, VecCollection};
26use mz_repr::{DatumVec, DatumVecBorrow, Diff, Row};
27use mz_timely_util::columnar::Column;
28use mz_timely_util::columnar::builder::ColumnBuilder;
29use mz_timely_util::columnar::chunk::{AccountedChunkBatcher, ChunkChunker};
30use mz_timely_util::columnar::columnar_consolidate_exchange;
31use mz_timely_util::operator::consolidate_pact;
32use timely::ContainerBuilder;
33use timely::container::{CapacityContainerBuilder, NoopBuilder};
34use timely::dataflow::channels::pact::{ExchangeCore, Pipeline};
35use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
36use timely::dataflow::operators::generic::{Operator, OutputBuilder};
37use timely::dataflow::{Scope, Stream, StreamVec};
38use timely::order::Product;
39use timely::progress::Antichain;
40
41use crate::render::RenderTimestamp;
42use crate::render::context::{ECB, Session};
43use crate::render::errors::DataflowErrorSer;
44
45/// A columnar collection of `(D, T, R)` updates traveling on a compute
46/// dataflow edge.
47///
48/// Mirrors differential's [`VecCollection<'scope, T, D, R>`]; the underlying
49/// container is [`Column<(D, T, R)>`] instead of `Vec<(D, T, R)>`.
50pub type ColumnarCollection<'scope, T, D, R> = Collection<'scope, T, Column<(D, T, R)>>;
51
52/// A columnar collection of `(Row, Diff)` updates, which is what a dataflow edge between
53/// Plan nodes carries.
54pub type ColCollection<'scope, T> = ColumnarCollection<'scope, T, Row, Diff>;
55
56/// Concatenates a collection of columnar edges.
57pub fn concat_many<'scope, T, I>(scope: Scope<'scope, T>, edges: I) -> ColCollection<'scope, T>
58where
59    T: RenderTimestamp,
60    I: IntoIterator<Item = ColCollection<'scope, T>>,
61{
62    let cols: Vec<_> = edges.into_iter().collect();
63    differential_dataflow::collection::concatenate(scope, cols)
64}
65
66/// Applies `logic` to each record in `edge`, exposing the record as a borrowed
67/// [`DatumVecBorrow`] and giving it ok and err output sessions.
68///
69/// `name` is the rendered operator's name. `max_demand` bounds the number of columns decoded
70/// per row. Pass `usize::MAX` to decode all columns.
71///
72/// This is the canonical entry point for "decoding consumers" (operators that
73/// read [`mz_repr::Datum`]s from each row anyway). It iterates the columnar
74/// batch directly without going through an owned [`Row`].
75pub fn flat_map_datums<'scope, T, DCB, L>(
76    edge: ColCollection<'scope, T>,
77    name: &str,
78    max_demand: usize,
79    mut logic: L,
80) -> (
81    Stream<'scope, T, DCB::Container>,
82    StreamVec<'scope, T, (DataflowErrorSer, T, Diff)>,
83)
84where
85    T: RenderTimestamp,
86    DCB: ContainerBuilder,
87    L: for<'a> FnMut(
88            &'a mut DatumVecBorrow<'_>,
89            T,
90            Diff,
91            &mut Session<T, DCB>,
92            &mut Session<T, ECB<T>>,
93        ) -> usize
94        + 'static,
95{
96    let scope = edge.inner.scope();
97    let mut builder = OperatorBuilder::new(name.to_string(), scope);
98    let (ok_output, ok_stream) = builder.new_output();
99    let mut ok_output = OutputBuilder::<_, DCB>::from(ok_output);
100    let (err_output, err_stream) = builder.new_output();
101    let mut err_output = OutputBuilder::<_, ECB<T>>::from(err_output);
102    let mut input = builder.new_input(edge.inner, Pipeline);
103    builder.build(move |_capabilities| {
104        let mut datums = DatumVec::new();
105        move |_frontiers| {
106            let mut ok_output = ok_output.activate();
107            let mut err_output = err_output.activate();
108            input.for_each(|time, data| {
109                // Retain the input capability to derive a `Capability` for each
110                // output. The `Session` type alias is fixed to `Capability<T>`.
111                let ok_cap = time.retain(0);
112                let err_cap = time.retain(1);
113                let mut ok_session = ok_output.session_with_builder(&ok_cap);
114                let mut err_session = err_output.session_with_builder(&err_cap);
115                // Rows are read from the borrowed column, never materialized as
116                // owned `Row`s.
117                for (v, t, d) in data.borrow().into_index_iter() {
118                    logic(
119                        &mut datums.borrow_with_limit(v, max_demand),
120                        Columnar::into_owned(t),
121                        Columnar::into_owned(d),
122                        &mut ok_session,
123                        &mut err_session,
124                    );
125                }
126            });
127        }
128    });
129    (ok_stream, err_stream)
130}
131
132/// Negates the diff of every record in `column`, rebuilding only the diff column.
133///
134/// A `Typed` input hands its row and time columns over untouched, so only the
135/// diffs are rebuilt, which is 8 bytes per record. A serialized input keeps all
136/// three columns in one buffer, so its rows and times are copied in bulk.
137///
138/// Negation stays checked: `Neg for Overflowing` runs `overflowing_neg` and
139/// reports overflow, which `-Diff::MIN` triggers.
140///
141/// TODO: Negate a `Typed` input's diffs in place rather than rebuilding them.
142/// This cannot go through `columnar::IndexMut`: `Overflows` stores the raw
143/// integer and materializes `Overflowing` on read, so there is no wrapper in
144/// memory to borrow mutably, and handing out `&mut` to the raw integer would let
145/// writes bypass the checked arithmetic. It needs a checked bulk operation on the
146/// container instead, keeping the overflow check inside.
147fn negate_column<T>(column: Column<(Row, T, Diff)>) -> Column<(Row, T, Diff)>
148where
149    T: RenderTimestamp,
150{
151    /// Collects the negation of every diff in `diffs` into a fresh column.
152    fn negated_diffs<'a, D>(diffs: &'a D) -> <Diff as Columnar>::Container
153    where
154        D: Len + Index<Ref = Diff> + 'a,
155    {
156        let mut negated = <Diff as Columnar>::Container::default();
157        for index in 0..diffs.len() {
158            negated.push(-diffs.get(index));
159        }
160        negated
161    }
162
163    match column {
164        Column::Typed((rows, times, diffs)) => {
165            let negated = negated_diffs(&diffs.borrow());
166            Column::Typed((rows, times, negated))
167        }
168        column => {
169            let view = column.borrow();
170            let len = view.len();
171            let mut negated = <(Row, T, Diff) as Columnar>::Container::default();
172            let (rows, times, diffs) = &mut negated;
173            rows.extend_from_self(view.0, 0..len);
174            times.extend_from_self(view.1, 0..len);
175            *diffs = negated_diffs(&view.2);
176            Column::Typed(negated)
177        }
178    }
179}
180
181/// The timestamp of a scope that carries iteration coordinates.
182pub type RecTimestamp = Product<mz_repr::Timestamp, PointStamp<u64>>;
183
184/// Truncates every time in `column` to `level - 1` iteration coordinates.
185///
186/// A `Typed` input hands its row and diff columns over untouched, because truncation only
187/// rewrites times. A serialized input keeps all three columns in one buffer, so its rows
188/// and diffs are copied in bulk rather than per record.
189fn truncate_times(
190    column: Column<(Row, RecTimestamp, Diff)>,
191    level: usize,
192) -> Column<(Row, RecTimestamp, Diff)> {
193    /// Collects the truncation of every time in `times` into a fresh column.
194    fn truncated<'a, C>(times: C, level: usize) -> <RecTimestamp as Columnar>::Container
195    where
196        C: Len + Index<Ref = columnar::Ref<'a, RecTimestamp>> + 'a,
197    {
198        let mut truncated = <RecTimestamp as Columnar>::Container::default();
199        let mut time = RecTimestamp::default();
200        for reference in times.into_index_iter() {
201            time.copy_from(reference);
202            let mut coordinates = std::mem::take(&mut time.inner).into_inner();
203            coordinates.truncate(level - 1);
204            time.inner = PointStamp::new(coordinates);
205            truncated.push(&time);
206        }
207        truncated
208    }
209
210    match column {
211        Column::Typed((rows, times, diffs)) => {
212            let times = truncated(times.borrow(), level);
213            Column::Typed((rows, times, diffs))
214        }
215        column => {
216            let view = column.borrow();
217            let len = view.len();
218            let mut out = <(Row, RecTimestamp, Diff) as Columnar>::Container::default();
219            let (rows, times, diffs) = &mut out;
220            rows.extend_from_self(view.0, 0..len);
221            *times = truncated(view.1, level);
222            diffs.extend_from_self(view.2, 0..len);
223            Column::Typed(out)
224        }
225    }
226}
227
228/// Leaves a dynamically created scope that has `level` iteration coordinates.
229///
230/// The columnar counterpart of differential's `leave_dynamic`, which it offers only for
231/// `Vec` collections. Keeping the recursive binding's result on the edge is what lets the
232/// feedback loop run without a decode and a re-encode per iteration.
233pub fn columnar_leave_dynamic<'scope>(
234    collection: ColumnarCollection<'scope, RecTimestamp, Row, Diff>,
235    level: usize,
236) -> ColumnarCollection<'scope, RecTimestamp, Row, Diff> {
237    let scope = collection.inner.scope();
238    let mut builder = OperatorBuilder::new("ColumnarLeaveDynamic".to_string(), scope);
239    let (output, stream) = builder.new_output();
240    let mut output =
241        OutputBuilder::<_, NoopBuilder<Column<(Row, RecTimestamp, Diff)>>>::from(output);
242    // The connection summary tells the scope that this operator drops all but `level - 1`
243    // coordinates, so a downstream frontier is not held back by the iteration it leaves.
244    let summary = Product {
245        outer: Default::default(),
246        inner: PointStampSummary {
247            retain: Some(level - 1),
248            actions: Vec::new(),
249        },
250    };
251    let mut input = builder.new_input_connection(
252        collection.inner,
253        Pipeline,
254        [(0, Antichain::from_elem(summary))],
255    );
256
257    builder.build(move |_capability| {
258        move |_frontier| {
259            let mut output = output.activate();
260            input.for_each(|cap, data| {
261                let mut time = cap.time().clone();
262                let mut coordinates = std::mem::take(&mut time.inner).into_inner();
263                coordinates.truncate(level - 1);
264                time.inner = PointStamp::new(coordinates);
265                let cap = cap.delayed(&time, 0);
266                let mut truncated = truncate_times(std::mem::take(data), level);
267                output
268                    .session_with_builder(&cap)
269                    .give_container(&mut truncated);
270            });
271        }
272    });
273
274    stream.as_collection()
275}
276
277/// Negates the diff of every record in a [`ColumnarCollection`].
278pub fn columnar_negate<'scope, T>(
279    collection: ColumnarCollection<'scope, T, Row, Diff>,
280) -> ColumnarCollection<'scope, T, Row, Diff>
281where
282    T: RenderTimestamp,
283{
284    collection
285        .inner
286        .unary::<NoopBuilder<Column<(Row, T, Diff)>>, _, _, _>(
287            Pipeline,
288            "ColumnarNegate",
289            |_cap, _info| {
290                move |input, output| {
291                    input.for_each(|time, data| {
292                        let mut negated = negate_column(std::mem::take(data));
293                        output
294                            .session_with_builder(&time)
295                            .give_container(&mut negated);
296                    });
297                }
298            },
299        )
300        .as_collection()
301}
302
303/// Consolidates a [`ColumnarCollection`] natively, without a row round-trip.
304///
305/// A [`ChunkChunker`] sorts and consolidates the input columns and an
306/// [`AccountedChunkBatcher`] merges them, both holding their data in columnar form, so
307/// nothing outside the exchange pact visits a record or materializes an owned [`Row`].
308/// The batcher's chains are chunks, so the process buffer pool spills them while the
309/// chunk spill gate is set, bounding what a consolidation holds resident.
310///
311/// Uses [`consolidate_pact`] rather than `mz_arrange_core`: a consolidate emits a
312/// consolidated collection, so building and reading back a maintained trace would be
313/// wasted work.
314pub fn columnar_consolidate<'scope, T>(
315    collection: ColumnarCollection<'scope, T, Row, Diff>,
316    name: &str,
317) -> ColumnarCollection<'scope, T, Row, Diff>
318where
319    T: RenderTimestamp,
320{
321    // TODO: This pact re-serializes every record into a per-destination `ColumnBuilder`,
322    // the one remaining full re-encode on this path. Bulk routing needs contiguous ranges
323    // of records sharing a destination, which a per-record hash cannot identify.
324    let exchange = ExchangeCore::<ColumnBuilder<_>, _>::new_core(
325        columnar_consolidate_exchange::<Row, T, Diff>,
326    );
327    let consolidated = consolidate_pact::<
328        ChunkChunker<Row, T, Diff>,
329        AccountedChunkBatcher<Row, T, Diff>,
330        _,
331        _,
332    >(collection.inner, exchange, name);
333
334    // Flatten the sealed chain into one container per chunk, loading a spilled body
335    // and moving a resident one, visiting no record either way.
336    //
337    // TODO: This ships a whole sealed snapshot in one activation, an un-fueled burst
338    // hazard on large consolidations. `consolidate_named`'s unpack does the same, so a
339    // fuel fix has to cover both.
340    consolidated
341        .unary::<NoopBuilder<Column<(Row, T, Diff)>>, _, _, _>(
342            Pipeline,
343            &format!("Flatten {name}"),
344            |_cap, _info| {
345                move |input, output| {
346                    input.for_each(|time, data| {
347                        let mut session = output.session_with_builder(&time);
348                        for chunk in data.drain(..).flatten() {
349                            let mut column = Column::from(chunk.into_body());
350                            session.give_container(&mut column);
351                        }
352                    });
353                }
354            },
355        )
356        .as_collection()
357}
358
359/// Repacks a row-based collection into columnar batches.
360///
361/// The leaf encode described in the module docs, named `VecToColumnar` in a rendered
362/// dataflow. Repacking copies row bytes and allocates no per-record `Row`.
363pub fn vec_to_columnar<'scope, T>(
364    collection: VecCollection<'scope, T, Row, Diff>,
365) -> ColumnarCollection<'scope, T, Row, Diff>
366where
367    T: RenderTimestamp,
368{
369    collection
370        .inner
371        .unary::<ColumnBuilder<(Row, T, Diff)>, _, _, _>(
372            Pipeline,
373            "VecToColumnar",
374            |_cap, _info| {
375                move |input, output| {
376                    input.for_each(|time, data| {
377                        let mut session = output.session_with_builder(&time);
378                        for (v, t, d) in data.drain(..) {
379                            session.give((&v, &t, &d));
380                        }
381                    });
382                }
383            },
384        )
385        .as_collection()
386}
387
388/// Decodes columnar batches into a row-based collection.
389///
390/// The leaf decode described in the module docs, named `ColumnarToVec` in a rendered
391/// dataflow. It allocates an owned [`Row`] per record, which is why it stays at those
392/// boundaries.
393pub fn columnar_to_vec<'scope, T>(
394    collection: ColumnarCollection<'scope, T, Row, Diff>,
395) -> VecCollection<'scope, T, Row, Diff>
396where
397    T: RenderTimestamp,
398{
399    collection
400        .inner
401        .unary::<CapacityContainerBuilder<Vec<(Row, T, Diff)>>, _, _, _>(
402            Pipeline,
403            "ColumnarToVec",
404            |_cap, _info| {
405                move |input, output| {
406                    input.for_each(|time, data| {
407                        let mut session = output.session(&time);
408                        for (v, t, d) in data.borrow().into_index_iter() {
409                            session.give((
410                                Columnar::into_owned(v),
411                                Columnar::into_owned(t),
412                                Columnar::into_owned(d),
413                            ));
414                        }
415                    });
416                }
417            },
418        )
419        .as_collection()
420}
421
422#[cfg(test)]
423mod tests {
424    use differential_dataflow::input::Input;
425    use mz_ore::cast::CastFrom;
426    use mz_repr::{Datum, Timestamp};
427    use timely::dataflow::operators::Capture;
428    use timely::dataflow::operators::capture::{Event, Extract};
429
430    use super::*;
431
432    type RowBuilder = CapacityContainerBuilder<Vec<(Row, Timestamp, Diff)>>;
433    type CapturedRows = std::sync::mpsc::Receiver<Event<Timestamp, Vec<(Row, Timestamp, Diff)>>>;
434
435    fn extract_sorted(captured: CapturedRows) -> Vec<(Row, Timestamp, Diff)> {
436        let mut updates: Vec<_> = captured
437            .extract()
438            .into_iter()
439            .flat_map(|(_, data)| data)
440            .collect();
441        updates.sort();
442        updates
443    }
444
445    fn test_rows() -> Vec<Row> {
446        vec![
447            Row::pack_slice(&[Datum::Int32(42), Datum::String("hello")]),
448            Row::pack_slice(&[Datum::Int64(100), Datum::Null]),
449            Row::pack_slice(&[Datum::True, Datum::False, Datum::Null]),
450            Row::default(),
451        ]
452    }
453
454    #[mz_ore::test]
455    fn round_trip_through_columnar() {
456        let rows = test_rows();
457        let expected: Vec<_> = {
458            let mut updates: Vec<_> = rows
459                .iter()
460                .enumerate()
461                .map(|(i, r)| (r.clone(), Timestamp::from(u64::cast_from(i / 2)), Diff::ONE))
462                .collect();
463            updates.sort();
464            updates
465        };
466        let captured = timely::execute_directly(move |worker| {
467            worker.dataflow::<Timestamp, _, _>(|scope| {
468                let (mut input, collection) = scope.new_collection();
469                let captured = columnar_to_vec(vec_to_columnar(collection)).inner.capture();
470                for (i, row) in rows.into_iter().enumerate() {
471                    input.advance_to(Timestamp::from(u64::cast_from(i / 2)));
472                    input.update(row, Diff::ONE);
473                }
474                input.advance_to(Timestamp::from(2_u64));
475                input.flush();
476                captured
477            })
478        });
479        assert_eq!(extract_sorted(captured), expected);
480    }
481
482    #[mz_ore::test]
483    fn columnar_negate_flips_diffs() {
484        let rows = test_rows();
485        let expected: Vec<_> = {
486            let mut updates: Vec<_> = rows
487                .iter()
488                .map(|r| (r.clone(), Timestamp::from(0_u64), -Diff::ONE))
489                .collect();
490            updates.sort();
491            updates
492        };
493        let captured = timely::execute_directly(move |worker| {
494            worker.dataflow::<Timestamp, _, _>(|scope| {
495                let (mut input, collection) = scope.new_collection();
496                let edge = columnar_negate(vec_to_columnar(collection));
497                let captured = columnar_to_vec(edge).inner.capture();
498                for row in rows {
499                    input.update(row, Diff::ONE);
500                }
501                input.advance_to(Timestamp::from(1_u64));
502                input.flush();
503                captured
504            })
505        });
506        assert_eq!(extract_sorted(captured), expected);
507    }
508
509    #[mz_ore::test]
510    fn concat_many_concatenates_columnar() {
511        let rows = test_rows();
512        let expected: Vec<_> = {
513            let mut updates: Vec<_> = rows
514                .iter()
515                .map(|r| (r.clone(), Timestamp::from(0_u64), Diff::ONE))
516                .collect();
517            // The first row arrives on both inputs.
518            updates.push((rows[0].clone(), Timestamp::from(0_u64), Diff::ONE));
519            updates.sort();
520            updates
521        };
522        let captured = timely::execute_directly(move |worker| {
523            worker.dataflow::<Timestamp, _, _>(|scope| {
524                let (mut input1, collection1) = scope.new_collection();
525                let (mut input2, collection2) = scope.new_collection();
526                let edge = concat_many(
527                    scope,
528                    [vec_to_columnar(collection1), vec_to_columnar(collection2)],
529                );
530                let captured = columnar_to_vec(edge).inner.capture();
531                let (first, rest) = rows.split_first().unwrap();
532                input1.update(first.clone(), Diff::ONE);
533                input2.update(first.clone(), Diff::ONE);
534                for row in rest {
535                    input1.update(row.clone(), Diff::ONE);
536                }
537                for input in [&mut input1, &mut input2] {
538                    input.advance_to(Timestamp::from(1_u64));
539                    input.flush();
540                }
541                captured
542            })
543        });
544        assert_eq!(extract_sorted(captured), expected);
545    }
546
547    #[mz_ore::test]
548    fn flat_map_datums_arms_agree() {
549        // Project the first datum of each row, exercising `max_demand`.
550        let rows = test_rows();
551        let captured = timely::execute_directly(move |worker| {
552            worker.dataflow::<Timestamp, _, _>(|scope| {
553                let (mut input, collection) = scope.new_collection();
554                let (oks, _errs) = flat_map_datums::<_, RowBuilder, _>(
555                    vec_to_columnar(collection),
556                    "test",
557                    1,
558                    |datums, t, d, ok_session, _err_session| {
559                        ok_session.give((Row::pack(datums.iter()), t, d));
560                        1
561                    },
562                );
563                let captured = oks.capture();
564                for row in rows {
565                    input.update(row, Diff::ONE);
566                }
567                input.advance_to(Timestamp::from(1_u64));
568                input.flush();
569                captured
570            })
571        });
572        let updates = extract_sorted(captured);
573        assert!(!updates.is_empty());
574        // Each output row retains at most the first datum of its input.
575        assert!(updates.iter().all(|(r, _, _)| r.iter().count() <= 1));
576    }
577
578    #[mz_ore::test]
579    fn columnar_consolidate_accumulates_and_cancels() {
580        let row1 = Row::pack_slice(&[Datum::Int32(1)]);
581        let row2 = Row::pack_slice(&[Datum::Int32(2)]);
582        let row3 = Row::pack_slice(&[Datum::Int32(3)]);
583        // `row1` accumulates at t=0 and again at t=1, kept apart by time. `row2` cancels
584        // at t=0 and `row3` at t=1, so neither reaches the output.
585        let expected = vec![
586            (row1.clone(), Timestamp::from(0_u64), Diff::from(2)),
587            (row1.clone(), Timestamp::from(1_u64), Diff::ONE),
588        ];
589
590        let captured = timely::execute_directly(move |worker| {
591            worker.dataflow::<Timestamp, _, _>(|scope| {
592                let (mut input, collection) = scope.new_collection();
593                let edge = columnar_consolidate(vec_to_columnar(collection), "Test");
594                let captured = columnar_to_vec(edge).inner.capture();
595                // t=0: row1 accumulates (+1, +1), row2 cancels (+1, -1).
596                input.advance_to(Timestamp::from(0_u64));
597                input.update(row1.clone(), Diff::ONE);
598                input.update(row1.clone(), Diff::ONE);
599                input.update(row2.clone(), Diff::ONE);
600                input.update(row2, -Diff::ONE);
601                // t=1: row1 survives (+1), row3 cancels (+1, -1).
602                input.advance_to(Timestamp::from(1_u64));
603                input.update(row1, Diff::ONE);
604                input.update(row3.clone(), Diff::ONE);
605                input.update(row3, -Diff::ONE);
606                input.advance_to(Timestamp::from(2_u64));
607                input.flush();
608                captured
609            })
610        });
611        assert_eq!(extract_sorted(captured), expected);
612    }
613
614    /// Consolidation over a payload big enough to spill: the chunks the batcher
615    /// commits land in the pool, and flattening loads them back whole.
616    #[mz_ore::test]
617    #[cfg_attr(miri, ignore)] // the pool's mmap and madvise calls are unsupported under miri
618    fn columnar_consolidate_spills_and_round_trips() {
619        use mz_ore::pool::Pool;
620        use mz_timely_util::columnar::chunk::set_spill_override;
621
622        // Enough bytes for a committed chunk to clear the pool's 64 KiB spill
623        // floor. Every row is distinct, so consolidation cancels nothing and the
624        // whole payload has to survive the round trip through the pool.
625        let rows: Vec<Row> = (0..40_000i64)
626            .map(|i| Row::pack_slice(&[Datum::Int64(i), Datum::String("a repeated string value")]))
627            .collect();
628        let expected = rows.len();
629
630        let pool = Pool::new().expect("pool creation");
631        // The override is per-thread and `execute_directly` runs the worker on
632        // this one, so it covers the dataflow below and nothing else.
633        set_spill_override(Some(pool.clone()));
634        let captured = timely::execute_directly(move |worker| {
635            worker.dataflow::<Timestamp, _, _>(|scope| {
636                let (mut input, collection) = scope.new_collection();
637                let edge = columnar_consolidate(vec_to_columnar(collection), "Test");
638                let captured = columnar_to_vec(edge).inner.capture();
639                input.advance_to(Timestamp::from(0_u64));
640                for row in rows {
641                    input.update(row, Diff::ONE);
642                }
643                input.advance_to(Timestamp::from(1_u64));
644                input.flush();
645                captured
646            })
647        });
648        set_spill_override(None);
649
650        assert!(
651            pool.stats().inserts > 0,
652            "the payload should have reached the pool"
653        );
654        let updates = extract_sorted(captured);
655        assert_eq!(updates.len(), expected);
656        assert!(
657            updates
658                .iter()
659                .all(|(_, time, diff)| *time == Timestamp::from(0_u64) && *diff == Diff::ONE),
660            "every row survives at its own time and multiplicity"
661        );
662    }
663}