1use std::collections::HashMap;
45use std::ops::Deref;
46use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as Atomic};
47use std::sync::{Arc, Mutex};
48
49use rudb_common::{Error, LogicalType, Result};
50use rudb_metrics::{LoadProfile, Stage};
51use rudb_storage::Range;
52use rudb_vector::{Bitmap, Chunk, Data, StringColumn, Validity, Vector};
53
54use super::{
55 ColumnStripe, DICTIONARY_CHECK_SEED, DICTIONARY_DECIDE_ROWS, DICTIONARY_DISTINCT_IN_TEN,
56 EncodedBlock, GlobalDictionary, MAX_ENCODE_WORKERS, MAX_PAGE, Part, PendingChunk, STRIPE_PARTS,
57 Spread, Unencoded, Writer, checksum, coded_page, invalid, push_validity, seeded_checksum,
58 stats, unique_codes, weight,
59};
60
61static BUSY: AtomicUsize = AtomicUsize::new(0);
68
69pub const DICTIONARY_CAP_BYTES: u64 = 512 * 1024 * 1024;
77
78#[derive(Debug)]
85pub(crate) struct Coding {
86 flags: Box<[AtomicBool]>,
87 growth: Box<[AtomicU64]>,
89 held: AtomicU64,
91 cap: AtomicU64,
92}
93
94impl Coding {
95 pub(crate) fn new(flags: impl IntoIterator<Item = bool>) -> Self {
96 let flags = flags.into_iter().map(AtomicBool::new).collect::<Box<[_]>>();
97 let growth = flags.iter().map(|_| AtomicU64::new(0)).collect();
98 Self { flags, growth, held: AtomicU64::new(0), cap: AtomicU64::new(DICTIONARY_CAP_BYTES) }
99 }
100
101 pub(crate) fn cap(&self, bytes: u64) {
103 self.cap.store(bytes, Atomic::Relaxed);
104 }
105
106 fn recount(&self, before: u64, now: u64) -> u64 {
109 if now >= before {
110 self.held.fetch_add(now - before, Atomic::Relaxed) + (now - before)
111 } else {
112 self.held.fetch_sub(before - now, Atomic::Relaxed).saturating_sub(before - now)
113 }
114 }
115
116 fn grew_most(&self, index: usize) -> bool {
123 let mine = self.growth[index].load(Atomic::Relaxed);
124 mine > 0 && self.growth.iter().all(|other| other.load(Atomic::Relaxed) <= mine)
125 }
126}
127
128impl Deref for Coding {
129 type Target = [AtomicBool];
130
131 fn deref(&self) -> &[AtomicBool] {
132 &self.flags
133 }
134}
135
136struct Share(usize);
138
139impl Share {
140 fn take(columns: usize, parts: usize) -> Self {
141 let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
142 let cores =
143 std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
144 let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
146 Self(workers)
147 }
148}
149
150impl Drop for Share {
151 fn drop(&mut self) {
152 BUSY.fetch_sub(1, Atomic::Relaxed);
153 }
154}
155
156#[derive(Debug, Clone)]
162pub struct Preparer {
163 types: Vec<LogicalType>,
164 coded: Arc<Coding>,
165 profile: Option<Arc<LoadProfile>>,
166}
167
168#[derive(Debug)]
170pub struct Prepared {
171 parts: Vec<Part>,
172 types: Vec<LogicalType>,
173 columns: Vec<Column>,
174 gathers: Vec<Option<stats::Gather>>,
175 profile: Option<Arc<LoadProfile>>,
176}
177
178#[derive(Debug)]
180pub struct Merged {
181 parts: Vec<Part>,
182 columns: Vec<Merge>,
183 blocks: Vec<Unencoded>,
184 profile: Option<Arc<LoadProfile>>,
185 counted: bool,
188}
189
190#[derive(Debug)]
192pub struct Paged {
193 parts: Vec<Part>,
194 columns: Vec<ColumnStripe>,
195 blocks: Vec<(usize, usize, EncodedBlock)>,
197 counted: bool,
198}
199
200enum Built {
202 Stripe(ColumnStripe),
203 Block(EncodedBlock),
204}
205
206#[derive(Debug)]
208enum Column {
209 Pages(ColumnStripe),
211 Coded(Local),
213}
214
215#[derive(Debug)]
217enum Merge {
218 Pages(ColumnStripe),
219 Codes {
221 parts: Vec<LocalPart>,
222 global: Vec<u32>,
223 },
224 Plain(Local),
228}
229
230const END: u32 = u32::MAX;
232
233#[derive(Debug, Default)]
239struct Local {
240 first: HashMap<u64, u32, Spread>,
242 next: Vec<u32>,
244 hashes: Vec<u64>,
245 checks: Vec<u64>,
246 bytes: Vec<u8>,
248 ends: Vec<usize>,
249 counts: Vec<u64>,
251 nulls: u64,
252 parts: Vec<LocalPart>,
253 blob: bool,
255}
256
257#[derive(Debug)]
259struct LocalPart {
260 codes: Vec<u32>,
261 validity: Vec<u8>,
263 range: Range,
264}
265
266impl Local {
267 fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
273 let mut local = Self::default();
274 let mut mapped = None;
275 for pending in held {
276 let column = pending.chunk.column(index)?;
277 local.blob = column.logical_type() == &LogicalType::Blob;
278 if let Some(codes) = local.code_dictionary(column, &mut mapped)? {
279 let mut validity = Vec::new();
280 push_validity(&mut validity, column);
281 local.parts.push(LocalPart { codes, validity, range: Range::of(column) });
282 continue;
283 }
284 let flat = column.flatten()?;
286 let mut codes = Vec::with_capacity(flat.len());
287 let mut last = None;
288 for row in 0..flat.len() {
289 let text = flat.bytes_at(row).unwrap_or(b"");
292 let code = match last {
295 Some(code) if local.value(code) == text => code,
296 _ => local.code(text)?,
297 };
298 last = Some(code);
299 if flat.is_null_at(row) {
300 local.nulls += 1;
301 } else {
302 local.counts[code as usize] += 1;
303 }
304 codes.push(code);
305 }
306 let mut validity = Vec::new();
307 push_validity(&mut validity, &flat);
308 local.parts.push(LocalPart { codes, validity, range: Range::of(column) });
309 }
310 local.first = HashMap::default();
313 local.next = Vec::new();
314 Ok(local)
315 }
316
317 fn code_dictionary(
332 &mut self,
333 column: &Vector,
334 mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
335 ) -> Result<Option<Vec<u32>>> {
336 let Some((codes, values)) = column.shared_dictionary_parts() else { return Ok(None) };
337 if !matches!(values.validity(), Validity::AllValid) {
338 return Ok(None);
339 }
340 let Some(codes) = codes.get(..column.len()) else { return Ok(None) };
341 let fresh = !matches!(mapped, Some((held, _)) if Arc::ptr_eq(held, values));
342 if fresh {
343 *mapped = Some((Arc::clone(values), vec![END; values.len()]));
344 }
345 let Some((_, map)) = mapped.as_mut() else { return Ok(None) };
346 let every = matches!(column.validity(), Validity::AllValid);
347 let mut coded = Vec::with_capacity(codes.len());
348 for (row, &code) in codes.iter().enumerate() {
349 if !every && !column.validity().is_valid(row) {
350 let code = self.code(b"")?;
352 self.nulls += 1;
353 coded.push(code);
354 continue;
355 }
356 let slot = map
357 .get_mut(code as usize)
358 .ok_or_else(|| invalid("a dictionary code is out of range"))?;
359 if *slot == END {
360 *slot = self.code(values.bytes_at(code as usize).unwrap_or(b""))?;
361 }
362 self.counts[*slot as usize] += 1;
363 coded.push(*slot);
364 }
365 Ok(Some(coded))
366 }
367
368 fn rows(&self) -> Result<Vec<Vector>> {
374 self.parts
375 .iter()
376 .map(|part| {
377 let len = part.codes.len();
378 let mut column = StringColumn::with_capacity(len);
379 for &code in &part.codes {
380 column.push_bytes(self.value(code));
381 }
382 let validity = match part.validity.split_first() {
383 Some((0, _)) => Validity::AllValid,
384 Some((1, _)) => Validity::AllInvalid,
385 Some((2, bits)) => {
386 let mut mask = Bitmap::all_valid(len);
387 for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
388 mask.set(row, false);
389 }
390 Validity::Mask(mask)
391 }
392 _ => return Err(Error::internal("a coded part has no validity")),
393 };
394 let ty = if self.blob { LogicalType::Blob } else { LogicalType::Varchar };
395 Ok(Vector::flat(ty, Data::Varlen(column))?.with_validity(validity))
396 })
397 .collect()
398 }
399
400 fn values(&self) -> usize {
401 self.ends.len()
402 }
403
404 fn value(&self, code: u32) -> &[u8] {
405 let code = code as usize;
406 let from = if code == 0 { 0 } else { self.ends[code - 1] };
407 &self.bytes[from..self.ends[code]]
408 }
409
410 fn code(&mut self, text: &[u8]) -> Result<u32> {
411 let hash = checksum(text);
412 let Some(&first) = self.first.get(&hash) else {
413 let code = self.push(text, hash)?;
414 self.first.insert(hash, code);
415 return Ok(code);
416 };
417 let mut at = first;
418 loop {
419 if self.value(at) == text {
420 return Ok(at);
421 }
422 match self.next[at as usize] {
423 END => break,
424 next => at = next,
425 }
426 }
427 let code = self.push(text, hash)?;
428 self.next[at as usize] = code;
429 Ok(code)
430 }
431
432 fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
433 let code = u32::try_from(self.ends.len())
434 .ok()
435 .filter(|&code| code != END)
436 .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
437 self.bytes.extend_from_slice(text);
438 self.ends.push(self.bytes.len());
439 self.next.push(END);
440 self.hashes.push(hash);
441 self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
442 self.counts.push(0);
443 Ok(code)
444 }
445
446 fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
449 let mut global = Vec::with_capacity(self.values());
450 for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
451 let text = self.value(code as u32);
452 let at = dictionary.code_hashed(text, hash, check)?;
453 let count = dictionary
454 .counts
455 .get_mut(at as usize)
456 .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
457 *count = count.saturating_add(self.counts[code]);
458 global.push(at);
459 }
460 dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
461 Ok(global)
462 }
463}
464
465fn drops_dictionary(rows: usize, distinct: usize) -> bool {
489 rows >= DICTIONARY_DECIDE_ROWS
490 && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
491}
492
493fn demotes(rows: usize, new: usize, total: u64, coding: &Coding, index: usize) -> bool {
506 drops_dictionary(rows, new)
507 || (total > coding.cap.load(Atomic::Relaxed) && coding.grew_most(index))
508}
509
510fn fan_out<T: Send>(
520 jobs: Vec<usize>,
521 workers: usize,
522 profile: Option<&LoadProfile>,
523 work: impl Fn(usize) -> Result<T> + Sync,
524) -> Result<Vec<(usize, T)>> {
525 if workers <= 1 || jobs.len() <= 1 {
526 let _span = profile.map(|profile| profile.span(Stage::Pages));
527 return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
528 }
529 let workers = workers.min(jobs.len());
530 let queue = Mutex::new(jobs);
531 let pieces = std::thread::scope(|scope| {
532 (0..workers)
533 .map(|_| {
534 scope.spawn(|| {
535 let _span = profile.map(|profile| profile.span(Stage::Pages));
536 let mut mine = Vec::new();
537 loop {
538 let taken = queue
539 .lock()
540 .map_err(|_| Error::internal("a native encode worker panicked"))?
541 .pop();
542 let Some(index) = taken else { break };
543 mine.push((index, work(index)?));
544 }
545 Ok(mine)
546 })
547 })
548 .collect::<Vec<_>>()
549 .into_iter()
550 .map(|handle| {
551 handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
552 })
553 .collect::<Result<Vec<Vec<_>>>>()
554 })?;
555 Ok(pieces.into_iter().flatten().collect())
556}
557
558fn column_of(held: &[PendingChunk], index: usize) -> Result<Vec<&Vector>> {
560 held.iter().map(|pending| pending.chunk.column(index)).collect()
561}
562
563fn in_order<T>(width: usize, done: Vec<(usize, T)>) -> Result<Vec<T>> {
565 let mut slots: Vec<Option<T>> = (0..width).map(|_| None).collect();
566 for (index, one) in done {
567 slots[index] = Some(one);
568 }
569 slots
570 .into_iter()
571 .map(|slot| slot.ok_or_else(|| Error::internal("a column was never encoded")))
572 .collect()
573}
574
575impl Preparer {
576 pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
587 if parts.len() > STRIPE_PARTS {
588 return Err(invalid("a stripe was handed more parts than it holds"));
589 }
590 let held = parts
591 .into_iter()
592 .filter(|(_, chunk)| !chunk.is_empty())
593 .map(|(order, chunk)| PendingChunk { order, chunk })
594 .collect::<Vec<_>>();
595 for pending in &held {
596 self.fits(&pending.chunk)?;
597 }
598 self.prepare_held(held)
599 }
600
601 fn fits(&self, chunk: &Chunk) -> Result<()> {
603 if chunk.width() != self.types.len() {
604 return Err(invalid("chunk width differs from table schema"));
605 }
606 for (index, ty) in self.types.iter().enumerate() {
607 if chunk.column(index)?.logical_type() != ty {
608 return Err(invalid("chunk type differs from table schema"));
609 }
610 }
611 Ok(())
612 }
613
614 pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
615 let width = self.types.len();
616 let key = held.first().map_or((0, 0), |pending| pending.order);
617 let share = Share::take(width, held.len());
618 let mut jobs = (0..width).collect::<Vec<_>>();
619 jobs.sort_by_key(|&index| weight(&self.types[index]));
620 let done = fan_out(jobs, share.0, self.profile.as_deref(), |index| {
621 let gather = stats::Gather::new(&self.types[index], 0)
624 .filter(|_| !held.is_empty())
625 .map(|mut gather| {
626 gather.stripe(
627 key,
628 held.iter().filter_map(|pending| pending.chunk.column(index).ok()),
629 );
630 gather
631 });
632 let column = if self.coded[index].load(Atomic::Relaxed) {
633 Column::Coded(Local::code_column(index, &held)?)
634 } else {
635 Column::Pages(Writer::encode_pages(&column_of(&held, index)?)?)
636 };
637 Ok((column, gather))
638 })?;
639 drop(share);
640 let (columns, gathers) = in_order(width, done)?.into_iter().unzip();
641 let parts = held.iter().map(Part::of).collect();
642 drop(held);
643 Ok(Prepared {
644 parts,
645 types: self.types.clone(),
646 columns,
647 gathers,
648 profile: self.profile.clone(),
649 })
650 }
651}
652
653enum Slot<'a> {
655 Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
657 Lent(&'a Mutex<LentColumn>, &'a Lent),
659}
660
661struct Step<'a> {
663 index: usize,
664 column: Column,
665 slot: Slot<'a>,
666 gather: Option<stats::Gather>,
668}
669
670impl Step<'_> {
671 fn cost(&self, coded: &Coding) -> usize {
675 match &self.column {
676 Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
677 local.values().saturating_add(1)
678 }
679 _ => 0,
680 }
681 }
682
683 fn run(
685 self,
686 rows: usize,
687 coded: &Coding,
688 profile: Option<&LoadProfile>,
689 ) -> Result<(usize, Merge, Vec<Unencoded>)> {
690 let Self { index, column, slot, gather } = self;
691 match slot {
692 Slot::Owned(dictionary, mine) => {
693 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
694 }
695 Slot::Lent(held, lent) => {
696 let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
697 if lent.reclaimed.load(Atomic::Acquire) {
700 return Err(Error::internal("a stripe was merged after its table was closed"));
701 }
702 let LentColumn { dictionary, gather: mine } = &mut *held;
703 merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
704 }
705 }
706 }
707}
708
709#[expect(clippy::too_many_arguments, reason = "one column's share of the stripe's merge state")]
711fn merge_column(
712 index: usize,
713 column: Column,
714 stripe: Option<stats::Gather>,
715 dictionary: &mut Option<GlobalDictionary>,
716 gather: &mut Option<stats::Gather>,
717 rows: usize,
718 coded: &Coding,
719 profile: Option<&LoadProfile>,
720) -> Result<(usize, Merge, Vec<Unencoded>)> {
721 if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
722 mine.absorb(stripe);
723 }
724 let mut new = None;
725 let merge = match (column, dictionary.as_mut()) {
726 (Column::Pages(stripe), None) => Merge::Pages(stripe),
727 (Column::Pages(stripe), Some(global)) if global.demoted => Merge::Pages(stripe),
728 (Column::Pages(_), Some(_)) => {
729 return Err(Error::internal(
730 "a column with a global dictionary was prepared without one",
731 ));
732 }
733 (Column::Coded(local), None) => Merge::Plain(local),
734 (Column::Coded(local), Some(global)) if global.demoted => Merge::Plain(local),
736 (Column::Coded(local), Some(global)) => {
737 if global.values() == 0 && drops_dictionary(rows, local.values()) {
740 if let Some(profile) = profile {
741 profile.release(global.charged);
742 }
743 coded.recount(global.charged, 0);
744 *dictionary = None;
745 coded[index].store(false, Atomic::Relaxed);
746 Merge::Plain(local)
747 } else {
748 let before = global.values();
749 let codes = local.merge_into(global)?;
750 new = Some(global.values() - before);
751 Merge::Codes { parts: local.parts, global: codes }
752 }
753 }
754 };
755 if let (Some(new), Some(global)) = (new, dictionary.as_mut()) {
757 let now = global.held_bytes();
758 coded.growth[index].store(now.saturating_sub(global.charged), Atomic::Relaxed);
759 let total =
760 coded.held.load(Atomic::Relaxed).saturating_add(now).saturating_sub(global.charged);
761 if demotes(rows, new, total, coded, index) {
762 global.demote();
763 coded[index].store(false, Atomic::Relaxed);
764 coded.growth[index].store(0, Atomic::Relaxed);
765 }
766 }
767 let blocks = match dictionary {
772 Some(dictionary) => {
773 dictionary.settle()?;
774 let blocks = dictionary.hand_out(index);
775 let (before, now) = dictionary.recharge(profile);
776 coded.recount(before, now);
777 blocks
778 }
779 None => Vec::new(),
780 };
781 Ok((index, merge, blocks))
782}
783
784fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &Coding) -> Result<Merged> {
792 let Prepared { parts, columns, gathers, profile, .. } = prepared;
793 let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
794 let rows: usize = parts.iter().map(|part| part.rows).sum();
795 let width = columns.len();
796 if slots.len() != width || gathers.len() != width {
797 return Err(Error::internal("a stripe was merged into a table of another width"));
798 }
799 let mut steps = columns
800 .into_iter()
801 .zip(gathers)
802 .zip(slots)
803 .enumerate()
804 .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
805 .collect::<Vec<_>>();
806 steps.sort_by_key(|step| step.cost(coded));
809 let workers = std::thread::available_parallelism()
810 .map_or(1, usize::from)
811 .min(MAX_ENCODE_WORKERS)
812 .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
813 .max(1);
814 let done = if workers <= 1 {
815 steps
816 .into_iter()
817 .map(|step| step.run(rows, coded, profile.as_deref()))
818 .collect::<Result<Vec<_>>>()?
819 } else {
820 let queue = Mutex::new(steps);
821 let pieces = std::thread::scope(|scope| {
822 (0..workers)
823 .map(|_| {
824 scope.spawn(|| {
825 let mut mine = Vec::new();
826 loop {
827 let taken = queue
828 .lock()
829 .map_err(|_| Error::internal("a merge worker panicked"))?
830 .pop();
831 let Some(step) = taken else { break };
832 mine.push(step.run(rows, coded, profile.as_deref())?);
833 }
834 Ok(mine)
835 })
836 })
837 .collect::<Vec<_>>()
838 .into_iter()
839 .map(|handle| {
840 handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
841 })
842 .collect::<Result<Vec<Vec<_>>>>()
843 })?;
844 pieces.into_iter().flatten().collect()
845 };
846 let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
847 for (index, merge, blocks) in done {
848 slots[index] = Some((merge, blocks));
849 }
850 let mut merged = Vec::with_capacity(width);
851 let mut blocks = Vec::new();
852 for slot in slots {
853 let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
854 merged.push(merge);
855 blocks.extend(handed);
856 }
857 drop(timing);
858 Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
859}
860
861#[derive(Debug)]
863pub(crate) struct Lent {
864 columns: Box<[Mutex<LentColumn>]>,
865 reclaimed: AtomicBool,
867}
868
869#[derive(Debug)]
871pub(crate) struct LentColumn {
872 pub(crate) dictionary: Option<GlobalDictionary>,
873 gather: Option<stats::Gather>,
874}
875
876impl Lent {
877 pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
878 &self.columns
879 }
880
881 fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
883 for (column, at, block) in blocks {
884 self.columns
885 .get(column)
886 .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
887 .lock()
888 .map_err(|_| Error::internal("a merge panicked"))?
889 .dictionary
890 .as_mut()
891 .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
892 .take_back(at, block)?;
893 }
894 Ok(())
895 }
896
897 #[allow(clippy::type_complexity)]
899 pub(crate) fn reclaim(
900 &self,
901 ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
902 self.reclaimed.store(true, Atomic::Release);
903 let mut dictionaries = Vec::with_capacity(self.columns.len());
904 let mut gathers = Vec::with_capacity(self.columns.len());
905 for column in &self.columns {
906 let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
907 dictionaries.push(held.dictionary.take());
908 gathers.push(held.gather.take());
909 }
910 Ok((dictionaries, gathers))
911 }
912}
913
914#[derive(Debug, Clone)]
921pub struct Merger {
922 lent: Arc<Lent>,
923 types: Vec<LogicalType>,
924 coded: Arc<Coding>,
925}
926
927impl Merger {
928 pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
934 if prepared.types != self.types {
935 return Err(invalid("a stripe was prepared for a table of other columns"));
936 }
937 let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
938 merge_columns(prepared, slots, &self.coded)
939 }
940
941 pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
948 self.lent.take_back(std::mem::take(&mut paged.blocks))
949 }
950}
951
952impl Merged {
953 pub fn pages(self) -> Result<Paged> {
961 let Self { parts, columns, blocks, profile, counted } = self;
962 let width = columns.len();
963 let mut jobs = (width..width + blocks.len())
967 .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
968 .collect::<Vec<_>>();
969 jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
971 let share = Share::take(jobs.len(), parts.len());
972 let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
973 let Some(column) = columns.get(index) else {
974 return Ok(Built::Block(blocks[index - width].encode()?));
975 };
976 Ok(Built::Stripe(match column {
977 Merge::Codes { parts, global } => code_pages(parts, global)?,
978 Merge::Plain(local) => {
979 Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
980 }
981 Merge::Pages(_) => {
982 return Err(Error::internal("a finished column was queued to be built"));
983 }
984 }))
985 })?;
986 drop(share);
987 let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
988 let mut encoded = Vec::with_capacity(blocks.len());
989 for (index, one) in built {
990 match one {
991 Built::Stripe(stripe) => slots[index] = Some(stripe),
992 Built::Block(block) => {
993 let (column, at) = blocks[index - width].place();
994 encoded.push((column, at, block));
995 }
996 }
997 }
998 let columns = columns
999 .into_iter()
1000 .zip(slots)
1001 .map(|(column, slot)| match (column, slot) {
1002 (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
1003 _ => Err(Error::internal("a column was never encoded")),
1004 })
1005 .collect::<Result<Vec<_>>>()?;
1006 Ok(Paged { parts, columns, blocks: encoded, counted })
1007 }
1008}
1009
1010fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
1012 let mut stripe = ColumnStripe {
1013 pages: Vec::with_capacity(parts.len()),
1014 codes: Vec::with_capacity(parts.len()),
1015 sieves: Vec::with_capacity(parts.len()),
1016 ranges: Vec::with_capacity(parts.len()),
1017 };
1018 for part in parts {
1019 let codes = part
1020 .codes
1021 .iter()
1022 .map(|&code| global.get(code as usize).copied())
1023 .collect::<Option<Vec<_>>>()
1024 .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
1025 let bytes = coded_page(&codes, &part.validity)?;
1026 if bytes.len() > MAX_PAGE {
1027 return Err(invalid("column page exceeds the configured bound"));
1028 }
1029 stripe.pages.push(bytes);
1030 stripe.codes.push(Some(unique_codes(&codes)));
1031 stripe.sieves.push(None);
1035 stripe.ranges.push(part.range.clone());
1036 }
1037 Ok(stripe)
1038}
1039
1040impl Writer {
1041 #[must_use]
1046 pub fn preparer(&self) -> Preparer {
1047 Preparer {
1048 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1049 coded: Arc::clone(&self.coded),
1050 profile: self.profile.clone(),
1051 }
1052 }
1053
1054 pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
1065 self.flush_pending()?;
1066 if prepared.columns.len() != self.table.fields.len()
1067 || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
1068 {
1069 return Err(invalid("a stripe was prepared for a table of other columns"));
1070 }
1071 self.table.rows = prepared
1072 .parts
1073 .iter()
1074 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1075 .ok_or_else(|| invalid("row count overflow"))?;
1076 self.merge_held(prepared)
1077 }
1078
1079 pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
1088 let slots = match &self.lent {
1089 Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
1090 None => self
1091 .dictionaries
1092 .iter_mut()
1093 .zip(self.gathers.iter_mut())
1094 .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
1095 .collect::<Vec<_>>(),
1096 };
1097 let mut merged = merge_columns(prepared, slots, &self.coded)?;
1098 merged.counted = true;
1099 Ok(merged)
1100 }
1101
1102 pub fn merger(&mut self) -> Result<Merger> {
1112 self.flush_pending()?;
1113 let lent = match &self.lent {
1114 Some(lent) => Arc::clone(lent),
1115 None => {
1116 let lent = Arc::new(Lent {
1117 columns: std::mem::take(&mut self.dictionaries)
1118 .into_iter()
1119 .zip(std::mem::take(&mut self.gathers))
1120 .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
1121 .collect(),
1122 reclaimed: AtomicBool::new(false),
1123 });
1124 self.lent = Some(Arc::clone(&lent));
1125 lent
1126 }
1127 };
1128 Ok(Merger {
1129 lent,
1130 types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1131 coded: Arc::clone(&self.coded),
1132 })
1133 }
1134
1135 pub fn write(&mut self, paged: Paged) -> Result<()> {
1141 self.write_paged(paged)
1142 }
1143
1144 pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
1145 let Paged { parts, columns, blocks, counted } = paged;
1146 if !counted {
1147 self.table.rows = parts
1148 .iter()
1149 .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1150 .ok_or_else(|| invalid("row count overflow"))?;
1151 }
1152 if let Some(lent) = &self.lent {
1153 lent.take_back(blocks)?;
1154 } else {
1155 for (column, at, block) in blocks {
1156 self.dictionaries
1157 .get_mut(column)
1158 .and_then(Option::as_mut)
1159 .ok_or_else(|| {
1160 Error::internal("a dictionary block came back to no dictionary")
1161 })?
1162 .take_back(at, block)?;
1163 }
1164 }
1165 if parts.is_empty() {
1166 return self.place_blocks();
1167 }
1168 self.write_stripe(&parts, columns)
1169 }
1170
1171 pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
1177 let merged = self.merge(prepared)?;
1178 let paged = merged.pages()?;
1179 self.write(paged)
1180 }
1181}
1182
1183#[cfg(test)]
1184mod tests {
1185 use std::fs;
1186 use std::path::PathBuf;
1187 use std::time::{SystemTime, UNIX_EPOCH};
1188
1189 use rudb_common::{Field, Value};
1190 use rudb_vector::Vector;
1191
1192 use super::*;
1193 use crate::Reader;
1194
1195 const PART: usize = 1_000;
1196
1197 #[test]
1203 fn dictionary_parts_code_as_their_rows_would() {
1204 let texts = |values: &[&str]| {
1205 Arc::new(
1206 Vector::from_values(
1207 LogicalType::Varchar,
1208 &values
1209 .iter()
1210 .map(|text| Value::Varchar((*text).to_string()))
1211 .collect::<Vec<_>>(),
1212 )
1213 .expect("a dictionary"),
1214 )
1215 };
1216 let shared = texts(&["b", "a", "", "c", "unused"]);
1217 let other = texts(&["c", "d", "a"]);
1218 let mut nulls = Bitmap::all_valid(6);
1219 nulls.set(1, false);
1220 nulls.set(4, false);
1221 let parts = [
1222 Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
1223 Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
1224 .expect("codes")
1225 .with_validity(Validity::Mask(nulls)),
1226 Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
1227 ];
1228 let held = |flat: bool| {
1229 parts
1230 .iter()
1231 .enumerate()
1232 .map(|(at, part)| PendingChunk {
1233 order: (at as u64, 0),
1234 chunk: Chunk::new(vec![if flat {
1235 part.flatten().expect("flat")
1236 } else {
1237 part.clone()
1238 }])
1239 .expect("a chunk"),
1240 })
1241 .collect::<Vec<_>>()
1242 };
1243 let parquet = held(false);
1244 let before = rudb_common::slow::here();
1245 let coded = Local::code_column(0, &parquet).expect("coded");
1246 assert_eq!(
1247 rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
1248 0,
1249 "a part that came in as codes was flattened",
1250 );
1251 let flat = Local::code_column(0, &held(true)).expect("coded");
1252 assert_eq!(coded.values(), flat.values());
1253 for code in 0..flat.values() as u32 {
1254 assert_eq!(coded.value(code), flat.value(code), "value {code}");
1255 }
1256 assert_eq!(coded.counts, flat.counts);
1257 assert_eq!(coded.nulls, flat.nulls);
1258 assert_eq!(coded.nulls, 2);
1259 assert_eq!(coded.hashes, flat.hashes);
1260 assert_eq!(coded.checks, flat.checks);
1261 assert_eq!(coded.parts.len(), flat.parts.len());
1262 for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
1263 assert_eq!(coded.codes, flat.codes);
1264 assert_eq!(coded.validity, flat.validity);
1265 }
1266 }
1267
1268 fn path(label: &str) -> PathBuf {
1269 let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1270 std::env::temp_dir()
1271 .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1272 }
1273
1274 fn fields() -> Vec<Field> {
1275 vec![
1276 Field::required("id", LogicalType::BigInt),
1277 Field::new("city", LogicalType::Varchar),
1278 Field::new("note", LogicalType::Varchar),
1279 ]
1280 }
1281
1282 fn row(id: usize) -> [Value; 3] {
1287 let city = if id % 11 == 0 {
1288 Value::Null
1289 } else {
1290 Value::Varchar(format!("city {}", (id / 7) % 13))
1291 };
1292 let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1293 [Value::BigInt(id as i64), city, note]
1294 }
1295
1296 fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1298 (first..first + parts)
1299 .map(|part| {
1300 let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1301 let column = |at: usize| {
1302 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1303 Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1304 };
1305 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1306 ((part as u64, 0), chunk)
1307 })
1308 .collect()
1309 }
1310
1311 fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1313 vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1314 }
1315
1316 fn check(path: &PathBuf) {
1317 let reader = Reader::open(path).expect("reopen");
1318 assert_eq!(reader.parts(), 13);
1319 for part in 0..13 {
1320 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1321 for at in [0, 17, PART - 1] {
1322 let want = row(part * PART + at);
1323 for (column, value) in want.iter().enumerate() {
1324 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1325 }
1326 }
1327 }
1328 }
1329
1330 #[test]
1338 fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1339 let alone = path("alone");
1340 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1341 for run in runs() {
1342 writer.append_stripe(run).expect("a stripe");
1343 }
1344 writer.finish().expect("commit");
1345
1346 let split = path("split");
1347 let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1348 let preparer = writer.preparer();
1349 let prepared = runs()
1350 .into_iter()
1351 .map(|run| preparer.prepare(run).expect("prepared"))
1352 .collect::<Vec<_>>();
1353 for one in prepared {
1354 writer.append_prepared(one).expect("a stripe");
1355 }
1356 assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1357 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1358 writer.finish().expect("commit");
1359
1360 assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1361 check(&split);
1362 fs::remove_file(alone).expect("remove");
1363 fs::remove_file(split).expect("remove");
1364 }
1365
1366 #[test]
1370 fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1371 let path = path("crossed");
1372 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1373 let preparer = writer.preparer();
1374 let mut merged = runs()
1375 .into_iter()
1376 .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1377 .map(|merged| merged.pages().expect("paged"))
1378 .collect::<Vec<_>>();
1379 merged.reverse();
1380 for paged in merged {
1381 writer.write(paged).expect("written");
1382 }
1383 writer.finish().expect("commit");
1384 check(&path);
1385 fs::remove_file(path).expect("remove");
1386 }
1387
1388 #[test]
1391 fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1392 let alone = path("alone-merger");
1393 let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1394 for run in runs() {
1395 writer.append_stripe(run).expect("a stripe");
1396 }
1397 writer.finish().expect("commit");
1398
1399 let lent = path("lent");
1400 let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1401 let preparer = writer.preparer();
1402 let merger = writer.merger().expect("a merger");
1403 for run in runs() {
1404 let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1405 let mut paged = merged.pages().expect("paged");
1406 merger.give_back(&mut paged).expect("given back");
1407 writer.write(paged).expect("written");
1408 }
1409 assert_eq!(writer.table.rows, 13 * PART);
1410 writer.finish().expect("commit");
1411
1412 assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1413 check(&lent);
1414 fs::remove_file(alone).expect("remove");
1415 fs::remove_file(lent).expect("remove");
1416 }
1417
1418 #[test]
1421 fn stripes_merged_on_several_threads_at_once_read_back() {
1422 let path = path("merged-at-once");
1423 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1424 let preparer = writer.preparer();
1425 let merger = writer.merger().expect("a merger");
1426 let writer = Mutex::new(writer);
1427 std::thread::scope(|scope| {
1428 for run in runs() {
1429 let (preparer, merger, writer) = (&preparer, &merger, &writer);
1430 scope.spawn(move || {
1431 let merged =
1432 merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1433 let mut paged = merged.pages().expect("paged");
1434 merger.give_back(&mut paged).expect("given back");
1435 writer.lock().expect("the writer").write(paged).expect("written");
1436 });
1437 }
1438 });
1439 writer.into_inner().expect("the writer").finish().expect("commit");
1440 check(&path);
1441 fs::remove_file(path).expect("remove");
1442 }
1443
1444 #[test]
1447 fn a_merge_after_the_table_is_closed_is_refused() {
1448 let path = path("late");
1449 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1450 let preparer = writer.preparer();
1451 let merger = writer.merger().expect("a merger");
1452 writer.finish().expect("commit");
1453 let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1454 assert!(merger.merge(prepared).is_err());
1455 fs::remove_file(path).expect("remove");
1456 }
1457
1458 #[test]
1461 fn a_stripe_of_another_table_is_refused_at_the_merge() {
1462 let path = path("refused");
1463 let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1464 let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1465 .expect("a file");
1466 let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1467 assert!(writer.merge(prepared).is_err());
1468 assert_eq!(writer.table.rows, 0);
1469 drop(other);
1470 fs::remove_file(path.with_extension("other")).expect("remove");
1471 fs::remove_file(path).expect("remove");
1472 }
1473
1474 fn turning(id: usize) -> [Value; 3] {
1478 let url = match id {
1479 _ if id % 13 == 0 => Value::Null,
1480 _ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
1481 _ => Value::Varchar(format!("https://example.com/page/{id}")),
1482 };
1483 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1484 }
1485
1486 fn turning_fields() -> Vec<Field> {
1487 vec![
1488 Field::required("id", LogicalType::BigInt),
1489 Field::new("city", LogicalType::Varchar),
1490 Field::new("url", LogicalType::Varchar),
1491 ]
1492 }
1493
1494 fn turning_stripe(
1495 rows: fn(usize) -> [Value; 3],
1496 first: usize,
1497 parts: usize,
1498 ) -> Vec<((u64, u64), Chunk)> {
1499 (first..first + parts)
1500 .map(|part| {
1501 let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
1502 let column = |at: usize| {
1503 let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1504 Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
1505 };
1506 let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1507 ((part as u64, 0), chunk)
1508 })
1509 .collect()
1510 }
1511
1512 fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1513 vec![
1514 turning_stripe(turning, 0, 5),
1515 turning_stripe(turning, 5, 5),
1516 turning_stripe(turning, 10, 3),
1517 ]
1518 }
1519
1520 fn check_turned(path: &PathBuf) {
1522 let reader = Reader::open(path).expect("reopen");
1523 assert_eq!(reader.parts(), 13);
1524 assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
1525 for part in 0..13 {
1526 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1527 let url = chunk.column(2).expect("url");
1528 assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
1529 for at in 0..PART {
1530 let want = turning(part * PART + at);
1531 for (column, value) in want.iter().enumerate() {
1532 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1533 }
1534 }
1535 }
1536 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1538 assert_eq!(reader.text_extremes(2).expect("asked"), None);
1539 assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
1540 assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
1541 assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
1542 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1544 assert!(reader.text_extremes(1).expect("asked").is_some());
1545 }
1546
1547 #[test]
1551 fn a_column_that_turns_unique_is_demoted_and_reads_back() {
1552 let alone = path("demoted-alone");
1553 let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
1554 let preparer = writer.preparer();
1555 let mut runs = turning_runs().into_iter();
1556 writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
1557 assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
1558 for run in runs {
1559 writer.append_stripe(run).expect("a stripe");
1560 }
1561 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
1562 assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1563 writer.finish().expect("commit");
1564 check_turned(&alone);
1565
1566 let split = path("demoted-split");
1567 let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
1568 let preparer = writer.preparer();
1569 let prepared = turning_runs()
1570 .into_iter()
1571 .map(|run| preparer.prepare(run).expect("prepared"))
1572 .collect::<Vec<_>>();
1573 for one in prepared {
1574 writer.append_prepared(one).expect("a stripe");
1575 }
1576 writer.finish().expect("commit");
1577 check_turned(&split);
1578
1579 fs::remove_file(alone).expect("remove");
1580 fs::remove_file(split).expect("remove");
1581 }
1582
1583 fn growing(id: usize) -> [Value; 3] {
1586 let url = Value::Varchar(format!("https://example.com/{}", id / 4));
1587 [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1588 }
1589
1590 #[test]
1593 fn the_dictionary_cap_demotes_the_column_that_grew_most() {
1594 let path = path("capped");
1595 let mut writer = Writer::create(&path, "t", turning_fields())
1596 .expect("a file")
1597 .with_dictionary_cap(1 << 30);
1598 let preparer = writer.preparer();
1599 writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
1600 assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
1601 assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
1602
1603 writer.coded.cap(1);
1604 writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
1605 assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
1606 assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
1607 writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
1608 assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
1609 writer.finish().expect("commit");
1610
1611 let reader = Reader::open(&path).expect("reopen");
1612 assert_eq!(reader.table().demoted, [false, false, true]);
1613 for part in 0..13 {
1614 let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1615 for at in 0..PART {
1616 let want = growing(part * PART + at);
1617 for (column, value) in want.iter().enumerate() {
1618 assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1619 }
1620 }
1621 }
1622 assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1623 assert_eq!(reader.distinct_values(2).expect("asked"), None);
1624 fs::remove_file(path).expect("remove");
1625 }
1626}