use timely::dataflow::Scope;
use crate::{Collection, ExchangeData, Hashable};
use crate::difference::Semigroup;
use crate::Data;
use crate::lattice::Lattice;
use crate::trace::{Batcher, Builder};
impl<G, D, R> Collection<G, D, R>
where
G: Scope,
G::Timestamp: Data+Lattice,
D: ExchangeData+Hashable,
R: Semigroup+ExchangeData,
{
pub fn consolidate(&self) -> Self {
use crate::trace::implementations::KeySpine;
self.consolidate_named::<KeySpine<_,_,_>>("Consolidate")
}
pub fn consolidate_named<Tr>(&self, name: &str) -> Self
where
Tr: crate::trace::Trace<KeyOwned = D,ValOwned = (),Time=G::Timestamp,Diff=R>+'static,
Tr::Batch: crate::trace::Batch,
Tr::Batcher: Batcher<Item = ((D,()),G::Timestamp,R), Time = G::Timestamp>,
Tr::Builder: Builder<Item = ((D,()),G::Timestamp,R), Time = G::Timestamp>,
{
use crate::operators::arrange::arrangement::Arrange;
use crate::trace::cursor::MyTrait;
self.map(|k| (k, ()))
.arrange_named::<Tr>(name)
.as_collection(|d, _| d.into_owned())
}
pub fn consolidate_stream(&self) -> Self {
use timely::dataflow::channels::pact::Pipeline;
use timely::dataflow::operators::Operator;
use crate::collection::AsCollection;
self.inner
.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_vec(&mut vector);
})
}
})
.as_collection()
}
}