1#![deny(missing_docs)]
19
20pub mod batcher;
21pub mod body;
22pub mod builder;
23pub mod builder_input;
24pub mod chunk;
25pub mod consolidate;
26pub mod merge_batcher;
27pub mod unload;
28
29use std::hash::{BuildHasher, Hash, Hasher};
30use std::sync::LazyLock;
31
32use columnar::Borrow;
33use columnar::bytes::indexed;
34use columnar::common::IterOwn;
35use columnar::{Clear, FromBytes, Index, Len};
36use columnar::{Columnar, Ref};
37use differential_dataflow::Hashable;
38use differential_dataflow::collection::containers::{Enter, ResultsIn};
39use differential_dataflow::trace::implementations::merge_batcher::MergeBatcher;
40use itertools::Itertools;
41use timely::Accountable;
42use timely::bytes::arc::Bytes;
43use timely::container::{DrainContainer, PushInto, SizableContainer};
44use timely::dataflow::channels::ContainerBytes;
45use timely::progress::Timestamp;
46use timely::progress::timestamp::Refines;
47
48use crate::columnation::ColInternalMerger;
49
50pub type Col2ValBatcher<K, V, T, R> = MergeBatcher<ColInternalMerger<(K, V), T, R>>;
57pub type Col2KeyBatcher<K, T, R> = Col2ValBatcher<K, (), T, R>;
59
60pub type Col2ValPagedBatcher<K, V, T, R> = merge_batcher::ColumnMergeBatcher<(K, V), T, R>;
72
73pub type Col2ValColBatcher<K, V, T, R> = MergeBatcher<batcher::ColumnMerger<(K, V), T, R>>;
81
82pub enum Column<C: Columnar> {
92 Typed(C::Container),
94 Bytes(Bytes),
96 Align(Vec<u64>),
103}
104
105impl<C: Columnar> Column<C> {
106 #[inline]
113 pub fn clear(&mut self) {
114 match self {
115 Column::Typed(t) => t.clear(),
116 Column::Bytes(_) | Column::Align(_) => *self = Default::default(),
117 }
118 }
119
120 #[inline]
126 pub fn is_empty(&self) -> bool {
127 match self {
128 Column::Typed(t) => t.is_empty(),
129 Column::Bytes(_) | Column::Align(_) => self.borrow().is_empty(),
130 }
131 }
132
133 #[inline(always)]
137 pub fn borrow(&self) -> <C::Container as Borrow>::Borrowed<'_> {
138 match self {
139 Column::Typed(t) => t.borrow(),
140 Column::Bytes(b) => <<C::Container as Borrow>::Borrowed<'_>>::from_bytes(
141 &mut indexed::decode(bytemuck::cast_slice(b)),
142 ),
143 Column::Align(a) => {
144 <<C::Container as Borrow>::Borrowed<'_>>::from_bytes(&mut indexed::decode(a))
145 }
146 }
147 }
148}
149
150impl<C: Columnar> Default for Column<C> {
151 fn default() -> Self {
152 Self::Typed(Default::default())
153 }
154}
155
156impl<C: Columnar> Clone for Column<C>
157where
158 C::Container: Clone,
159{
160 fn clone(&self) -> Self {
161 match self {
162 Column::Typed(t) => Column::Typed(t.clone()),
165 Column::Bytes(b) => {
166 assert_eq!(b.len() % 8, 0);
167 Self::Align(bytemuck::allocation::pod_collect_to_vec(b))
168 }
169 Column::Align(a) => Column::Align(a.clone()),
170 }
171 }
172}
173
174impl<C: Columnar> Accountable for Column<C> {
175 #[inline]
176 fn record_count(&self) -> i64 {
177 self.borrow().len().try_into().expect("Must fit")
178 }
179}
180impl<C: Columnar> DrainContainer for Column<C> {
181 type Item<'a> = Ref<'a, C>;
182 type DrainIter<'a> = IterOwn<<C::Container as Borrow>::Borrowed<'a>>;
183 #[inline]
184 fn drain(&mut self) -> Self::DrainIter<'_> {
185 self.borrow().into_index_iter()
186 }
187}
188
189impl<C: Columnar, T> PushInto<T> for Column<C>
190where
191 C::Container: columnar::Push<T>,
192{
193 #[inline]
194 fn push_into(&mut self, item: T) {
195 use columnar::Push;
196 match self {
197 Column::Typed(t) => t.push(item),
198 Column::Align(_) | Column::Bytes(_) => {
199 unimplemented!("Pushing into Column::Bytes without first clearing");
202 }
203 }
204 }
205}
206
207impl<D, T1, T2, R> Enter<T1, T2> for Column<(D, T1, R)>
210where
211 D: Columnar,
212 T1: Columnar + Timestamp,
213 T2: Columnar + Refines<T1>,
214 R: Columnar,
215 (D, T1, R): Columnar<Container = (D::Container, T1::Container, R::Container)>,
216 (D, T2, R): Columnar<Container = (D::Container, T2::Container, R::Container)>,
217 for<'a> D::Container: columnar::Push<Ref<'a, D>>,
218 for<'a> T2::Container: columnar::Push<&'a T2>,
219 for<'a> R::Container: columnar::Push<Ref<'a, R>>,
220{
221 type InnerContainer = Column<(D, T2, R)>;
222
223 fn enter(self) -> Self::InnerContainer {
224 use columnar::Push;
225 match self {
226 Column::Typed((data, times, diffs)) => {
228 let mut inner = T2::Container::default();
229 for time in times.borrow().into_index_iter() {
230 inner.push(&T2::to_inner(T1::into_owned(time)));
231 }
232 Column::Typed((data, inner, diffs))
233 }
234 serialized => {
239 let (borrowed_data, borrowed_times, borrowed_diffs) = serialized.borrow();
240 let mut times = T2::Container::default();
241 for time in borrowed_times.into_index_iter() {
242 times.push(&T2::to_inner(T1::into_owned(time)));
243 }
244 let view = (borrowed_data, times.borrow(), borrowed_diffs);
245 let words = indexed::length_in_words(&view);
246 let mut alloc: Vec<u64> = Vec::with_capacity(words);
247 indexed::encode(&mut alloc, &view);
248 Column::Align(alloc)
249 }
250 }
251 }
252}
253
254impl<D, T, R> ResultsIn<T::Summary> for Column<(D, T, R)>
261where
262 D: Columnar,
263 T: Columnar + Timestamp,
264 R: Columnar,
265 (D, T, R): Columnar<Container = (D::Container, T::Container, R::Container)>,
266 for<'a> D::Container: columnar::Push<Ref<'a, D>>,
267 for<'a> T::Container: columnar::Push<&'a T>,
268 for<'a> R::Container: columnar::Push<Ref<'a, R>>,
269{
270 fn results_in(self, step: &T::Summary) -> Self {
271 use columnar::Push;
272 use timely::progress::PathSummary;
273
274 let (times, kept) = {
278 let (_, borrowed_times, _) = self.borrow();
279 let mut times = T::Container::default();
280 let mut time = T::minimum();
281 let kept: Vec<bool> = borrowed_times
282 .into_index_iter()
283 .map(|reference| {
284 time.copy_from(reference);
285 match step.results_in(&time) {
286 Some(time) => {
287 times.push(&time);
288 true
289 }
290 None => false,
291 }
292 })
293 .collect();
294 (times, kept)
295 };
296
297 if kept.iter().all(|kept| *kept) {
298 match self {
299 Column::Typed((data, _, diffs)) => Column::Typed((data, times, diffs)),
301 serialized => {
304 let (borrowed_data, _, borrowed_diffs) = serialized.borrow();
305 let view = (borrowed_data, times.borrow(), borrowed_diffs);
306 let words = indexed::length_in_words(&view);
307 let mut alloc: Vec<u64> = Vec::with_capacity(words);
308 indexed::encode(&mut alloc, &view);
309 Column::Align(alloc)
310 }
311 }
312 } else {
313 let (borrowed_data, _, borrowed_diffs) = self.borrow();
314 let mut data = D::Container::default();
315 let mut diffs = R::Container::default();
316 let records = borrowed_data
317 .into_index_iter()
318 .zip_eq(borrowed_diffs.into_index_iter())
319 .zip_eq(kept.iter());
320 for ((datum, diff), _) in records.filter(|(_, kept)| **kept) {
321 data.push(datum);
322 diffs.push(diff);
323 }
324 Column::Typed((data, times, diffs))
325 }
326 }
327}
328
329const SHIP_WORDS: usize = 1 << 18;
334
335#[inline]
346pub(crate) fn at_serialized_capacity<'a, A>(borrow: &A) -> bool
347where
348 A: columnar::AsBytes<'a>,
349{
350 indexed::length_in_words(borrow) >= SHIP_WORDS - SHIP_WORDS / 10
351}
352
353impl<C: Columnar> SizableContainer for Column<C> {
354 fn at_capacity(&self) -> bool {
355 match self {
363 Column::Typed(c) => at_serialized_capacity(&c.borrow()),
364 Column::Bytes(_) | Column::Align(_) => true,
365 }
366 }
367
368 fn ensure_capacity(&mut self, _stash: &mut Option<Self>) {
369 }
376}
377
378impl<C: Columnar> ContainerBytes for Column<C> {
379 #[inline]
380 fn from_bytes(bytes: Bytes) -> Self {
381 assert_eq!(bytes.len() % 8, 0);
387 if let Ok(_) = bytemuck::try_cast_slice::<_, u64>(&bytes) {
388 Self::Bytes(bytes)
389 } else {
390 Self::Align(bytemuck::allocation::pod_collect_to_vec(&bytes[..]))
393 }
394 }
395
396 #[inline]
397 fn length_in_bytes(&self) -> usize {
398 match self {
399 Column::Typed(t) => indexed::length_in_bytes(&t.borrow()),
400 Column::Bytes(b) => b.len(),
401 Column::Align(a) => 8 * a.len(),
402 }
403 }
404
405 #[inline]
406 fn into_bytes<W: ::std::io::Write>(&self, writer: &mut W) {
407 match self {
408 Column::Typed(t) => indexed::write(writer, &t.borrow()).unwrap(),
409 Column::Bytes(b) => writer.write_all(b).unwrap(),
410 Column::Align(a) => writer.write_all(bytemuck::cast_slice(a)).unwrap(),
411 }
412 }
413}
414
415#[inline(always)]
419pub fn columnar_exchange<K, V, T, D>(((k, _), _, _): &Ref<'_, ((K, V), T, D)>) -> u64
420where
421 K: Columnar,
422 for<'a> Ref<'a, K>: Hash,
423 V: Columnar,
424 D: Columnar,
425 T: Columnar,
426{
427 k.hashed()
428}
429
430pub fn columnar_exchange_data<D, T, R>((d, _, _): &Ref<'_, (D, T, R)>) -> u64
436where
437 D: Columnar,
438 for<'a> Ref<'a, D>: Hash,
439 T: Columnar,
440 R: Columnar,
441{
442 d.hashed()
443}
444
445pub fn columnar_consolidate_exchange<D, T, R>((d, _, _): &Ref<'_, (D, T, R)>) -> u64
454where
455 D: Columnar,
456 for<'a> Ref<'a, D>: Hash,
457 T: Columnar,
458 R: Columnar,
459{
460 static STATE: LazyLock<ahash::RandomState> = LazyLock::new(crate::hash::fixed_state);
461 let mut hasher = STATE.build_hasher();
462 d.hash(&mut hasher);
463 hasher.finish()
464}
465
466#[cfg(test)]
467mod tests {
468 use timely::bytes::arc::BytesMut;
469 use timely::container::PushInto;
470 use timely::dataflow::channels::ContainerBytes;
471
472 use super::*;
473
474 fn raw_columnar_bytes() -> Vec<u8> {
476 let mut raw = Vec::new();
477 raw.extend(16_u64.to_le_bytes()); raw.extend(28_u64.to_le_bytes()); raw.extend(1_i32.to_le_bytes());
480 raw.extend(2_i32.to_le_bytes());
481 raw.extend(3_i32.to_le_bytes());
482 raw.extend([0, 0, 0, 0]); raw
484 }
485
486 #[mz_ore::test]
487 fn test_column_clone() {
488 let columns = Columnar::as_columns([1, 2, 3].iter());
489 let column_typed: Column<i32> = Column::Typed(columns);
490 let column_typed2 = column_typed.clone();
491
492 assert_eq!(
493 column_typed2.borrow().into_index_iter().collect::<Vec<_>>(),
494 vec![&1, &2, &3]
495 );
496
497 let bytes = BytesMut::from(raw_columnar_bytes()).freeze();
498 let column_bytes: Column<i32> = Column::Bytes(bytes);
499 let column_bytes2 = column_bytes.clone();
500
501 assert_eq!(
502 column_bytes2.borrow().into_index_iter().collect::<Vec<_>>(),
503 vec![&1, &2, &3]
504 );
505
506 let raw = raw_columnar_bytes();
507 let mut region: Vec<u64> = vec![0; raw.len() / 8];
508 let region_bytes = bytemuck::cast_slice_mut(&mut region[..]);
509 region_bytes[..raw.len()].copy_from_slice(&raw);
510 let column_align: Column<i32> = Column::Align(region);
511 let column_align2 = column_align.clone();
512
513 assert_eq!(
514 column_align2.borrow().into_index_iter().collect::<Vec<_>>(),
515 vec![&1, &2, &3]
516 );
517 }
518
519 #[mz_ore::test]
522 fn test_column_known_bytes() {
523 let mut column: Column<i32> = Default::default();
524 column.push_into(1);
525 column.push_into(2);
526 column.push_into(3);
527 let mut data = Vec::new();
528 column.into_bytes(&mut std::io::Cursor::new(&mut data));
529 assert_eq!(data, raw_columnar_bytes());
530 }
531
532 #[mz_ore::test]
533 fn test_column_from_bytes() {
534 let raw = raw_columnar_bytes();
535
536 let buf = vec![0; raw.len() + 8];
537 let align = buf.as_ptr().align_offset(std::mem::size_of::<u64>());
538 let mut bytes_mut = BytesMut::from(buf);
539 let _ = bytes_mut.extract_to(align);
540 bytes_mut[..raw.len()].copy_from_slice(&raw);
541 let aligned_bytes = bytes_mut.extract_to(raw.len());
542
543 let column: Column<i32> = Column::from_bytes(aligned_bytes);
544 assert!(matches!(column, Column::Bytes(_)));
545 assert_eq!(
546 column.borrow().into_index_iter().collect::<Vec<_>>(),
547 vec![&1, &2, &3]
548 );
549
550 let buf = vec![0; raw.len() + 8];
551 let align = buf.as_ptr().align_offset(std::mem::size_of::<u64>());
552 let mut bytes_mut = BytesMut::from(buf);
553 let _ = bytes_mut.extract_to(align + 1);
554 bytes_mut[..raw.len()].copy_from_slice(&raw);
555 let unaligned_bytes = bytes_mut.extract_to(raw.len());
556
557 let column: Column<i32> = Column::from_bytes(unaligned_bytes);
558 assert!(matches!(column, Column::Align(_)));
559 assert_eq!(
560 column.borrow().into_index_iter().collect::<Vec<_>>(),
561 vec![&1, &2, &3]
562 );
563 }
564
565 #[mz_ore::test]
568 #[cfg_attr(miri, ignore)] fn ship_threshold_monotone() {
570 use columnar::Push;
571 let mut container = <Vec<u64> as Columnar>::Container::default();
572 let wide: Vec<u64> = vec![0u64; 50_000];
574 let mut fired = false;
575 for pushes in 1..=64 {
576 container.push(&wide);
577 let now = at_serialized_capacity(&container.borrow());
578 if fired {
579 assert!(now, "ship signal un-fired at {pushes} records");
580 }
581 fired = fired || now;
582 }
583 assert!(fired, "ship signal never fired");
584 }
585}