1use std::collections::VecDeque;
34
35use columnar::bytes::indexed;
36use columnar::{Borrow, Columnar, Push};
37use differential_dataflow::Data;
38use differential_dataflow::consolidation::consolidate_updates;
39use differential_dataflow::difference::Semigroup;
40use timely::container::{ContainerBuilder, PushInto};
41
42use crate::columnar::Column;
43
44const STAGING_BUFFER_BYTES: usize = 8 * 1024;
47
48fn default_staging_cap<D, T, R>() -> usize {
52 let elem = std::mem::size_of::<(D, T, R)>().max(1);
53 (2 * STAGING_BUFFER_BYTES / elem).max(2)
55}
56const OUTPUT_TARGET_WORDS: usize = 1 << 18;
59const FLUSH_THRESHOLD_WORDS: usize = OUTPUT_TARGET_WORDS - OUTPUT_TARGET_WORDS / 10;
63const DRAIN_CHUNK_ROWS: usize = 16;
69
70pub struct ConsolidatingColumnBuilder<D, T, R>
75where
76 D: Columnar,
77 T: Columnar,
78 R: Columnar,
79{
80 staging: Vec<(D, T, R)>,
82 staging_cap: usize,
84 cur_d: D::Container,
86 cur_t: T::Container,
87 cur_r: R::Container,
88 cur_len: usize,
90 pending: VecDeque<Column<(D, T, R)>>,
92 finished: Option<Column<(D, T, R)>>,
94}
95
96impl<D, T, R> Default for ConsolidatingColumnBuilder<D, T, R>
97where
98 D: Columnar,
99 T: Columnar,
100 R: Columnar,
101{
102 fn default() -> Self {
103 let cap = default_staging_cap::<D, T, R>();
104 Self {
105 staging: Vec::with_capacity(cap),
108 staging_cap: cap,
109 cur_d: D::Container::default(),
110 cur_t: T::Container::default(),
111 cur_r: R::Container::default(),
112 cur_len: 0,
113 pending: VecDeque::new(),
114 finished: None,
115 }
116 }
117}
118
119impl<D, T, R> ConsolidatingColumnBuilder<D, T, R>
120where
121 D: Data + Columnar,
122 T: Data + Columnar,
123 R: Semigroup + Columnar + 'static,
124 (D, T, R): Columnar<Container = (D::Container, T::Container, R::Container)>,
125{
126 #[cold]
132 fn consolidate_and_drain(&mut self, grain: usize) {
133 consolidate_updates(&mut self.staging);
134 let drain_n = (self.staging.len() / grain) * grain;
135 if drain_n == 0 {
136 return;
137 }
138
139 let mut consumed = 0;
142 while consumed < drain_n {
143 let take = (drain_n - consumed).min(DRAIN_CHUNK_ROWS);
144 let head = &self.staging[consumed..consumed + take];
145 for (d, _, _) in head {
146 self.cur_d.push(d);
147 }
148 for (_, t, _) in head {
149 self.cur_t.push(t);
150 }
151 for (_, _, r) in head {
152 self.cur_r.push(r);
153 }
154 self.cur_len += take;
155 consumed += take;
156
157 let words = {
158 let view = (
159 self.cur_d.borrow(),
160 self.cur_t.borrow(),
161 self.cur_r.borrow(),
162 );
163 indexed::length_in_words(&view)
164 };
165 if words >= FLUSH_THRESHOLD_WORDS {
166 self.flush_aligned();
167 }
168 }
169 self.staging.drain(..consumed);
170 }
171
172 #[cold]
176 fn flush_aligned(&mut self) {
177 if self.cur_len == 0 {
178 return;
179 }
180 let cur: <(D, T, R) as Columnar>::Container = (
181 std::mem::take(&mut self.cur_d),
182 std::mem::take(&mut self.cur_t),
183 std::mem::take(&mut self.cur_r),
184 );
185 self.cur_len = 0;
186
187 let mut buf: Vec<u64> = Vec::with_capacity(indexed::length_in_words(&cur.borrow()));
188 indexed::encode(&mut buf, &cur.borrow());
189 self.pending.push_back(Column::Align(buf));
190 }
191}
192
193impl<D, T, R> PushInto<(D, T, R)> for ConsolidatingColumnBuilder<D, T, R>
194where
195 D: Data + Columnar,
196 T: Data + Columnar,
197 R: Semigroup + Columnar + 'static,
198 (D, T, R): Columnar<Container = (D::Container, T::Container, R::Container)>,
199{
200 #[inline]
202 fn push_into(&mut self, item: (D, T, R)) {
203 self.staging.push(item);
204 if self.staging.len() == self.staging_cap {
205 self.consolidate_and_drain(self.staging_cap / 2);
206 }
207 }
208}
209
210impl<D, T, R> ContainerBuilder for ConsolidatingColumnBuilder<D, T, R>
211where
212 D: Data + Columnar,
213 T: Data + Columnar,
214 R: Semigroup + Columnar + 'static,
215 (D, T, R): Columnar<Container = (D::Container, T::Container, R::Container)>,
216 <(D, T, R) as Columnar>::Container: Clone,
217{
218 type Container = Column<(D, T, R)>;
219
220 #[inline]
221 fn extract(&mut self) -> Option<&mut Self::Container> {
222 if let Some(c) = self.pending.pop_front() {
223 self.finished = Some(c);
224 self.finished.as_mut()
225 } else {
226 None
227 }
228 }
229
230 #[inline]
231 fn finish(&mut self) -> Option<&mut Self::Container> {
232 if !self.staging.is_empty() {
233 self.consolidate_and_drain(1);
235 }
236 if self.cur_len > 0 {
238 let cur: <(D, T, R) as Columnar>::Container = (
239 std::mem::take(&mut self.cur_d),
240 std::mem::take(&mut self.cur_t),
241 std::mem::take(&mut self.cur_r),
242 );
243 self.cur_len = 0;
244 self.pending.push_back(Column::Typed(cur));
245 }
246 self.extract()
247 }
248}
249
250#[cfg(test)]
251mod tests {
252 use columnar::Index;
253 use columnar::Len;
254 use timely::container::{ContainerBuilder, PushInto};
255
256 use super::*;
257
258 fn rows(column: &Column<(u64, u64, i64)>) -> Vec<(u64, u64, i64)> {
260 let borrow = column.borrow();
261 (0..borrow.len())
262 .map(|i| {
263 let r = borrow.get(i);
264 (*r.0, *r.1, *r.2)
265 })
266 .collect()
267 }
268
269 fn drain(mut builder: ConsolidatingColumnBuilder<u64, u64, i64>) -> Vec<(u64, u64, i64)> {
271 let mut out: Vec<(u64, u64, i64)> = Vec::new();
272 while let Some(c) = builder.extract() {
273 for r in rows(c) {
274 out.push(r);
275 }
276 }
277 if let Some(c) = builder.finish() {
278 for r in rows(c) {
279 out.push(r);
280 }
281 }
282 out
283 }
284
285 #[mz_ore::test]
286 fn empty_finish_yields_none() {
287 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
288 assert!(builder.extract().is_none());
289 assert!(builder.finish().is_none());
290 }
291
292 #[mz_ore::test]
293 fn single_push_finish_yields_one() {
294 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
295 builder.push_into((1u64, 0u64, 1i64));
296 let column = builder.finish().expect("one container");
297 assert_eq!(rows(column), vec![(1, 0, 1)]);
298 assert!(builder.finish().is_none());
299 }
300
301 #[mz_ore::test]
302 fn consolidates_on_threshold() {
303 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
304 let cap = default_staging_cap::<u64, u64, i64>();
308 for _ in 0..(cap * 4) {
309 builder.push_into((7u64, 0u64, 1i64));
310 builder.push_into((7u64, 0u64, -1i64));
311 }
312 assert!(drain(builder).is_empty());
313 }
314
315 #[mz_ore::test]
316 fn cross_batch_consolidation() {
317 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
321 let n: i64 = 100_000;
322 for _ in 0..n {
323 builder.push_into((42u64, 0u64, 1i64));
324 }
325 let out = drain(builder);
326 assert_eq!(out, vec![(42, 0, n)]);
327 }
328
329 #[mz_ore::test]
330 fn multiple_distinct_keys() {
331 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
332 builder.push_into((1u64, 0u64, 1i64));
333 builder.push_into((2u64, 0u64, 1i64));
334 builder.push_into((1u64, 0u64, 1i64));
335 let mut out = drain(builder);
336 out.sort();
337 assert_eq!(out, vec![(1, 0, 2), (2, 0, 1)]);
338 }
339
340 #[mz_ore::test]
341 #[cfg_attr(miri, ignore)] fn emits_multiple_containers() {
343 let mut builder: ConsolidatingColumnBuilder<u64, u64, i64> = Default::default();
344 let n: u64 = 300_000;
348 for d in 0..n {
349 builder.push_into((d, 0u64, 1i64));
350 }
351
352 let mut containers = 0;
353 let mut out: Vec<(u64, u64, i64)> = Vec::new();
354 while let Some(c) = builder.extract() {
355 containers += 1;
356 for r in rows(c) {
357 out.push(r);
358 }
359 }
360 if let Some(c) = builder.finish() {
361 containers += 1;
362 for r in rows(c) {
363 out.push(r);
364 }
365 }
366 assert!(
367 containers > 1,
368 "expected multiple containers, got {containers}"
369 );
370 out.sort();
371 let expected: Vec<(u64, u64, i64)> = (0..n).map(|d| (d, 0u64, 1i64)).collect();
372 assert_eq!(out, expected);
373 }
374}