From 56b92e85b035de928c3421118685c31851abd33c Mon Sep 17 00:00:00 2001 From: Moritz Hoffmann Date: Wed, 15 May 2024 12:09:36 -0400 Subject: [PATCH] Use consolidating builder in consolidate_stream This switches from consolidating each input individually to consolidating at the output. The benefits are that it can yield better consolidation performance because it can consolidate across input containers instead of only being able to consolidate each individual input container. Signed-off-by: Moritz Hoffmann --- src/operators/consolidate.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/operators/consolidate.rs b/src/operators/consolidate.rs index ac0ef7f8f..2235c4df9 100644 --- a/src/operators/consolidate.rs +++ b/src/operators/consolidate.rs @@ -9,6 +9,7 @@ use timely::dataflow::Scope; use crate::{Collection, ExchangeData, Hashable}; +use crate::consolidation::ConsolidatingContainerBuilder; use crate::difference::Semigroup; use crate::Data; @@ -92,14 +93,13 @@ where use crate::collection::AsCollection; self.inner - .unary(Pipeline, "ConsolidateStream", |_cap, _info| { + .unary::, _, _, _>(Pipeline, "ConsolidateStream", |_cap, _info| { let mut vector = Vec::new(); move |input, output| { input.for_each(|time, data| { data.swap(&mut vector); - crate::consolidation::consolidate_updates(&mut vector); - output.session(&time).give_container(&mut vector); + output.session_with_builder(&time).give_container(&mut vector); }) } })