1use 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
45pub type ColumnarCollection<'scope, T, D, R> = Collection<'scope, T, Column<(D, T, R)>>;
51
52pub type ColCollection<'scope, T> = ColumnarCollection<'scope, T, Row, Diff>;
55
56pub 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
66pub 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 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 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
132fn negate_column<T>(column: Column<(Row, T, Diff)>) -> Column<(Row, T, Diff)>
148where
149 T: RenderTimestamp,
150{
151 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
181pub type RecTimestamp = Product<mz_repr::Timestamp, PointStamp<u64>>;
183
184fn truncate_times(
190 column: Column<(Row, RecTimestamp, Diff)>,
191 level: usize,
192) -> Column<(Row, RecTimestamp, Diff)> {
193 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
228pub 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 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
277pub 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
303pub 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 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 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
359pub 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
388pub 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 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 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 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 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 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 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 #[mz_ore::test]
617 #[cfg_attr(miri, ignore)] 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 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 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}