@@ -18,6 +18,7 @@ use mz_expr::{Eval, MfpPlan};
1818use mz_repr:: { DatumVec , RowArena , SharedRow } ;
1919use mz_repr:: { Diff , Row , RowRef , Timestamp } ;
2020use mz_timely_util:: columnar:: Column ;
21+ use mz_timely_util:: columnar:: consolidate:: ConsolidatingColumnBuilder ;
2122use mz_timely_util:: operator:: StreamExt ;
2223use timely:: Container ;
2324use timely:: container:: DrainContainer ;
@@ -76,15 +77,19 @@ impl<'scope, T: crate::render::RenderTimestamp> Context<'scope, T> {
7677 } ;
7778
7879 use differential_dataflow:: AsCollection ;
79- let ok_collection = oks. as_collection ( ) ;
80+ let ok_collection = CollectionEdge :: Columnar ( oks. as_collection ( ) ) ;
8081 let new_err_collection = errs. as_collection ( ) ;
8182 let err_collection = err_collection. concat ( new_err_collection) ;
82- CollectionBundle :: from_collections ( ok_collection, err_collection)
83+ CollectionBundle :: from_edge ( ok_collection, err_collection)
8384 }
8485}
8586
8687/// Output ok-session container builder for [`flat_map_stage`].
87- type FlatMapOk < T > = ConsolidatingContainerBuilder < Vec < ( Row , T , Diff ) > > ;
88+ ///
89+ /// Consolidating like the err builder, but emits `Column<(Row, T, Diff)>` so
90+ /// the FlatMap output travels as the columnar edge. Output rows are freshly
91+ /// built by the mfp, so the owned give into staging is a move, not a new alloc.
92+ type FlatMapOk < T > = ConsolidatingColumnBuilder < Row , T , Diff > ;
8893/// Output err-session container builder for [`flat_map_stage`].
8994type FlatMapErr < T > = ConsolidatingContainerBuilder < Vec < ( DataflowErrorSer , T , Diff ) > > ;
9095
@@ -137,7 +142,7 @@ fn flat_map_stage<'scope, T, C>(
137142 until : Antichain < Timestamp > ,
138143 budget : usize ,
139144) -> (
140- Stream < ' scope , T , Vec < ( Row , T , Diff ) > > ,
145+ Stream < ' scope , T , Column < ( Row , T , Diff ) > > ,
141146 Stream < ' scope , T , Vec < ( DataflowErrorSer , T , Diff ) > > ,
142147)
143148where
@@ -370,13 +375,14 @@ mod tests {
370375 let scope = stream. scope ( ) ;
371376 let ( oks, _errs) =
372377 flat_map_stage ( stream, scope, exprs, func, mfp, Antichain :: new ( ) , budget) ;
373- // Count through the container rather than per record. A
374- // per-record `inspect` needs `&Container: IntoIterator`, and
375- // resolving that on macOS recurses through `objc2`'s blanket
376- // impls until the trait solver overflows.
378+ // The columnar output ships whole `Column`s, so count records
379+ // through the container rather than per raw element. A
380+ // per-record `inspect` also needs `&Container: IntoIterator`,
381+ // and resolving that on macOS recurses through `objc2`'s
382+ // blanket impls until the trait solver overflows.
377383 oks. inspect_container ( move |event| {
378384 if let Ok ( ( _time, data) ) = event {
379- * sink. borrow_mut ( ) += data. len ( ) ;
385+ * sink. borrow_mut ( ) += data. borrow ( ) . into_index_iter ( ) . count ( ) ;
380386 }
381387 } ) ;
382388 input
@@ -487,21 +493,105 @@ mod tests {
487493 } )
488494 } ) ;
489495
490- let extract_sorted = |captured : std:: sync:: mpsc:: Receiver < _ > | {
491- let mut updates: Vec < ( Row , Timestamp , Diff ) > = captured
492- . extract ( )
493- . into_iter ( )
494- . flat_map ( |( _, data) | data)
495- . collect ( ) ;
496- updates. sort ( ) ;
497- updates
498- } ;
499- let vec_updates = extract_sorted ( vec_captured) ;
496+ let vec_updates = extract_sorted_columns ( vec_captured) ;
500497 assert ! ( !vec_updates. is_empty( ) ) ;
501498 assert ! (
502499 vec_updates. iter( ) . any( |( _, _, d) | * d < Diff :: ZERO ) ,
503500 "the retraction must survive as a negative diff"
504501 ) ;
505- assert_eq ! ( vec_updates, extract_sorted( col_captured) ) ;
502+ assert_eq ! ( vec_updates, extract_sorted_columns( col_captured) ) ;
503+ }
504+
505+ /// Decodes a capture of the columnar FlatMap output into sorted owned
506+ /// `(row, time, diff)` updates.
507+ fn extract_sorted_columns (
508+ captured : std:: sync:: mpsc:: Receiver <
509+ timely:: dataflow:: operators:: capture:: Event < Timestamp , Column < ( Row , Timestamp , Diff ) > > ,
510+ > ,
511+ ) -> Vec < ( Row , Timestamp , Diff ) > {
512+ let mut updates: Vec < ( Row , Timestamp , Diff ) > = captured
513+ . extract ( )
514+ . into_iter ( )
515+ . flat_map ( |( _, col) | {
516+ col. borrow ( )
517+ . into_index_iter ( )
518+ . map ( |( v, t, d) | {
519+ (
520+ Columnar :: into_owned ( v) ,
521+ Columnar :: into_owned ( t) ,
522+ Columnar :: into_owned ( d) ,
523+ )
524+ } )
525+ . collect :: < Vec < _ > > ( )
526+ } )
527+ . collect ( ) ;
528+ updates. sort ( ) ;
529+ updates
530+ }
531+
532+ #[ mz_ore:: test]
533+ fn flat_map_output_consolidates_within_batch ( ) {
534+ // Two distinct input rows whose table-function expansions overlap once
535+ // the mfp projects away the differing `stop` column. The overlapping
536+ // output rows land at the same time in one batch, so the consolidating
537+ // columnar output builder must fold them into summed diffs.
538+ let captured = timely:: execute_directly ( move |worker| {
539+ worker. dataflow :: < Timestamp , _ , _ > ( |scope| {
540+ let ( mut input, collection) = scope. new_collection ( ) ;
541+ let exprs = vec ! [
542+ LirScalarExpr :: column( 0 ) ,
543+ LirScalarExpr :: column( 1 ) ,
544+ LirScalarExpr :: literal_ok( Datum :: Int64 ( 1 ) , ReprScalarType :: Int64 ) ,
545+ ] ;
546+ let func = TableFunc :: GenerateSeriesInt64 ;
547+ // Project to only the generated value (column 2), collapsing the
548+ // two input rows' distinct (start, stop) prefixes.
549+ let mfp = MapFilterProject :: < LirScalarExpr > :: new ( 3 )
550+ . project ( vec ! [ 2 ] )
551+ . into_plan ( )
552+ . expect ( "project mfp" ) ;
553+ let stream = collection. inner ;
554+ let scope = stream. scope ( ) ;
555+ let ( oks, _errs) = flat_map_stage (
556+ stream,
557+ scope,
558+ exprs,
559+ func,
560+ mfp,
561+ Antichain :: new ( ) ,
562+ usize:: MAX ,
563+ ) ;
564+ let captured = oks. capture ( ) ;
565+ // Both rows at t=0: generate_series(1, 2) -> {1, 2},
566+ // generate_series(1, 3) -> {1, 2, 3}. Generated 1 and 2 appear on
567+ // both, so they must fold to a diff of two.
568+ input. advance_to ( Timestamp :: from ( 0_u64 ) ) ;
569+ input. update ( input_row ( 2 ) , Diff :: ONE ) ;
570+ input. update ( input_row ( 3 ) , Diff :: ONE ) ;
571+ input. advance_to ( Timestamp :: from ( 1_u64 ) ) ;
572+ input. flush ( ) ;
573+ captured
574+ } )
575+ } ) ;
576+
577+ let updates = extract_sorted_columns ( captured) ;
578+ let expected = vec ! [
579+ (
580+ Row :: pack_slice( & [ Datum :: Int64 ( 1 ) ] ) ,
581+ Timestamp :: from( 0_u64 ) ,
582+ Diff :: from( 2 ) ,
583+ ) ,
584+ (
585+ Row :: pack_slice( & [ Datum :: Int64 ( 2 ) ] ) ,
586+ Timestamp :: from( 0_u64 ) ,
587+ Diff :: from( 2 ) ,
588+ ) ,
589+ (
590+ Row :: pack_slice( & [ Datum :: Int64 ( 3 ) ] ) ,
591+ Timestamp :: from( 0_u64 ) ,
592+ Diff :: ONE ,
593+ ) ,
594+ ] ;
595+ assert_eq ! ( updates, expected) ;
506596 }
507597}
0 commit comments