Skip to content

Commit 2f94da1

Browse files
authored
start, emit an empty 'RecordedUpdates' in the columnar 'ValDistributor', when the input consolidates to nothing (#834)
1 parent aa96fb3 commit 2f94da1

2 files changed

Lines changed: 97 additions & 0 deletions

File tree

‎differential-dataflow/src/columnar/collection/exchange.rs‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ pub struct ValDistributor<U: Update, H> {
2222
marker: std::marker::PhantomData<U>,
2323
hashfunc: H,
2424
pre_lens: Vec<usize>,
25+
worker: usize,
2526
}
2627

2728
impl<U: Update, H: for<'a> FnMut(columnar::Ref<'a, U::Key>)->u64> Distributor<RecordedUpdates<U>> for ValDistributor<U, H> {
@@ -68,6 +69,21 @@ impl<U: Update, H: for<'a> FnMut(columnar::Ref<'a, U::Key>)->u64> Distributor<Re
6869
// Distribute the input's record count across non-empty outputs.
6970
let total_records = container.records;
7071
let non_empty: usize = outputs.iter().filter(|o| !o.keys.values.is_empty()).count();
72+
// N.B. An input whose updates all consolidated away needs to emit a message
73+
// carrying a `records` count for timely's bookkeeping.
74+
if non_empty == 0 && total_records > 0 {
75+
let mut recorded = RecordedUpdates::<U> {
76+
updates: Default::default(),
77+
records: total_records,
78+
consolidated: true,
79+
};
80+
// Push the empty update to the worker that produced the original
81+
// values so the send stays local. Not needed for correctness, but
82+
// a reasonable choice.
83+
Message::push_at(&mut recorded, time.clone(), &mut pushers[self.worker % pushers.len()]);
84+
return;
85+
}
86+
7187
let mut first_records = total_records.saturating_sub(non_empty.saturating_sub(1));
7288
for (pusher, output) in pushers.iter_mut().zip(outputs) {
7389
if !output.keys.values.is_empty() {
@@ -108,6 +124,7 @@ where
108124
marker: std::marker::PhantomData,
109125
hashfunc: self.hashfunc,
110126
pre_lens: Vec::new(),
127+
worker: worker.index(),
111128
};
112129
(Exchange::new(senders, distributor), LogPuller::new(receiver, worker.index(), identifier, logging.clone()))
113130
}
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
use std::time::{Duration, Instant};
2+
3+
use timely::dataflow::operators::{Input, Probe};
4+
use timely::dataflow::{InputHandle, ProbeHandle};
5+
use timely::Config;
6+
7+
use differential_dataflow::columnar::collection;
8+
use differential_dataflow::columnar::trace::{Batcher, Builder, Chunker, Spine};
9+
use differential_dataflow::operators::arrange::arrangement::arrange_core;
10+
11+
type Upd = (u64, (), u64, i64);
12+
13+
/// Feed the dataflow with a single (key, val) at the provided `diffs` the
14+
/// advance the frontier to 1.
15+
///
16+
/// Returns false if the dataflow hangs for longer than 5s.
17+
fn arrange_reaches_one(config: Config, diffs: &'static [i64]) -> bool {
18+
let guards = timely::execute(config, move |worker| {
19+
let index = worker.index();
20+
let mut probe = ProbeHandle::new();
21+
let mut input = <InputHandle<u64, collection::Builder<Upd>>>::new_with_builder();
22+
23+
worker.dataflow::<u64, _, _>(|scope| {
24+
let stream = scope.input_from(&mut input);
25+
let pact = collection::Pact {
26+
hashfunc: |k: columnar::Ref<'_, u64>| *k,
27+
};
28+
29+
arrange_core::<
30+
_,
31+
_,
32+
Chunker<Upd>,
33+
Batcher<u64, (), u64, i64>,
34+
Builder<u64, (), u64, i64>,
35+
Spine<u64, (), u64, i64>,
36+
>(stream, pact, "Arrange")
37+
.stream
38+
.probe_with(&mut probe);
39+
});
40+
41+
for &diff in diffs {
42+
input.send((index as u64, (), 0, diff));
43+
}
44+
input.advance_to(1);
45+
input.flush();
46+
47+
let start = Instant::now();
48+
while probe.less_than(&1) {
49+
worker.step();
50+
if start.elapsed() > Duration::from_secs(5) {
51+
for id in worker.installed_dataflows() {
52+
worker.drop_dataflow(id);
53+
}
54+
return false;
55+
}
56+
}
57+
true
58+
})
59+
.expect("timely execute");
60+
61+
guards.join().into_iter().all(|r| r.unwrap_or(false))
62+
}
63+
64+
#[test]
65+
fn columnar_exchange_net_empty_container() {
66+
// [+1, -1] diffs should consolidate to nothing, but the exchange must
67+
// still deliver a message with `records = 2`.
68+
//
69+
// Previously the exchange did not emit a message and the dataflow froze.
70+
assert!(
71+
arrange_reaches_one(Config::process(2), &[1, -1]),
72+
"frontier stalled: record count lost for an all-cancelled container"
73+
);
74+
}
75+
76+
#[test]
77+
fn columnar_exchange_consolidated_container() {
78+
assert!(arrange_reaches_one(Config::process(2), &[1, 1, 1]));
79+
assert!(arrange_reaches_one(Config::process(3), &[2, -1, 1]));
80+
}

0 commit comments

Comments
 (0)