Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions diagnostics/src/logging.rs
Original file line number Diff line number Diff line change
Expand Up @@ -448,7 +448,7 @@ fn construct_timely<'scope>(
let mut chs = ch_act.session(&cap);
let mut els = el_act.session(&cap);
let mut msgs = msg_act.session(&cap);
let ts = *cap.time();
let Some(&ts) = cap.least() else { return };

for (event_time, event) in data.drain(..) {
match event {
Expand Down Expand Up @@ -587,7 +587,7 @@ fn construct_differential<'scope>(
let mut b_sz = bs_act.session(&cap);
let mut b_cap = bc_act.session(&cap);
let mut b_alloc = ba_act.session(&cap);
let ts = *cap.time();
let Some(&ts) = cap.least() else { return };

for (_event_time, event) in data.drain(..) {
match event {
Expand Down
8 changes: 4 additions & 4 deletions differential-dataflow/src/capture.rs
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,7 @@ pub mod source {
use std::rc::Rc;
use std::marker::{Send, Sync};
use std::sync::Arc;
use timely::dataflow::{Scope, Stream, operators::{Capability, CapabilitySet}};
use timely::dataflow::{Scope, Stream, operators::CapabilitySet};
use timely::dataflow::operators::generic::OutputBuilder;
use timely::progress::Timestamp;
use timely::scheduling::SyncActivator;
Expand Down Expand Up @@ -462,17 +462,17 @@ pub mod source {

// If the frontier changes we need a capability to express that.
// Any capability should work; the downstream listener doesn't care.
let mut capability: Option<Capability<T>> = None;
let mut capability: Option<CapabilitySet<T>> = None;

// Drain all relevant update counts in to the mutable antichain tracking its frontier.
counts.for_each(|cap, counts| {
updates_frontier.update_iter(counts.iter().cloned());
capability = Some(cap.retain(0));
capability = Some(cap.retain_stamp(0));
});
// Drain all progress statements into the queue out of which we will work.
input.for_each(|cap, progress| {
progress_queue.extend(progress.iter().map(|x| (x.1).clone()));
capability = Some(cap.retain(0));
capability = Some(cap.retain_stamp(0));
});

// Extract and act on actionable progress messages.
Expand Down
43 changes: 24 additions & 19 deletions differential-dataflow/src/collection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,15 +128,15 @@ impl<'scope, T: Timestamp, C: Container> Collection<'scope, T, C> {
/// scope.new_collection_from(1 .. 10).1
/// .map_in_place(|x| *x *= 2)
/// .filter(|x| x % 2 == 1)
/// .inspect_container(|event| println!("event: {:?}", event));
/// .inspect_core(|event| println!("event: {:?}", event));
/// });
/// ```
pub fn inspect_container<F>(self, func: F) -> Self
pub fn inspect_core<F>(self, func: F) -> Self
where
F: FnMut(Result<(&T, &C), &[T]>)+'static,
F: FnMut(Result<(&timely::progress::Stamp<T>, &C), &[T]>)+'static,
{
self.inner
.inspect_container(func)
.inspect_core(func)
.as_collection()
}
/// Attaches a timely dataflow probe to the output of a Collection.
Expand Down Expand Up @@ -235,9 +235,6 @@ impl<'scope, T: Timestamp, C: Container> Collection<'scope, T, C> {
///
/// # Examples
/// ```
/// use timely::dataflow::Scope;
/// use timely::dataflow::operators::{ToStream, Concat, Inspect, vec::BranchWhen};
///
/// use differential_dataflow::input::Input;
///
/// timely::example(|scope| {
Expand Down Expand Up @@ -562,17 +559,25 @@ pub mod vec {
/// ordered, they should have the same order or compare equal once `func` is applied to them (this
/// is because we advance the timely capability with the same logic, and it must remain `less_equal`
/// to all of the data timestamps).
pub fn delay<F>(self, func: F) -> Collection<'scope, T, D, R>
pub fn delay<F>(self, mut func: F) -> Collection<'scope, T, D, R>
where
T: Hash,
F: FnMut(&T) -> T + Clone + 'static,
F: FnMut(&T) -> T + 'static,
{
let mut func1 = func.clone();
let mut func2 = func.clone();

use timely::dataflow::channels::pact::Pipeline;
use timely::dataflow::operators::CapabilitySet;
self.inner
.delay_batch(move |x| func1(x))
.map_in_place(move |x| x.1 = func2(&x.1))
.unary(Pipeline, "Delay", move |_, _| move |input, output| {
input.for_each_stamp(|time, data| {
// Each element of the stamp is delayed by `func`, and each update's time likewise;
// by monotonicity the delayed capabilities still cover the delayed updates.
let caps: CapabilitySet<T> = time.stamp().iter().map(|t| time.delayed(&func(t), output.output_index())).collect();
let mut session = output.session(&caps);
for data in data {
for (_, t, _) in data.iter_mut() { *t = func(t); }
session.give_container(data);
}
});
})
.as_collection()
}

Expand Down Expand Up @@ -609,8 +614,8 @@ pub mod vec {
}
/// Applies a supplied function to each batch of updates.
///
/// This method is analogous to `inspect`, but operates on batches and reveals the timestamp of the
/// timely dataflow capability associated with the batch of updates. The observed batching depends
/// This method is analogous to `inspect`, but operates on batches and reveals the stamp under which
/// the batch of updates travels: the timestamps of the timely dataflow capabilities associated with it. The observed batching depends
/// on how the system executes, and may vary run to run.
///
/// # Examples
Expand All @@ -627,10 +632,10 @@ pub mod vec {
/// ```
pub fn inspect_batch<F>(self, mut func: F) -> Collection<'scope, T, D, R>
where
F: FnMut(&T, &[(D, T, R)])+'static,
F: FnMut(&timely::progress::Stamp<T>, &[(D, T, R)])+'static,
{
self.inner
.inspect_batch(move |time, data| func(time, data))
.inspect_core(move |event| if let Ok((stamp, data)) = event { func(stamp, data) })
.as_collection()
}

Expand Down
19 changes: 13 additions & 6 deletions differential-dataflow/src/columnar/collection/operators.rs
Original file line number Diff line number Diff line change
Expand Up @@ -116,12 +116,19 @@ where
move |_frontier| {
let mut output = output.activate();
op_input.for_each(|cap, data| {
// Truncate the capability's timestamp.
let mut new_time = cap.time().clone();
let mut vec = std::mem::take(&mut new_time.inner).into_inner();
vec.truncate(level - 1);
new_time.inner = PointStamp::new(vec);
let new_cap = cap.delayed(&new_time, 0);
// A message may carry several timestamps (a multi-element stamp): hold a
// capability for each, truncated exactly as the updates are.
let new_cap: timely::dataflow::operators::CapabilitySet<_> = cap
.stamp()
.iter()
.map(|t| {
let mut new_time = t.clone();
let mut vec = std::mem::take(&mut new_time.inner).into_inner();
vec.truncate(level - 1);
new_time.inner = PointStamp::new(vec);
cap.delayed(&new_time, 0)
})
.collect();
// Push updates with truncated times into the builder.
// The builder's form call on flush sorts and consolidates,
// handling the duplicate times that truncation can produce.
Expand Down
23 changes: 14 additions & 9 deletions differential-dataflow/src/dynamic/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ pub mod pointstamp;
use timely::order::Product;
use timely::progress::Timestamp;
use timely::dataflow::operators::generic::{OutputBuilder, builder_rc::OperatorBuilder};
use timely::dataflow::operators::CapabilitySet;
use timely::dataflow::channels::pact::Pipeline;
use timely::progress::Antichain;

Expand Down Expand Up @@ -47,17 +48,21 @@ where
builder.build(move |_capability| move |_frontier| {
let mut output = output.activate();
input.for_each(|cap, data| {
let mut new_time = cap.time().clone();
let mut vec = std::mem::take(&mut new_time.inner).into_inner();
vec.truncate(level - 1);
new_time.inner = PointStamp::new(vec);
let new_cap = cap.delayed(&new_time, 0);
for (_data, time, _diff) in data.iter_mut() {
let mut vec = std::mem::take(&mut time.inner).into_inner();
// A message may carry several timestamps (a multi-element stamp, e.g. late
// iterations of one epoch alongside early ones of the next): hold a capability
// for each, truncated exactly as the records are.
let truncate = |time: &Product<TOuter, PointStamp<T>>| {
let mut new_time = time.clone();
let mut vec = std::mem::take(&mut new_time.inner).into_inner();
vec.truncate(level - 1);
time.inner = PointStamp::new(vec);
new_time.inner = PointStamp::new(vec);
new_time
};
let caps: CapabilitySet<_> = cap.stamp().iter().map(|t| cap.delayed(&truncate(t), 0)).collect();
for (_data, time, _diff) in data.iter_mut() {
*time = truncate(time);
}
output.session(&new_cap).give_container(data);
output.session(&caps).give_container(data);
});
});

Expand Down
2 changes: 1 addition & 1 deletion differential-dataflow/src/operators/arrange/arrangement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -186,7 +186,7 @@
while let Some(key) = cursor.get_key(batch) {
while let Some(val) = cursor.get_val(batch) {
for datum in logic(key, val) {
cursor.map_times(batch, |time, diff| {

Check warning on line 189 in differential-dataflow/src/operators/arrange/arrangement.rs

View workflow job for this annotation

GitHub Actions / Cargo clippy

`time` shadows a previous, unrelated binding
session.give((datum.clone(), <BatchCursor<Tr> as Cursor>::owned_time(time), <BatchCursor<Tr> as Cursor>::owned_diff(diff)));
});
}
Expand Down Expand Up @@ -333,7 +333,7 @@

let mut reader: Option<TraceAgent<Tr>> = None;

// fabricate a data-parallel operator using the `unary_notify` pattern.
// fabricate a data-parallel operator that holds capabilities and consults its input frontier.
let reader_ref = &mut reader;
let scope = stream.scope();

Expand Down
4 changes: 2 additions & 2 deletions differential-dataflow/src/operators/arrange/upsert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ where
{
let mut reader: Option<TraceAgent<Tr>> = None;

// fabricate a data-parallel operator using the `unary_notify` pattern.
// fabricate a data-parallel operator that holds capabilities and consults its input frontier.
let stream = {

let reader = &mut reader;
Expand Down Expand Up @@ -180,7 +180,7 @@ where

// Stash capabilities and associated data (ordered by time).
input.for_each(|cap, data| {
capabilities.insert(cap.retain(0));
if let Some(cap) = cap.retain_least(0) { capabilities.insert(cap); }
for (key, val, time) in data.drain(..) {
priority_queue.push(std::cmp::Reverse((time, key, val)))
}
Expand Down
2 changes: 1 addition & 1 deletion differential-dataflow/src/operators/reduce.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ where
{
let mut result_trace = None;

// fabricate a data-parallel operator using the `unary_notify` pattern.
// fabricate a data-parallel operator that holds capabilities and consults its input frontier.
let stream = {

let mut source_trace = trace.trace;
Expand Down
5 changes: 4 additions & 1 deletion dogsdogsdogs/examples/dogsdogsdogs.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use timely::dataflow::operators::{ToStream, vec::{Map, Partition, count::Accumulate}, Inspect, Probe};
use timely::dataflow::operators::{ToStream, vec::{Map, Partition}, Inspect, Probe};
use timely::dataflow::operators::probe::Handle;
use differential_dataflow::{Collection, AsCollection};
use differential_dataflow::input::Input;
Expand Down Expand Up @@ -74,7 +74,10 @@ fn main() {
.concat(validate1)
// Delay updates to the payload time worked out while extending.
.inner.map(|((extended, payload), _time, r)| (extended, payload, r))
.as_collection()
.map(|_| ())
.count()
.inner
.inspect(move |x| println!("{:?}", x))
// .inspect(move |x| println!("{:?}:\t{:?}", timer.elapsed(), x))
.probe_with(&mut probe);
Expand Down
3 changes: 1 addition & 2 deletions dogsdogsdogs/examples/ngo.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
use std::hash::Hash;
use timely::dataflow::operators::*;
use timely::dataflow::operators::vec::count::Accumulate;

use differential_dataflow::VecCollection;
use differential_dataflow::lattice::Lattice;
Expand Down Expand Up @@ -31,7 +30,7 @@ fn main() {
println!("loaded {} nodes, {} edges", nodes, edges.len());

worker.dataflow::<(),_,_>(|scope| {
triangles(VecCollection::new(edges.to_stream(scope))).inner.count().inspect(|x| println!("{:?}", x));
triangles(VecCollection::new(edges.to_stream(scope))).map(|_| ()).count().inspect(|x| println!("{:?}", x));
});

}).unwrap();
Expand Down
3 changes: 1 addition & 2 deletions dogsdogsdogs/src/operators/half_join.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,9 +109,8 @@ where
let output: &mut OutputBuilderSession<'_, Tr::Time, NoopBuilder<C>> = output;

// Stage all arriving updates, retaining capabilities that cover them.
// TODO: Tolerate multi-capability inputs.
input1.for_each(|capability, data| {
caps.insert(capability.retain(0));
for cap in capability.retain_stamp(0).iter() { caps.insert(cap.clone()); }
batcher.insert(data);
});

Expand Down
5 changes: 2 additions & 3 deletions dogsdogsdogs/tests/wcoj_partial_order.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,9 +83,8 @@ fn triangles(deltas: &[((u32, u32), Time, isize)]) -> Vec<((u32, u32, u32), Time
.inner.map(|((d, payload), _time, r)| (d, payload, r)).as_collection();

let left = triangles
.inspect_batch(move |_t, xs| {
let mut v = found_outer.lock().unwrap();
v.extend(xs.iter().cloned());
.inspect(move |x| {
found_outer.lock().unwrap().push(x.clone());
})
.leave(scope);

Expand Down
4 changes: 2 additions & 2 deletions experiments/src/bin/graphs-static.rs
Original file line number Diff line number Diff line change
Expand Up @@ -182,15 +182,15 @@ fn connected_components<'s>(
let f_prop = labels.clone().join_core(forward, |_k,l,d| Some((*d,*l)));
let r_prop = labels.join_core(reverse, |_k,l,d| Some((*d,*l)));

use timely::dataflow::operators::vec::{Map, Delay};
use timely::dataflow::operators::vec::Map;
use timely::dataflow::operators::Concat;

// Records carry their advanced times; the downstream reduce acts on those, not on capabilities.
let result =
nodes
.inner
.map_in_place(|dtr| (dtr.1).inner = 256 * ((((::std::mem::size_of::<Node>() * 8) as u32) - (dtr.0).1.leading_zeros())))
.concat(inner_collection.filter(|_| false).inner)
.delay(|dtr,_| dtr.1.clone())
.as_collection()
.concat(f_prop)
.concat(r_prop)
Expand Down
18 changes: 13 additions & 5 deletions interactive/src/backend/corgi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -366,11 +366,19 @@ impl Backend for CorgiBackend {
builder.build(move |_capability| move |_frontier| {
let mut output = output.activate();
input.for_each(|cap, data| {
let mut new_time = cap.time().clone();
let mut v = std::mem::take(&mut new_time.inner).into_inner();
v.truncate(level - 1);
new_time.inner = PointStamp::new(v);
let new_cap = cap.delayed(&new_time, 0);
// A message may carry several timestamps (a multi-element stamp): hold a
// capability for each, truncated exactly as the rows are.
let new_cap: timely::dataflow::operators::CapabilitySet<_> = cap
.stamp()
.iter()
.map(|t| {
let mut new_time = t.clone();
let mut v = std::mem::take(&mut new_time.inner).into_inner();
v.truncate(level - 1);
new_time.inner = PointStamp::new(v);
cap.delayed(&new_time, 0)
})
.collect();
for t in data.times.iter_mut() {
let mut v = std::mem::take(&mut t.inner).into_inner();
v.truncate(level - 1);
Expand Down
Loading