1use std::collections::BTreeMap;
11use std::rc::Rc;
12use std::sync::{Arc, Weak};
13
14use differential_dataflow::difference::Semigroup;
15use differential_dataflow::lattice::Lattice;
16use differential_dataflow::operators::arrange::arrangement::arrange_core;
17use differential_dataflow::operators::arrange::{Arranged, TraceAgent};
18use differential_dataflow::trace::implementations::BatchContainer;
19use differential_dataflow::trace::implementations::spine_fueled::Spine;
20use differential_dataflow::trace::{Batch, Batcher, Builder, Trace, TraceReader};
21use differential_dataflow::{Collection, Data, ExchangeData, Hashable, VecCollection};
22use mz_compute_types::dyncfgs::{ENABLE_COLUMN_PAGED_BATCHER, ENABLE_COLUMNAR_MERGE_BATCHER};
23use mz_dyncfg::ConfigSet;
24use mz_row_spine::ArcBatch;
25use mz_timely_util::containers::HeapSize;
26use timely::Container;
27use timely::container::{ContainerBuilder, PushInto};
28use timely::dataflow::Stream;
29use timely::dataflow::channels::pact::{Exchange, ParallelizationContract, Pipeline};
30use timely::dataflow::operators::Operator;
31use timely::progress::Timestamp;
32
33use crate::logging::compute::{
34 ArrangementHeapAllocations, ArrangementHeapCapacity, ArrangementHeapSize,
35 ArrangementHeapSizeOperator, ComputeEvent, ComputeEventBuilder,
36};
37use crate::typedefs::{
38 KeyAgent, KeyValAgent, MzArrangeData, MzData, MzTimestamp, RowAgent, RowRowAgent, RowValAgent,
39};
40
41pub enum ArrangementBatcher {
48 Columnation,
51 Columnar,
54 Chunked,
59}
60
61impl ArrangementBatcher {
62 pub fn from_config(config: &ConfigSet) -> Self {
70 if ENABLE_COLUMN_PAGED_BATCHER.get(config) {
71 Self::Chunked
72 } else if ENABLE_COLUMNAR_MERGE_BATCHER.get(config) {
73 Self::Columnar
74 } else {
75 Self::Columnation
76 }
77 }
78}
79
80pub trait MzArrange<'scope>: MzArrangeCore<'scope> {
82 fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
88 where
89 Ba: Batcher<Time = Self::Timestamp> + 'static,
90 Chu: ContainerBuilder<Container = Ba::Output>
91 + for<'a> PushInto<&'a mut Self::Input>
92 + 'static,
93 Bu: Builder<Time = Self::Timestamp, Input = Ba::Output, Output = Tr::Batch>,
94 Tr: Trace + TraceReader<Time = Self::Timestamp> + 'static,
95 Tr::Batch: Batch,
96 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize;
97}
98
99pub trait MzArrangeCore<'scope> {
101 type Timestamp: Timestamp + Lattice;
103 type Input: Container + Clone + 'static;
105
106 fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
113 self,
114 pact: P,
115 name: &str,
116 ) -> Arranged<'scope, TraceAgent<Tr>>
117 where
118 P: ParallelizationContract<Self::Timestamp, Self::Input>,
119 Ba: Batcher<Time = Self::Timestamp> + 'static,
120 Chu: ContainerBuilder<Container = Ba::Output>
121 + for<'a> PushInto<&'a mut Self::Input>
122 + 'static,
123 Bu: Builder<Time = Self::Timestamp, Input = Ba::Output, Output = Tr::Batch>,
124 Tr: Trace + TraceReader<Time = Self::Timestamp> + 'static,
125 Tr::Batch: Batch,
126 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize;
127}
128
129impl<'scope, T, C> MzArrangeCore<'scope> for Stream<'scope, T, C>
130where
131 T: Timestamp + Lattice,
132 C: Container + Clone + 'static,
133{
134 type Timestamp = T;
135 type Input = C;
136
137 fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
138 self,
139 pact: P,
140 name: &str,
141 ) -> Arranged<'scope, TraceAgent<Tr>>
142 where
143 P: ParallelizationContract<T, Self::Input>,
144 Ba: Batcher<Time = T> + 'static,
145 Chu: ContainerBuilder<Container = Ba::Output>
146 + for<'a> PushInto<&'a mut Self::Input>
147 + 'static,
148 Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
149 Tr: Trace + TraceReader<Time = T> + 'static,
150 Tr::Batch: Batch,
151 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
152 {
153 #[allow(clippy::disallowed_methods)]
155 arrange_core::<_, _, Chu, Ba, Bu, _>(self, pact, name).log_arrangement_size()
156 }
157}
158
159impl<'scope, T, K, V, R> MzArrange<'scope> for VecCollection<'scope, T, (K, V), R>
160where
161 T: Timestamp + Lattice,
162 K: ExchangeData + Hashable,
163 V: ExchangeData,
164 R: ExchangeData,
165{
166 fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
167 where
168 Ba: Batcher<Time = T> + 'static,
169 Chu: ContainerBuilder<Container = Ba::Output>
170 + for<'a> PushInto<&'a mut Self::Input>
171 + 'static,
172 Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
173 Tr: Trace + TraceReader<Time = T> + 'static,
174 Tr::Batch: Batch,
175 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
176 {
177 let exchange = Exchange::new(move |update: &((K, V), T, R)| (update.0).0.hashed().into());
178 self.mz_arrange_core::<_, Chu, Ba, Bu, _>(exchange, name)
179 }
180}
181
182impl<'scope, T, C> MzArrangeCore<'scope> for Collection<'scope, T, C>
183where
184 T: Timestamp + Lattice,
185 C: Container + Clone + 'static,
186{
187 type Timestamp = T;
188 type Input = C;
189
190 fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
191 self,
192 pact: P,
193 name: &str,
194 ) -> Arranged<'scope, TraceAgent<Tr>>
195 where
196 P: ParallelizationContract<T, Self::Input>,
197 Ba: Batcher<Time = T> + 'static,
198 Chu: ContainerBuilder<Container = Ba::Output>
199 + for<'a> PushInto<&'a mut Self::Input>
200 + 'static,
201 Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
202 Tr: Trace + TraceReader<Time = T> + 'static,
203 Tr::Batch: Batch,
204 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
205 {
206 self.inner.mz_arrange_core::<_, Chu, Ba, Bu, _>(pact, name)
207 }
208}
209
210pub struct KeyCollection<'scope, T: Timestamp, K: 'static, R: 'static = usize>(
214 VecCollection<'scope, T, K, R>,
215);
216
217impl<'scope, T: Timestamp, K, R: Semigroup> From<VecCollection<'scope, T, K, R>>
218 for KeyCollection<'scope, T, K, R>
219{
220 fn from(value: VecCollection<'scope, T, K, R>) -> Self {
221 KeyCollection(value)
222 }
223}
224
225impl<'scope, T, K, R> MzArrange<'scope> for KeyCollection<'scope, T, K, R>
226where
227 T: Timestamp + Lattice,
228 K: ExchangeData + Hashable,
229 R: ExchangeData,
230{
231 fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
232 where
233 Ba: Batcher<Time = T> + 'static,
234 Chu: ContainerBuilder<Container = Ba::Output>
235 + for<'a> PushInto<&'a mut Self::Input>
236 + 'static,
237 Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
238 Tr: Trace + TraceReader<Time = T> + 'static,
239 Tr::Batch: Batch,
240 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
241 {
242 self.0.map(|d| (d, ())).mz_arrange::<Chu, Ba, Bu, _>(name)
243 }
244}
245
246impl<'scope, T, K, R> MzArrangeCore<'scope> for KeyCollection<'scope, T, K, R>
247where
248 T: Timestamp + Lattice,
249 K: Clone + 'static,
250 R: Clone + 'static,
251{
252 type Timestamp = T;
253 type Input = Vec<((K, ()), T, R)>;
254
255 fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
256 self,
257 pact: P,
258 name: &str,
259 ) -> Arranged<'scope, TraceAgent<Tr>>
260 where
261 P: ParallelizationContract<T, Self::Input>,
262 Ba: Batcher<Time = T> + 'static,
263 Chu: ContainerBuilder<Container = Ba::Output>
264 + for<'a> PushInto<&'a mut Self::Input>
265 + 'static,
266 Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
267 Tr: Trace + TraceReader<Time = T> + 'static,
268 Tr::Batch: Batch,
269 Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
270 {
271 self.0
272 .map(|d| (d, ()))
273 .mz_arrange_core::<_, Chu, Ba, Bu, _>(pact, name)
274 }
275}
276
277pub trait ArrangementSize {
279 fn log_arrangement_size(self) -> Self;
281}
282
283fn log_arrangement_size_inner<'scope, B, L>(
293 arranged: Arranged<'scope, TraceAgent<Spine<ArcBatch<B>>>>,
294 mut logic: L,
295) -> Arranged<'scope, TraceAgent<Spine<ArcBatch<B>>>>
296where
297 B: Batch + 'static,
298 L: FnMut(&B) -> (usize, usize, usize) + 'static,
299{
300 let scope = arranged.stream.scope();
301 let Some(logger) = scope
302 .worker()
303 .logger_for::<ComputeEventBuilder>("materialize/compute")
304 else {
305 return arranged;
306 };
307 let operator_id = arranged.trace.operator().global_id;
308 let trace = Rc::downgrade(&arranged.trace.trace_box_unstable());
309
310 let (mut old_size, mut old_capacity, mut old_allocations) = (0isize, 0isize, 0isize);
311
312 let stream = arranged
313 .stream
314 .unary(Pipeline, "ArrangementSize", |_cap, info| {
315 let address = info.address;
316 logger.log(&ComputeEvent::ArrangementHeapSizeOperator(
317 ArrangementHeapSizeOperator {
318 operator_id,
319 address: address.to_vec(),
320 },
321 ));
322
323 let mut batches: BTreeMap<*const B, (Weak<B>, (usize, usize, usize))> = BTreeMap::new();
329
330 move |input, output| {
331 input.for_each(|time, data| {
332 for batch in data.iter() {
333 batches
334 .entry(Arc::as_ptr(&batch.0))
335 .or_insert_with(|| (Arc::downgrade(&batch.0), logic(&batch.0)));
336 }
337 output.session(&time).give_container(data);
338 });
339 let Some(trace) = trace.upgrade() else {
340 batches.clear();
347 return;
348 };
349
350 trace.borrow().trace().map_batches(|batch| {
351 batches
352 .entry(Arc::as_ptr(&batch.0))
353 .or_insert_with(|| (Arc::downgrade(&batch.0), logic(&batch.0)));
354 });
355
356 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
357 batches.retain(|_, (weak, cached)| {
358 if weak.strong_count() > 0 {
359 let (sz, c, a) = *cached;
360 (size += sz, capacity += c, allocations += a);
361 true
362 } else {
363 false
364 }
365 });
366
367 let size = size.try_into().expect("must fit");
368 if size != old_size {
369 logger.log(&ComputeEvent::ArrangementHeapSize(ArrangementHeapSize {
370 operator_id,
371 delta_size: size - old_size,
372 }));
373 }
374
375 let capacity = capacity.try_into().expect("must fit");
376 if capacity != old_capacity {
377 logger.log(&ComputeEvent::ArrangementHeapCapacity(
378 ArrangementHeapCapacity {
379 operator_id,
380 delta_capacity: capacity - old_capacity,
381 },
382 ));
383 }
384
385 let allocations = allocations.try_into().expect("must fit");
386 if allocations != old_allocations {
387 logger.log(&ComputeEvent::ArrangementHeapAllocations(
388 ArrangementHeapAllocations {
389 operator_id,
390 delta_allocations: allocations - old_allocations,
391 },
392 ));
393 }
394
395 old_size = size;
396 old_capacity = capacity;
397 old_allocations = allocations;
398 }
399 });
400 Arranged {
401 trace: arranged.trace,
402 stream,
403 }
404}
405
406impl<'scope, T, K, V, R> ArrangementSize for Arranged<'scope, KeyValAgent<K, V, T, R>>
407where
408 T: MzTimestamp,
409 K: Data + MzData,
410 V: Data + MzData,
411 R: Semigroup + Ord + MzData + 'static,
412{
413 fn log_arrangement_size(self) -> Self {
414 log_arrangement_size_inner(self, |batch| {
415 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
416 let mut callback = |siz, cap| {
417 size += siz;
418 capacity += cap;
419 allocations += usize::from(cap > 0);
420 };
421 batch.storage.keys.heap_size(&mut callback);
422 batch.storage.vals.offs.heap_size(&mut callback);
423 batch.storage.vals.vals.heap_size(&mut callback);
424 batch.storage.upds.offs.heap_size(&mut callback);
425 batch.storage.upds.times.heap_size(&mut callback);
426 batch.storage.upds.diffs.heap_size(&mut callback);
427 (size, capacity, allocations)
428 })
429 }
430}
431
432impl<'scope, T, K, R> ArrangementSize for Arranged<'scope, KeyAgent<K, T, R>>
433where
434 T: MzTimestamp,
435 K: Data + MzArrangeData,
436 R: Semigroup + Ord + MzData + 'static,
437{
438 fn log_arrangement_size(self) -> Self {
439 log_arrangement_size_inner(self, |batch| {
440 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
441 let mut callback = |siz, cap| {
442 size += siz;
443 capacity += cap;
444 allocations += usize::from(cap > 0);
445 };
446 batch.storage.keys.heap_size(&mut callback);
447 batch.storage.upds.offs.heap_size(&mut callback);
448 batch.storage.upds.times.heap_size(&mut callback);
449 batch.storage.upds.diffs.heap_size(&mut callback);
450 (size, capacity, allocations)
451 })
452 }
453}
454
455impl<'scope, T, V, R> ArrangementSize for Arranged<'scope, RowValAgent<V, T, R>>
456where
457 T: MzTimestamp,
458 V: Data + MzArrangeData,
459 R: Semigroup + Ord + MzArrangeData + 'static,
460{
461 fn log_arrangement_size(self) -> Self {
462 log_arrangement_size_inner(self, |batch| {
463 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
464 let mut callback = |siz, cap| {
465 size += siz;
466 capacity += cap;
467 allocations += usize::from(cap > 0);
468 };
469 batch.storage.keys.heap_size(&mut callback);
470 batch.storage.vals.offs.heap_size(&mut callback);
471 batch.storage.vals.vals.heap_size(&mut callback);
472 batch.storage.upds.offs.heap_size(&mut callback);
473 batch.storage.upds.times.heap_size(&mut callback);
474 batch.storage.upds.diffs.heap_size(&mut callback);
475 (size, capacity, allocations)
476 })
477 }
478}
479
480impl<'scope, T, R> ArrangementSize for Arranged<'scope, RowRowAgent<T, R>>
481where
482 T: MzTimestamp,
483 R: Semigroup + Ord + MzArrangeData + 'static,
484{
485 fn log_arrangement_size(self) -> Self {
486 log_arrangement_size_inner(self, |batch| {
487 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
488 let mut callback = |siz, cap| {
489 size += siz;
490 capacity += cap;
491 allocations += usize::from(cap > 0);
492 };
493 batch.storage.keys.heap_size(&mut callback);
494 batch.storage.vals.offs.heap_size(&mut callback);
495 batch.storage.vals.vals.heap_size(&mut callback);
496 batch.storage.upds.offs.heap_size(&mut callback);
497 batch.storage.upds.times.heap_size(&mut callback);
498 batch.storage.upds.diffs.heap_size(&mut callback);
499 (size, capacity, allocations)
500 })
501 }
502}
503
504impl<'scope, T, DC> ArrangementSize for Arranged<'scope, RowAgent<T, DC::Owned, DC>>
505where
506 T: MzTimestamp,
507 DC: BatchContainer<Owned: Semigroup + 'static> + HeapSize,
508{
509 fn log_arrangement_size(self) -> Self {
510 log_arrangement_size_inner(self, |batch| {
511 let (mut size, mut capacity, mut allocations) = (0, 0, 0);
512 let mut callback = |siz, cap| {
513 size += siz;
514 capacity += cap;
515 allocations += usize::from(cap > 0);
516 };
517 batch.storage.keys.heap_size(&mut callback);
518 batch.storage.upds.offs.heap_size(&mut callback);
519 batch.storage.upds.times.heap_size(&mut callback);
520 batch.storage.upds.diffs.heap_size(&mut callback);
521 (size, capacity, allocations)
522 })
523 }
524}