diff --git a/differential-dataflow/examples/spines.rs b/differential-dataflow/examples/spines.rs index 6bfe3b224..ff17a281a 100644 --- a/differential-dataflow/examples/spines.rs +++ b/differential-dataflow/examples/spines.rs @@ -52,15 +52,14 @@ fn main() { let mut probe = Handle::new(); let mut workload: Workload = worker.dataflow(|scope| { - use differential_dataflow::operators::arrange::Arrange; match mode.as_str() { "key" => { use differential_dataflow::trace::implementations::ord_neu::{OrdKeyBatcher, VecOrdKeyBuilder, OrdKeySpine}; let (data_input, data) = scope.new_collection::(); let (keys_input, keys) = scope.new_collection::(); - let data = data.arrange::, VecOrdKeyBuilder, OrdKeySpine>(); - let keys = keys.arrange::, VecOrdKeyBuilder, OrdKeySpine>(); + let data = data.map(|k| (k, ())).arrange::, VecOrdKeyBuilder, OrdKeySpine>(); + let keys = keys.map(|k| (k, ())).arrange::, VecOrdKeyBuilder, OrdKeySpine>(); keys.join_core(data, |_k, &(), &()| Option::<()>::None) .probe_with(&mut probe); Workload { data_input, keys_input } diff --git a/differential-dataflow/src/collection.rs b/differential-dataflow/src/collection.rs index 53ce59fdb..bab5ae4a6 100644 --- a/differential-dataflow/src/collection.rs +++ b/differential-dataflow/src/collection.rs @@ -971,7 +971,6 @@ pub mod vec { Bu: crate::trace::Builder, Output: Into>, F: Fn(BatchKey<'_, Tr>, BatchVal<'_, Tr>) -> D + 'static, { - use crate::operators::arrange::arrangement::Arrange; self.map(|k| (k, ())) .arrange_named::(name) .as_collection(reify) @@ -1024,42 +1023,40 @@ pub mod vec { use crate::trace::implementations::{ValSpine, ValBatcher, ValBuilder}; use crate::trace::implementations::{KeySpine, KeyBatcher, KeyBuilder}; use crate::trace::implementations::ContainerChunker; - use crate::operators::arrange::Arrange; - - impl<'scope, T, K, V, R> Arrange<'scope, T, Vec<((K, V), T, R)>> for Collection<'scope, T, (K, V), R> + impl<'scope, T, K, V, R> Collection<'scope, T, (K, V), R> where T: Timestamp + Lattice, K: crate::ExchangeData + Hashable, V: crate::ExchangeData, R: crate::ExchangeData + Semigroup, { - fn arrange_named(self, name: &str) -> Arranged<'scope, TraceAgent> + /// Arranges updates into a shared trace, exchanged by a hash of the key. + /// + /// The batcher's output container must equal the stream container; the default chunker + /// only consolidates same-type containers. For chunker setups that convert between + /// container types (e.g. columnar layouts), call + /// [`arrange_core`](crate::operators::arrange::arrangement::arrange_core) directly. + pub fn arrange(self) -> Arranged<'scope, TraceAgent> where Ba: crate::trace::Batcher, Time=T> + 'static, Bu: crate::trace::Builder, Output: Into>, Tr: crate::trace::Trace + 'static, { - let exchange = timely::dataflow::channels::pact::Exchange::new(move |update: &((K,V),T,R)| (update.0).0.hashed().into()); - crate::operators::arrange::arrangement::arrange_core::<_, _, ContainerChunker>, Ba, Bu, _>(self.inner, exchange, name) + self.arrange_named::("Arrange") } - } - impl<'scope, T, K: crate::ExchangeData+Hashable, R: crate::ExchangeData+Semigroup> Arrange<'scope, T, Vec<((K, ()), T, R)>> for Collection<'scope, T, K, R> - where - T: Timestamp + Lattice + Ord, - { - fn arrange_named(self, name: &str) -> Arranged<'scope, TraceAgent> + /// As [`Collection::arrange`] but with the ability to name the operator. + pub fn arrange_named(self, name: &str) -> Arranged<'scope, TraceAgent> where - Ba: crate::trace::Batcher, Time=T> + 'static, - Bu: crate::trace::Builder, Output: Into>, + Ba: crate::trace::Batcher, Time=T> + 'static, + Bu: crate::trace::Builder, Output: Into>, Tr: crate::trace::Trace + 'static, { - let exchange = timely::dataflow::channels::pact::Exchange::new(move |update: &((K,()),T,R)| (update.0).0.hashed().into()); - crate::operators::arrange::arrangement::arrange_core::<_, _, ContainerChunker>, Ba, Bu, _>(self.map(|k| (k, ())).inner, exchange, name) + let exchange = timely::dataflow::channels::pact::Exchange::new(move |update: &((K,V),T,R)| (update.0).0.hashed().into()); + crate::operators::arrange::arrangement::arrange_core::<_, _, ContainerChunker>, Ba, Bu, _>(self.inner, exchange, name) } } - impl<'scope, T, K: crate::ExchangeData+Hashable, V: crate::ExchangeData, R: crate::ExchangeData+Semigroup> Collection<'scope, T, (K,V), R> where T: Timestamp + Lattice + Ord, diff --git a/differential-dataflow/src/operators/arrange/arrangement.rs b/differential-dataflow/src/operators/arrange/arrangement.rs index 97d81f7dd..795411c11 100644 --- a/differential-dataflow/src/operators/arrange/arrangement.rs +++ b/differential-dataflow/src/operators/arrange/arrangement.rs @@ -297,33 +297,6 @@ impl<'scope, Tr: TraceReader> Arranged<'scope, Tr> { } } -/// A type that can be arranged as if a collection of updates. -pub trait Arrange<'scope, T: Timestamp+Lattice, C> : Sized { - /// Arranges updates into a shared trace. - /// - /// The batcher's output container must equal the stream container `C`; the default - /// chunker only consolidates same-type containers. For chunker setups that convert - /// between container types (e.g. columnar layouts), call [`arrange_core`] directly. - fn arrange(self) -> Arranged<'scope, TraceAgent> - where - Ba: Batcher + 'static, - Bu: Builder>, - Tr: Trace + 'static, - { - self.arrange_named::("Arrange") - } - - /// Arranges updates into a shared trace, with a supplied name. - /// - /// See [`Arrange::arrange`] for constraints on the batcher's output container. - fn arrange_named(self, name: &str) -> Arranged<'scope, TraceAgent> - where - Ba: Batcher + 'static, - Bu: Builder>, - Tr: Trace + 'static, - ; -} - /// Arranges a stream of updates by a key, configured with a name and a parallelization contract. /// /// This operator arranges a stream of values into a shared trace, whose contents it maintains. diff --git a/differential-dataflow/src/operators/arrange/mod.rs b/differential-dataflow/src/operators/arrange/mod.rs index f0afe1781..7652f6ec6 100644 --- a/differential-dataflow/src/operators/arrange/mod.rs +++ b/differential-dataflow/src/operators/arrange/mod.rs @@ -71,4 +71,4 @@ pub mod upsert; pub use self::writer::TraceWriter; pub use self::agent::{TraceAgent, ShutdownButton}; -pub use self::arrangement::{Arranged, Arrange}; +pub use self::arrangement::Arranged; diff --git a/experiments/src/bin/deals.rs b/experiments/src/bin/deals.rs index f20e3678f..045dd623a 100644 --- a/experiments/src/bin/deals.rs +++ b/experiments/src/bin/deals.rs @@ -4,10 +4,9 @@ use differential_dataflow::input::Input; use differential_dataflow::VecCollection; use differential_dataflow::operators::*; -use differential_dataflow::trace::implementations::{ValSpine, KeySpine, KeyBatcher, KeyBuilder, ValBatcher, ValBuilder}; +use differential_dataflow::trace::implementations::{ValSpine, ValBatcher, ValBuilder}; use differential_dataflow::operators::arrange::TraceAgent; use differential_dataflow::operators::arrange::Arranged; -use differential_dataflow::operators::arrange::Arrange; use differential_dataflow::operators::iterate::Variable; use differential_dataflow::lattice::Lattice; use differential_dataflow::difference::Present; @@ -97,7 +96,7 @@ fn tc<'s, T: timely::progress::Timestamp + Lattice + Default + timely::order::Em .arrange::, ValBuilder<_,_,_,_>, ValSpine<_,_,_,_>>() .join_core(edges.clone(), |_y,&x,&z| Some((x, z))) .concat(edges.as_collection(|&k,&v| (k,v))) - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ; @@ -127,7 +126,7 @@ fn sg<'s, T: timely::progress::Timestamp + Lattice + Default + timely::order::Em .arrange::, ValBuilder<_,_,_,_>, ValSpine<_,_,_,_>>() .join_core(edges, |_,&x,&z| Some((x, z))) .concat(peers) - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ; diff --git a/experiments/src/bin/graspan1.rs b/experiments/src/bin/graspan1.rs index ce04f2202..b7b7cb356 100644 --- a/experiments/src/bin/graspan1.rs +++ b/experiments/src/bin/graspan1.rs @@ -7,7 +7,6 @@ use differential_dataflow::difference::Present; use differential_dataflow::input::Input; use differential_dataflow::trace::implementations::{ValBatcher, ValBuilder, ValSpine}; use differential_dataflow::operators::*; -use differential_dataflow::operators::arrange::Arrange; use differential_dataflow::operators::iterate::Variable; type Node = u32; @@ -46,6 +45,7 @@ fn main() { let next = labels_collection.join_core(edges, |_b, a, c| Some((*c, *a))) .concat(nodes) + .map(|k| (k, ())) .arrange::, ValBuilder<_,_,_,_>, ValSpine<_,_,_,_>>() // .distinct_total_core::(); .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }); diff --git a/experiments/src/bin/graspan2.rs b/experiments/src/bin/graspan2.rs index fcd05ee65..07783d244 100644 --- a/experiments/src/bin/graspan2.rs +++ b/experiments/src/bin/graspan2.rs @@ -8,8 +8,7 @@ use differential_dataflow::operators::iterate::Variable; use differential_dataflow::VecCollection; use differential_dataflow::input::Input; use differential_dataflow::operators::*; -use differential_dataflow::operators::arrange::Arrange; -use differential_dataflow::trace::implementations::{ValSpine, KeySpine, ValBatcher, KeyBatcher, ValBuilder, KeyBuilder}; +use differential_dataflow::trace::implementations::{ValSpine, ValBatcher, ValBuilder}; use differential_dataflow::difference::Present; type Node = u32; @@ -85,7 +84,7 @@ fn unoptimized() { let value_flow_next = value_flow_next - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() // .distinct_total_core::() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ; @@ -99,7 +98,7 @@ fn unoptimized() { let memory_alias_next: VecCollection<_,_,Present> = memory_alias_next - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() // .distinct_total_core::() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ; @@ -199,7 +198,7 @@ fn optimized() { .arrange::, ValBuilder<_,_,_,_>, ValSpine<_,_,_,_>>() .join_core(value_flow_arranged, |_,&a,&b| Some((a,b))) .concat(nodes.map(|n| (n,n))) - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() // .distinct_total_core::() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ; @@ -224,7 +223,7 @@ fn optimized() { .arrange::, ValBuilder<_,_,_,_>, ValSpine<_,_,_,_>>() .join_core(value_flow_deref, |_y,&a,&b| Some((a,b))) .concat(memory_alias_next) - .arrange::, KeyBuilder<_,_,_>, KeySpine<_,_,_>>() + .arrange_by_self() // .distinct_total_core::() .threshold_semigroup(|_,_,x: Option<&Present>| if x.is_none() { Some(Present) } else { None }) ;