Skip to content
Open
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
122 changes: 61 additions & 61 deletions diagnostics/src/logging.rs
Original file line number Diff line number Diff line change
Expand Up @@ -448,46 +448,46 @@ 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();

for (event_time, event) in data.drain(..) {
match event {
TimelyEvent::Operates(e) => {
ops.give(((e.id, e.name.clone(), e.addr.clone()), ts, 1i64));
state.operators.insert(e.id, e);
}
TimelyEvent::Shutdown(e) => {
if let Some(op) = state.operators.remove(&e.id) {
ops.give(((op.id, op.name, op.addr), ts, -1i64));
if let Some(&ts) = cap.least() {
for (event_time, event) in data.drain(..) {
match event {
TimelyEvent::Operates(e) => {
ops.give(((e.id, e.name.clone(), e.addr.clone()), ts, 1i64));
state.operators.insert(e.id, e);
}
}
TimelyEvent::Channels(e) => {
chs.give((
(e.id, e.scope_addr.clone(), e.source, e.target),
ts,
1i64,
));
}
TimelyEvent::Schedule(e) => match e.start_stop {
StartStop::Start => {
state.schedule_starts.insert(e.id, event_time);
TimelyEvent::Shutdown(e) => {
if let Some(op) = state.operators.remove(&e.id) {
ops.give(((op.id, op.name, op.addr), ts, -1i64));
}
}
StartStop::Stop => {
if let Some(start) = state.schedule_starts.remove(&e.id) {
let elapsed_ns =
event_time.saturating_sub(start).as_nanos() as i64;
if elapsed_ns > 0 {
els.give((e.id, ts, elapsed_ns));
TimelyEvent::Channels(e) => {
chs.give((
(e.id, e.scope_addr.clone(), e.source, e.target),
ts,
1i64,
));
}
TimelyEvent::Schedule(e) => match e.start_stop {
StartStop::Start => {
state.schedule_starts.insert(e.id, event_time);
}
StartStop::Stop => {
if let Some(start) = state.schedule_starts.remove(&e.id) {
let elapsed_ns =
event_time.saturating_sub(start).as_nanos() as i64;
if elapsed_ns > 0 {
els.give((e.id, ts, elapsed_ns));
}
}
}
},
TimelyEvent::Messages(e) => {
if e.is_send {
msgs.give((e.channel, ts, e.record_count as i64));
}
}
},
TimelyEvent::Messages(e) => {
if e.is_send {
msgs.give((e.channel, ts, e.record_count as i64));
}
_ => {}
}
_ => {}
}
}
});
Expand Down Expand Up @@ -587,40 +587,40 @@ 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();

for (_event_time, event) in data.drain(..) {
match event {
DifferentialEvent::Batch(e) => {
bat.give((e.operator, ts, 1i64));
rec.give((e.operator, ts, e.length as i64));
}
DifferentialEvent::Merge(e) => {
if let Some(complete) = e.complete {
if let Some(&ts) = cap.least() {
for (_event_time, event) in data.drain(..) {
match event {
DifferentialEvent::Batch(e) => {
bat.give((e.operator, ts, 1i64));
rec.give((e.operator, ts, e.length as i64));
}
DifferentialEvent::Merge(e) => {
if let Some(complete) = e.complete {
bat.give((e.operator, ts, -1i64));
let diff = complete as i64 - (e.length1 + e.length2) as i64;
if diff != 0 {
rec.give((e.operator, ts, diff));
}
}
}
DifferentialEvent::Drop(e) => {
bat.give((e.operator, ts, -1i64));
let diff = complete as i64 - (e.length1 + e.length2) as i64;
let diff = -(e.length as i64);
if diff != 0 {
rec.give((e.operator, ts, diff));
}
}
}
DifferentialEvent::Drop(e) => {
bat.give((e.operator, ts, -1i64));
let diff = -(e.length as i64);
if diff != 0 {
rec.give((e.operator, ts, diff));
DifferentialEvent::TraceShare(e) => {
shr.give((e.operator, ts, e.diff as i64));
}
DifferentialEvent::Batcher(e) => {
b_rec.give((e.operator, ts, e.records_diff as i64));
b_sz.give((e.operator, ts, e.size_diff as i64));
b_cap.give((e.operator, ts, e.capacity_diff as i64));
b_alloc.give((e.operator, ts, e.allocations_diff as i64));
}
_ => {}
}
DifferentialEvent::TraceShare(e) => {
shr.give((e.operator, ts, e.diff as i64));
}
DifferentialEvent::Batcher(e) => {
b_rec.give((e.operator, ts, e.records_diff as i64));
b_sz.give((e.operator, ts, e.size_diff as i64));
b_cap.give((e.operator, ts, e.capacity_diff as i64));
b_alloc.give((e.operator, ts, e.allocations_diff as i64));
}
_ => {}
}
}
});
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
79 changes: 60 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 Expand Up @@ -1252,6 +1257,42 @@ pub mod vec {
}
}

/// Routes the records of one message to singleton capabilities when the message's stamp has
/// several elements.
///
/// A batch is shipped under the set of capabilities it retires, so a message of batches can
/// carry several timestamps. Records each have one time, and an operator over records reads a
/// message's time as one (`InputCapability::time`); an operator that turns batches into records
/// keeps that true by sending each record under the first element of the stamp at or before its
/// time. The elements cover every record: a batch's times are at or beyond one of the
/// capabilities it retired under.
pub struct StampRouter<T: Timestamp> {
caps: Vec<timely::dataflow::operators::Capability<T>>,
}

impl<T: Timestamp> StampRouter<T> {
/// One capability per element of the message's stamp, for output `port`.
pub fn new(cap: &timely::dataflow::operators::InputCapability<T>, port: usize) -> Self {
Self { caps: cap.stamp().iter().map(|t| cap.delayed(t, port)).collect() }
}
/// One capability per element of a set.
pub fn from_set(set: &timely::dataflow::operators::CapabilitySet<T>) -> Self {
Self { caps: set.iter().cloned().collect() }
}
/// The capabilities, in the order `index` names them.
pub fn capabilities(&self) -> &[timely::dataflow::operators::Capability<T>] {
&self.caps
}
/// Which capability a record at `time` goes under.
pub fn index(&self, time: &T) -> usize {
self.caps.iter().position(|c| c.time().less_equal(time)).expect("a record's time is at or beyond an element of its message's stamp")
}
/// Empty buffers, one per capability, to route records into and then give under each.
pub fn buffers<D>(&self) -> Vec<Vec<D>> {
(0..self.caps.len()).map(|_| Vec::new()).collect()
}
}

/// Conversion to a differential dataflow Collection.
pub trait AsCollection<'scope, T: Timestamp, C> {
/// Converts the type to a differential dataflow collection.
Expand Down
55 changes: 36 additions & 19 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 Expand Up @@ -167,20 +174,30 @@ where
.unary::<Builder<U>, _, _, _>(Pipeline, "AsRecordedUpdates", |_, _| {
move |input, output| {
input.for_each(|time, batches| {
let mut session = output.session_with_builder(&time);
for batch in batches.drain(..) {
let Some(batch) = batch.inner else { continue };
let mut cursor = batch.cursor();
while cursor.key_valid(&batch) {
while cursor.val_valid(&batch) {
let key = cursor.key(&batch);
let val = cursor.val(&batch);
cursor.map_times(&batch, |time, diff| {
session.give((key, val, time, diff));
});
cursor.step_val(&batch);
// A message of batches may carry several timestamps; each record goes out
// under the one at or before its time (a pass per element, the rare case),
// so the collection's messages carry one each.
let router = crate::collection::StampRouter::new(&time, 0);
let batches: Vec<_> = batches.drain(..).filter_map(|b| b.inner).collect();
let mut owned_time = U::Time::default();
for (index, cap) in router.capabilities().iter().enumerate() {
let mut session = output.session_with_builder(cap);
for batch in batches.iter() {
let mut cursor = batch.cursor();
while cursor.key_valid(batch) {
while cursor.val_valid(batch) {
let key = cursor.key(batch);
let val = cursor.val(batch);
cursor.map_times(batch, |time, diff| {
columnar::Columnar::copy_from(&mut owned_time, time);
if router.index(&owned_time) == index {
session.give((key, val, time, diff));
}
});
cursor.step_val(batch);
}
cursor.step_key(batch);
}
cursor.step_key(&batch);
}
}
});
Expand Down
Loading