From e1f17a49f35811095de11b01b6b7ea62cf33d3e7 Mon Sep 17 00:00:00 2001 From: Ben Pfaff Date: Fri, 14 Jun 2024 14:58:49 -0700 Subject: [PATCH] dbsp: Change threshold for spilling to storage from rows to bytes. A row can be any old size but a byte is well-defined. Signed-off-by: Ben Pfaff --- benchmark/feldera-sql/run.py | 10 +- crates/adapters/src/controller/mod.rs | 8 +- crates/dbsp/src/circuit/dbsp_handle.rs | 12 +- crates/dbsp/src/circuit/runtime.rs | 16 +- crates/dbsp/src/storage/file/reader.rs | 9 + crates/dbsp/src/trace/mod.rs | 13 ++ .../src/trace/ord/fallback/indexed_wset.rs | 204 ++++++++++-------- .../dbsp/src/trace/ord/fallback/key_batch.rs | 190 ++++++++++------ crates/dbsp/src/trace/ord/fallback/utils.rs | 104 ++++++++- .../dbsp/src/trace/ord/fallback/val_batch.rs | 186 ++++++++++------ crates/dbsp/src/trace/ord/fallback/wset.rs | 200 ++++++++++------- .../src/trace/ord/file/indexed_wset_batch.rs | 46 ++-- crates/dbsp/src/trace/ord/file/key_batch.rs | 46 ++-- crates/dbsp/src/trace/ord/file/val_batch.rs | 49 ++--- crates/dbsp/src/trace/ord/file/wset_batch.rs | 43 ++-- .../src/trace/ord/vec/indexed_wset_batch.rs | 5 + crates/dbsp/src/trace/ord/vec/key_batch.rs | 4 + crates/dbsp/src/trace/ord/vec/val_batch.rs | 4 + crates/dbsp/src/trace/ord/vec/wset_batch.rs | 5 + crates/dbsp/src/trace/spine_async/mod.rs | 4 + crates/dbsp/src/trace/spine_fueled.rs | 4 + crates/dbsp/src/trace/test/test_batch.rs | 4 + crates/nexmark/benches/nexmark/main.rs | 6 +- crates/nexmark/src/config.rs | 10 +- crates/pipeline-types/src/config.rs | 9 +- crates/pipeline_manager/src/db/test.rs | 6 +- openapi.json | 10 +- .../sql/streaming/StreamingTests.java | 2 +- .../services/manager/models/PipelineConfig.ts | 9 +- .../services/manager/models/RuntimeConfig.ts | 9 +- 30 files changed, 736 insertions(+), 491 deletions(-) diff --git a/benchmark/feldera-sql/run.py b/benchmark/feldera-sql/run.py index fba58cc2780..cb4fe4f38bc 100755 --- a/benchmark/feldera-sql/run.py +++ b/benchmark/feldera-sql/run.py @@ -291,7 +291,7 @@ def main(): parser.add_argument('--output', action=argparse.BooleanOptionalAction, help='whether to write query output back to Kafka (default: --no-output)') parser.add_argument('--merge', action=argparse.BooleanOptionalAction, help='whether to merge all the queries into one program (default: --no-merge)') parser.add_argument('--storage', action=argparse.BooleanOptionalAction, help='whether to enable storage (default: --no-storage)') - parser.add_argument('--min-storage-rows', type=int, help='If storage is enabled, the minimum number of rows to write a batch to storage.') + parser.add_argument('--min-storage-bytes', type=int, help='If storage is enabled, the minimum number of bytes to write a batch to storage.') parser.add_argument('--query', action='append', help='queries to run (by default, all queries), specify one or more of: ' + ','.join(sort_queries(QUERY_SQL.keys()))) parser.add_argument('--input-topic-suffix', help='suffix to apply to input topic names (by default, "")') parser.add_argument('--csv', help='File to write results in .csv format') @@ -313,9 +313,9 @@ def main(): queries = sort_queries(parse_queries(parser.parse_args().query)) cores = int(parser.parse_args().cores) storage = parser.parse_args().storage - min_storage_rows = parser.parse_args().min_storage_rows - if min_storage_rows is not None: - min_storage_rows = int(min_storage_rows) + min_storage_bytes = parser.parse_args().min_storage_bytes + if min_storage_bytes is not None: + min_storage_bytes = int(min_storage_bytes) suffix = parser.parse_args().input_topic_suffix or '' csvfile = parser.parse_args().csv csvmetricsfile = parser.parse_args().csv_metrics @@ -375,7 +375,7 @@ def main(): "config": { "workers": cores, "storage": storage, - "min_storage_rows": min_storage_rows, + "min_storage_bytes": min_storage_bytes, "cpu_profiler": True, "resources": { # "cpu_cores_min": 0, diff --git a/crates/adapters/src/controller/mod.rs b/crates/adapters/src/controller/mod.rs index 8b6ef0d5d43..f2e8b2bd459 100644 --- a/crates/adapters/src/controller/mod.rs +++ b/crates/adapters/src/controller/mod.rs @@ -447,21 +447,21 @@ impl Controller { -> Result<(Box, Box), ControllerError>, { let mut start: Option = None; - let min_storage_rows = if controller.status.pipeline_config.global.storage { + let min_storage_bytes = if controller.status.pipeline_config.global.storage { // This reduces the files stored on disk to a reasonable number. controller .status .pipeline_config .global - .min_storage_rows - .unwrap_or(1000) + .min_storage_bytes + .unwrap_or(1024 * 1024) } else { usize::MAX }; let config = CircuitConfig { layout: Layout::new_solo(controller.status.pipeline_config.global.workers as usize), storage: controller.status.pipeline_config.storage_config.clone(), - min_storage_rows, + min_storage_bytes, init_checkpoint: Uuid::nil(), }; let mut circuit = match circuit_factory(config) { diff --git a/crates/dbsp/src/circuit/dbsp_handle.rs b/crates/dbsp/src/circuit/dbsp_handle.rs index b89a009946d..9c651c06766 100644 --- a/crates/dbsp/src/circuit/dbsp_handle.rs +++ b/crates/dbsp/src/circuit/dbsp_handle.rs @@ -215,11 +215,11 @@ pub struct CircuitConfig { pub layout: Layout, /// Storage configuration (if storage is enabled). pub storage: Option, - /// Minimum number of rows in a persistent trace to spill it to storage. If - /// this is 0, then all traces will be stored on disk; if it is - /// `usize::MAX`, then all traces will be kept in memory; and intermediate + /// Estimated minimum number of bytes in a data batch to spill it to + /// storage. If this is 0, then all batches will be stored on disk; if it is + /// `usize::MAX`, then all batches will be kept in memory; and intermediate /// values specify a threshold. - pub min_storage_rows: usize, + pub min_storage_bytes: usize, /// The initial checkpoint to start the circuit from. /// /// In case of a new circuit, this should be `Uuid::nil()`. @@ -239,7 +239,7 @@ impl CircuitConfig { Self { layout: Layout::new_solo(n), storage: None, - min_storage_rows: usize::MAX, + min_storage_bytes: usize::MAX, init_checkpoint: Uuid::nil(), } } @@ -1050,7 +1050,7 @@ mod tests { path: temp.path().to_str().unwrap().to_string(), cache: StorageCacheConfig::default(), }), - min_storage_rows: 0, + min_storage_bytes: 0, init_checkpoint: Uuid::nil(), }; (temp, cconf) diff --git a/crates/dbsp/src/circuit/runtime.rs b/crates/dbsp/src/circuit/runtime.rs index 6b3582902df..4b2d66da11a 100644 --- a/crates/dbsp/src/circuit/runtime.rs +++ b/crates/dbsp/src/circuit/runtime.rs @@ -202,7 +202,7 @@ struct RuntimeInner { layout: Layout, storage: PathBuf, cache: StorageCacheConfig, - min_storage_rows: usize, + min_storage_bytes: usize, store: LocalStore, // Panic info collected from failed worker threads. panic_info: Vec>>, @@ -272,7 +272,7 @@ impl RuntimeInner { layout: config.layout, cache, storage, - min_storage_rows: config.min_storage_rows, + min_storage_bytes: config.min_storage_bytes, store: TypedDashMap::new(), panic_info, }) @@ -606,14 +606,14 @@ impl Runtime { } } - /// Returns the minimum number of rows of a trace to spill it to + /// Returns the minimum number of bytes in a batch to spill it to /// storage. For threads that run without a runtime, this method returns /// `usize::MAX`. - pub fn min_storage_rows() -> usize { + pub fn min_storage_bytes() -> usize { RUNTIME.with(|rt| { rt.borrow() .as_ref() - .map_or(usize::MAX, |runtime| runtime.0.min_storage_rows) + .map_or(usize::MAX, |runtime| runtime.0.min_storage_bytes) }) } @@ -941,7 +941,7 @@ mod tests { path: path.to_str().unwrap().to_string(), cache: StorageCacheConfig::default(), }), - min_storage_rows: usize::MAX, + min_storage_bytes: usize::MAX, init_checkpoint: Uuid::nil(), }; @@ -963,7 +963,7 @@ mod tests { let cconf = CircuitConfig { layout: Layout::new_solo(4), storage: None, - min_storage_rows: usize::MAX, + min_storage_bytes: usize::MAX, init_checkpoint: Uuid::nil(), }; let storage_path_clone = storage_path.clone(); @@ -986,7 +986,7 @@ mod tests { let cconf = CircuitConfig { layout: Layout::new_solo(4), storage: None, - min_storage_rows: usize::MAX, + min_storage_bytes: usize::MAX, init_checkpoint: Uuid::nil(), }; let storage_path_clone = storage_path.clone(); diff --git a/crates/dbsp/src/storage/file/reader.rs b/crates/dbsp/src/storage/file/reader.rs index 6e698a6daef..ae06e9757e8 100644 --- a/crates/dbsp/src/storage/file/reader.rs +++ b/crates/dbsp/src/storage/file/reader.rs @@ -1201,6 +1201,15 @@ where { pub fn path(&self) -> PathBuf { self.0.file.as_ref().path.clone() } + + /// Returns the size of the underlying file in bytes. + pub fn byte_size(&self) -> Result { + Ok(self + .0 + .file + .cache + .get_size(self.0.file.file_handle.as_ref().unwrap())?) + } } impl Clone for Reader { diff --git a/crates/dbsp/src/trace/mod.rs b/crates/dbsp/src/trace/mod.rs index 949d8e3a17d..89f76431224 100644 --- a/crates/dbsp/src/trace/mod.rs +++ b/crates/dbsp/src/trace/mod.rs @@ -319,6 +319,19 @@ where /// The number of updates in the batch. fn len(&self) -> usize; + /// The memory or storage size of the batch in bytes. + /// + /// This can be an approximation, such as the size of an on-disk file for a + /// stored batch. + /// + /// Implementations of this function can be expensive because they might + /// require iterating through all the data in a batch. Currently this is + /// only used to decide whether to keep the result of a merge in memory or + /// on storage. For this case, the merge will visit and copy all the data + /// in the batch. The batch will be discarded afterward, which means that + /// the implementation need not attempt to cache the return value. + fn approximate_byte_size(&self) -> usize; + /// True if the batch is empty. fn is_empty(&self) -> bool { self.len() == 0 diff --git a/crates/dbsp/src/trace/ord/fallback/indexed_wset.rs b/crates/dbsp/src/trace/ord/fallback/indexed_wset.rs index c6f425d6872..0864ee4b00d 100644 --- a/crates/dbsp/src/trace/ord/fallback/indexed_wset.rs +++ b/crates/dbsp/src/trace/ord/fallback/indexed_wset.rs @@ -18,16 +18,16 @@ use crate::{ Batch, BatchFactories, BatchReader, BatchReaderFactories, Builder, FileIndexedWSet, FileIndexedWSetFactories, Filter, Merger, WeightedItem, }, - DBData, DBWeight, NumEntries, Runtime, + DBData, DBWeight, NumEntries, }; use rand::Rng; use rkyv::{ser::Serializer, Archive, Archived, Deserialize, Fallible, Serialize}; use size_of::SizeOf; use std::fmt::{self, Debug}; -use std::path::Path; +use std::{mem::replace, path::Path}; use std::{ops::Neg, path::PathBuf}; -use super::utils::GenericMerger; +use super::utils::{copy_to_builder, BuildTo, GenericMerger, MergeTo}; pub struct FallbackIndexedWSetFactories where @@ -107,16 +107,19 @@ where } } +#[derive(SizeOf)] pub struct FallbackIndexedWSet where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackIndexedWSetFactories, inner: Inner, } +#[derive(SizeOf)] #[allow(clippy::large_enum_variant)] enum Inner where @@ -326,6 +329,14 @@ where } } + #[inline] + fn approximate_byte_size(&self) -> usize { + match &self.inner { + Inner::File(file) => file.approximate_byte_size(), + Inner::Vec(vec) => vec.approximate_byte_size(), + } + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) @@ -389,16 +400,19 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FallbackIndexedWSetMerger where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackIndexedWSetFactories, inner: MergerInner, } +#[derive(SizeOf)] enum MergerInner where K: DataTrait + ?Sized, @@ -426,27 +440,22 @@ where ) -> Self { Self { factories: batch1.factories.clone(), - inner: if batch1.len() + batch2.len() < Runtime::min_storage_rows() { - match (&batch1.inner, &batch2.inner) { - (Inner::Vec(vec1), Inner::Vec(vec2)) => { - MergerInner::AllVec(VecIndexedWSetMerger::new_merger(vec1, vec2)) - } - _ => MergerInner::ToVec(GenericMerger::new( - &batch1.factories.vec, - batch1, - batch2, - )), + inner: match ( + MergeTo::from((batch1, batch2)), + &batch1.inner, + &batch2.inner, + ) { + (MergeTo::Memory, Inner::Vec(vec1), Inner::Vec(vec2)) => { + MergerInner::AllVec(VecIndexedWSetMerger::new_merger(vec1, vec2)) + } + (MergeTo::Memory, _, _) => { + MergerInner::ToVec(GenericMerger::new(&batch1.factories.vec, batch1, batch2)) } - } else { - match (&batch1.inner, &batch2.inner) { - (Inner::File(file1), Inner::File(file2)) => { - MergerInner::AllFile(FileIndexedWSetMerger::new_merger(file1, file2)) - } - _ => MergerInner::ToFile(GenericMerger::new( - &batch1.factories.file, - batch1, - batch2, - )), + (MergeTo::Storage, Inner::File(file1), Inner::File(file2)) => { + MergerInner::AllFile(FileIndexedWSetMerger::new_merger(file1, file2)) + } + (MergeTo::Storage, _, _) => { + MergerInner::ToFile(GenericMerger::new(&batch1.factories.file, batch1, batch2)) } }, } @@ -516,33 +525,20 @@ where } } -impl SizeOf for FallbackIndexedWSetMerger -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - MergerInner::AllFile(file) => file.size_of_children(context), - MergerInner::AllVec(vec) => vec.size_of_children(context), - MergerInner::ToFile(merger) => merger.size_of_children(context), - MergerInner::ToVec(merger) => merger.size_of_children(context), - } - } -} - /// A builder for batches from ordered update tuples. +#[derive(SizeOf)] pub struct FallbackIndexedWSetBuilder where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackIndexedWSetFactories, inner: BuilderInner, } +#[derive(SizeOf)] #[allow(clippy::large_enum_variant)] enum BuilderInner where @@ -550,8 +546,36 @@ where V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { - File(FileIndexedWSetBuilder), + /// In-memory. Vec(VecIndexedWSetBuilder), + + /// On-storage. + File(FileIndexedWSetBuilder), + + /// In-memory as long as we don't exceed a maximum threshold size. + Threshold { + vec: VecIndexedWSetBuilder, + + /// Bytes left to add until the threshold is exceeded. + remaining: usize, + }, +} + +impl FallbackIndexedWSetBuilder +where + Self: SizeOf, + K: DataTrait + ?Sized, + V: DataTrait + ?Sized, + R: WeightTrait + ?Sized, +{ + /// We ran out of the bytes threshold for `BuilderInner::Threshold`. Spill + /// to storage as `BuilderInner::File`, writing `vec` as the initial + /// contents. + fn spill(&mut self, vec: VecIndexedWSet) { + let mut file = FileIndexedWSetBuilder::with_capacity(&self.factories.file, (), 0); + copy_to_builder(&mut file, vec.cursor()); + self.inner = BuilderInner::File(file); + } } impl Builder> for FallbackIndexedWSetBuilder @@ -574,18 +598,17 @@ where ) -> Self { Self { factories: factories.clone(), - inner: if capacity < Runtime::min_storage_rows() { - BuilderInner::Vec(VecIndexedWSetBuilder::with_capacity( - &factories.vec, - time, - capacity, - )) - } else { - BuilderInner::File(FileIndexedWSetBuilder::with_capacity( - &factories.file, - time, - capacity, - )) + inner: match BuildTo::for_capacity( + &factories.vec, + &factories.file, + time, + capacity, + VecIndexedWSetBuilder::with_capacity, + FileIndexedWSetBuilder::with_capacity, + ) { + BuildTo::Memory(vec) => BuilderInner::Vec(vec), + BuildTo::Storage(file) => BuilderInner::File(file), + BuildTo::Threshold(vec, remaining) => BuilderInner::Threshold { vec, remaining }, }, } } @@ -593,27 +616,66 @@ where #[inline] fn reserve(&mut self, _additional: usize) {} - #[inline] fn push(&mut self, item: &mut DynPair, R>) { match &mut self.inner { BuilderInner::File(file) => file.push(item), BuilderInner::Vec(vec) => vec.push(item), + BuilderInner::Threshold { vec, remaining } => { + let size = item.size_of().total_bytes(); + vec.push(item); + if size > *remaining { + let vec = replace( + vec, + VecIndexedWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } - #[inline] fn push_refs(&mut self, key: &K, val: &V, weight: &R) { match &mut self.inner { BuilderInner::File(file) => file.push_refs(key, val, weight), BuilderInner::Vec(vec) => vec.push_refs(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key, val, weight).size_of().total_bytes(); + vec.push_refs(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecIndexedWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } - #[inline] fn push_vals(&mut self, key: &mut K, val: &mut V, weight: &mut R) { match &mut self.inner { BuilderInner::File(file) => file.push_vals(key, val, weight), BuilderInner::Vec(vec) => vec.push_vals(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key as &K, val as &V, weight as &R).size_of().total_bytes(); + vec.push_vals(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecIndexedWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -623,40 +685,14 @@ where factories: self.factories, inner: match self.inner { BuilderInner::File(file) => Inner::File(file.done()), - BuilderInner::Vec(vec) => Inner::Vec(vec.done()), + BuilderInner::Vec(vec) | BuilderInner::Threshold { vec, .. } => { + Inner::Vec(vec.done()) + } }, } } } -impl SizeOf for FallbackIndexedWSetBuilder -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - BuilderInner::File(file) => file.size_of_children(context), - BuilderInner::Vec(vec) => vec.size_of_children(context), - } - } -} - -impl SizeOf for FallbackIndexedWSet -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - Inner::Vec(vec) => vec.size_of_children(context), - Inner::File(file) => file.size_of_children(context), - } - } -} - impl Archive for FallbackIndexedWSet where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/fallback/key_batch.rs b/crates/dbsp/src/trace/ord/fallback/key_batch.rs index fb2b8bc8043..71556e46a6e 100644 --- a/crates/dbsp/src/trace/ord/fallback/key_batch.rs +++ b/crates/dbsp/src/trace/ord/fallback/key_batch.rs @@ -10,19 +10,20 @@ use crate::{ FileKeyBatch, OrdKeyBatch, }, Batch, BatchFactories, BatchReader, BatchReaderFactories, Builder, FileKeyBatchFactories, - Filter, Merger, OrdKeyBatchFactories, WeightedItem, + Filter, Merger, OrdKeyBatchFactories, TimedBuilder, WeightedItem, }, - DBData, DBWeight, NumEntries, Runtime, Timestamp, + DBData, DBWeight, NumEntries, Timestamp, }; use rand::Rng; use rkyv::{ser::Serializer, Archive, Archived, Deserialize, Fallible, Serialize}; use size_of::SizeOf; use std::{ fmt::{self, Debug}, + mem::replace, path::PathBuf, }; -use super::utils::GenericMerger; +use super::utils::{copy_to_builder, BuildTo, GenericMerger, MergeTo}; pub struct FallbackKeyBatchFactories where @@ -108,16 +109,19 @@ where /// /// Each tuple in `FallbackKeyBatch` has key type `K`, value type `()`, /// weight type `R`, and time type `R`. +#[derive(SizeOf)] pub struct FallbackKeyBatch where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackKeyBatchFactories, inner: Inner, } +#[derive(SizeOf)] enum Inner where K: DataTrait + ?Sized, @@ -242,6 +246,14 @@ where } } + #[inline] + fn approximate_byte_size(&self) -> usize { + match &self.inner { + Inner::Vec(vec) => vec.approximate_byte_size(), + Inner::File(file) => file.approximate_byte_size(), + } + } + fn lower(&self) -> AntichainRef<'_, T> { match &self.inner { Inner::Vec(vec) => vec.lower(), @@ -305,16 +317,19 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FallbackKeyMerger where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackKeyBatchFactories, inner: MergerInner, } +#[derive(SizeOf)] enum MergerInner where K: DataTrait + ?Sized, @@ -337,27 +352,22 @@ where fn new_merger(batch1: &FallbackKeyBatch, batch2: &FallbackKeyBatch) -> Self { FallbackKeyMerger { factories: batch1.factories.clone(), - inner: if batch1.len() + batch2.len() < Runtime::min_storage_rows() { - match (&batch1.inner, &batch2.inner) { - (Inner::Vec(vec1), Inner::Vec(vec2)) => { - MergerInner::AllVec(VecKeyMerger::new_merger(vec1, vec2)) - } - _ => MergerInner::ToVec(GenericMerger::new( - &batch1.factories.vec, - batch1, - batch2, - )), + inner: match ( + MergeTo::from((batch1, batch2)), + &batch1.inner, + &batch2.inner, + ) { + (MergeTo::Memory, Inner::Vec(vec1), Inner::Vec(vec2)) => { + MergerInner::AllVec(VecKeyMerger::new_merger(vec1, vec2)) + } + (MergeTo::Memory, _, _) => { + MergerInner::ToVec(GenericMerger::new(&batch1.factories.vec, batch1, batch2)) + } + (MergeTo::Storage, Inner::File(file1), Inner::File(file2)) => { + MergerInner::AllFile(FileKeyMerger::new_merger(file1, file2)) } - } else { - match (&batch1.inner, &batch2.inner) { - (Inner::File(file1), Inner::File(file2)) => { - MergerInner::AllFile(FileKeyMerger::new_merger(file1, file2)) - } - _ => MergerInner::ToFile(GenericMerger::new( - &batch1.factories.file, - batch1, - batch2, - )), + (MergeTo::Storage, _, _) => { + MergerInner::ToFile(GenericMerger::new(&batch1.factories.file, batch1, batch2)) } }, } @@ -426,41 +436,56 @@ where } } -impl SizeOf for FallbackKeyMerger +/// A builder for creating layers from unsorted update tuples. +#[derive(SizeOf)] +pub struct FallbackKeyBuilder where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - MergerInner::AllFile(file) => file.size_of_children(context), - MergerInner::AllVec(vec) => vec.size_of_children(context), - MergerInner::ToFile(merger) => merger.size_of_children(context), - MergerInner::ToVec(merger) => merger.size_of_children(context), - } - } + #[size_of(skip)] + factories: FallbackKeyBatchFactories, + inner: BuilderInner, } -/// A builder for creating layers from unsorted update tuples. -pub struct FallbackKeyBuilder +#[derive(SizeOf)] +enum BuilderInner where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - factories: FallbackKeyBatchFactories, - inner: BuilderInner, + /// In-memory. + Vec(VecKeyBuilder), + + /// On-storage. + File(FileKeyBuilder), + + /// In-memory as long as we don't exceed a maximum threshold size. + Threshold { + vec: VecKeyBuilder, + + /// Bytes left to add until the threshold is exceeded. + remaining: usize, + }, } -enum BuilderInner +impl FallbackKeyBuilder where + Self: SizeOf, K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - File(FileKeyBuilder), - Vec(VecKeyBuilder), + /// We ran out of the bytes threshold for `BuilderInner::Threshold`. Spill + /// to storage as `BuilderInner::File`, writing `vec` as the initial + /// contents. + fn spill(&mut self, vec: OrdKeyBatch) { + let mut file = FileKeyBuilder::timed_with_capacity(&self.factories.file, 0); + copy_to_builder(&mut file, vec.cursor()); + self.inner = BuilderInner::File(file); + } } impl Builder> for FallbackKeyBuilder @@ -483,14 +508,17 @@ where ) -> Self { Self { factories: factories.clone(), - inner: if capacity < Runtime::min_storage_rows() { - BuilderInner::Vec(VecKeyBuilder::with_capacity(&factories.vec, time, capacity)) - } else { - BuilderInner::File(FileKeyBuilder::with_capacity( - &factories.file, - time, - capacity, - )) + inner: match BuildTo::for_capacity( + &factories.vec, + &factories.file, + time, + capacity, + VecKeyBuilder::with_capacity, + FileKeyBuilder::with_capacity, + ) { + BuildTo::Memory(vec) => BuilderInner::Vec(vec), + BuildTo::Storage(file) => BuilderInner::File(file), + BuildTo::Threshold(vec, remaining) => BuilderInner::Threshold { vec, remaining }, }, } } @@ -503,6 +531,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push(item), BuilderInner::Vec(vec) => vec.push(item), + BuilderInner::Threshold { vec, remaining } => { + let size = item.size_of().total_bytes(); + vec.push(item); + if size > *remaining { + let vec = replace( + vec, + VecKeyBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -511,6 +553,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_refs(key, val, weight), BuilderInner::Vec(vec) => vec.push_refs(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key, val, weight).size_of().total_bytes(); + vec.push_refs(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecKeyBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -519,6 +575,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_vals(key, val, weight), BuilderInner::Vec(vec) => vec.push_vals(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key as &K, weight as &R).size_of().total_bytes(); + vec.push_vals(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecKeyBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -528,34 +598,14 @@ where factories: self.factories, inner: match self.inner { BuilderInner::File(file) => Inner::File(file.done()), - BuilderInner::Vec(vec) => Inner::Vec(vec.done()), + BuilderInner::Vec(vec) | BuilderInner::Threshold { vec, .. } => { + Inner::Vec(vec.done()) + } }, } } } -impl SizeOf for FallbackKeyBuilder -where - K: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FallbackKeyBatch -where - K: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - impl Archive for FallbackKeyBatch where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/fallback/utils.rs b/crates/dbsp/src/trace/ord/fallback/utils.rs index f4e93f90152..ad945d1d788 100644 --- a/crates/dbsp/src/trace/ord/fallback/utils.rs +++ b/crates/dbsp/src/trace/ord/fallback/utils.rs @@ -11,7 +11,7 @@ use crate::{ ord::filter, Batch, BatchReader, BatchReaderFactories, Builder, Cursor, Filter, TimedBuilder, }, - Timestamp, + Runtime, Timestamp, }; /// The row position of a [`Cursor`], regardless of the underlying type of the @@ -20,6 +20,7 @@ use crate::{ /// [`GenericMerger`] uses this to save and restore positions in the batches /// it's merging, since it can't keep a cursor around from one run to another /// because of lifetime issues. +#[derive(SizeOf)] enum Position where K: DataTrait + ?Sized, @@ -64,6 +65,7 @@ where } } +#[derive(SizeOf)] pub(super) struct GenericMerger where K: DataTrait + ?Sized, @@ -319,16 +321,100 @@ where } } -impl SizeOf for GenericMerger +/// Reads all of the data from `cursor` and writes it to `builder`. +pub(super) fn copy_to_builder(builder: &mut B, mut cursor: C) where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, + B: TimedBuilder, + Output: Batch, + C: HasTimeDiffCursor, + K: ?Sized, + V: ?Sized, R: WeightTrait + ?Sized, - O: Batch + BatchReader, - O::Builder: TimedBuilder, { - fn size_of_children(&self, context: &mut size_of::Context) { - self.builder.size_of_children(context) + let mut tmp = cursor.weight_factory().default_box(); + while cursor.key_valid() { + while cursor.val_valid() { + let mut td_cursor = cursor.time_diff_cursor(); + while let Some((time, diff)) = td_cursor.current(&mut tmp) { + builder.push_time(cursor.key(), cursor.val(), time, diff); + td_cursor.step(); + } + drop(td_cursor); + cursor.step_val(); + } + cursor.step_key(); + } +} + +pub(super) enum MergeTo { + Memory, + Storage, +} + +impl From<(&B, &B)> for MergeTo +where + B: BatchReader, +{ + fn from((batch1, batch2): (&B, &B)) -> Self { + // This is equivalent to `batch1.byte_size() + batch2.byte_size() >= + // Runtime::min_storage_bytes()` but it avoids calling `byte_size()` any + // more than necessary since it can be expensive. + let spill = match Runtime::min_storage_bytes() { + 0 => true, + usize::MAX => false, + min_storage_bytes => { + let size1 = batch1.approximate_byte_size(); + size1 >= min_storage_bytes || { + let size2 = batch2.approximate_byte_size(); + size1 + size2 >= min_storage_bytes + } + } + }; + + if spill { + Self::Storage + } else { + Self::Memory + } + } +} + +pub(super) enum BuildTo { + Memory(M), + Storage(S), + Threshold(M, usize), +} + +impl BuildTo { + pub fn for_capacity( + vf: MF, + sf: SF, + time: T, + capacity: usize, + mc: MC, + sc: SC, + ) -> Self + where + MC: Fn(MF, T, usize) -> M, + SC: Fn(SF, T, usize) -> S, + { + match Runtime::min_storage_bytes() { + usize::MAX => { + // Storage is disabled. + Self::Memory(mc(vf, time, capacity)) + } + + min_storage_bytes if capacity >= min_storage_bytes => { + // If `capacity` is filled up then we'll have at least 1 byte + // per item (as a bottom of the barrel estimate) so we might as + // well start out on storage. + Self::Storage(sc(sf, time, capacity)) + } + min_storage_bytes => { + // Start out in memory and spill to storage if + // `min_storage_bytes` is used. + Self::Threshold(mc(vf, time, capacity), min_storage_bytes) + } + } } } diff --git a/crates/dbsp/src/trace/ord/fallback/val_batch.rs b/crates/dbsp/src/trace/ord/fallback/val_batch.rs index 4057f06b690..204436fd6d2 100644 --- a/crates/dbsp/src/trace/ord/fallback/val_batch.rs +++ b/crates/dbsp/src/trace/ord/fallback/val_batch.rs @@ -1,9 +1,11 @@ use std::fmt::{Display, Formatter}; +use std::mem::replace; use std::path::PathBuf; use crate::trace::cursor::DelegatingCursor; use crate::trace::ord::file::val_batch::FileValBuilder; use crate::trace::ord::vec::val_batch::VecValBuilder; +use crate::trace::TimedBuilder; use crate::{ dynamic::{DataTrait, DynPair, DynVec, DynWeightedPairs, Erase, Factory, WeightTrait}, time::AntichainRef, @@ -15,13 +17,13 @@ use crate::{ Batch, BatchFactories, BatchReader, BatchReaderFactories, Builder, FileValBatch, FileValBatchFactories, Filter, Merger, OrdValBatch, OrdValBatchFactories, WeightedItem, }, - DBData, DBWeight, NumEntries, Runtime, Timestamp, + DBData, DBWeight, NumEntries, Timestamp, }; use rand::Rng; use rkyv::{ser::Serializer, Archive, Archived, Deserialize, Fallible, Serialize}; use size_of::SizeOf; -use super::utils::GenericMerger; +use super::utils::{copy_to_builder, BuildTo, GenericMerger, MergeTo}; pub struct FallbackValBatchFactories where @@ -107,6 +109,7 @@ where /// An immutable collection of update tuples, from a contiguous interval of /// logical times. +#[derive(SizeOf)] pub struct FallbackValBatch where K: DataTrait + ?Sized, @@ -114,10 +117,12 @@ where T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackValBatchFactories, inner: Inner, } +#[derive(SizeOf)] enum Inner where K: DataTrait + ?Sized, @@ -250,6 +255,14 @@ where } } + #[inline] + fn approximate_byte_size(&self) -> usize { + match &self.inner { + Inner::File(file) => file.approximate_byte_size(), + Inner::Vec(vec) => vec.approximate_byte_size(), + } + } + fn lower(&self) -> AntichainRef<'_, T> { match &self.inner { Inner::Vec(vec) => vec.lower(), @@ -314,6 +327,7 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FallbackValMerger where K: DataTrait + ?Sized, @@ -321,10 +335,12 @@ where T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackValBatchFactories, inner: MergerInner, } +#[derive(SizeOf)] enum MergerInner where K: DataTrait + ?Sized, @@ -352,27 +368,22 @@ where ) -> Self { FallbackValMerger { factories: batch1.factories.clone(), - inner: if batch1.len() + batch2.len() < Runtime::min_storage_rows() { - match (&batch1.inner, &batch2.inner) { - (Inner::Vec(vec1), Inner::Vec(vec2)) => { - MergerInner::AllVec(VecValMerger::new_merger(vec1, vec2)) - } - _ => MergerInner::ToVec(GenericMerger::new( - &batch1.factories.vec, - batch1, - batch2, - )), + inner: match ( + MergeTo::from((batch1, batch2)), + &batch1.inner, + &batch2.inner, + ) { + (MergeTo::Memory, Inner::Vec(vec1), Inner::Vec(vec2)) => { + MergerInner::AllVec(VecValMerger::new_merger(vec1, vec2)) + } + (MergeTo::Memory, _, _) => { + MergerInner::ToVec(GenericMerger::new(&batch1.factories.vec, batch1, batch2)) } - } else { - match (&batch1.inner, &batch2.inner) { - (Inner::File(file1), Inner::File(file2)) => { - MergerInner::AllFile(FileValMerger::new_merger(file1, file2)) - } - _ => MergerInner::ToFile(GenericMerger::new( - &batch1.factories.file, - batch1, - batch2, - )), + (MergeTo::Storage, Inner::File(file1), Inner::File(file2)) => { + MergerInner::AllFile(FileValMerger::new_merger(file1, file2)) + } + (MergeTo::Storage, _, _) => { + MergerInner::ToFile(GenericMerger::new(&batch1.factories.file, batch1, batch2)) } }, } @@ -441,39 +452,59 @@ where } } -impl SizeOf for FallbackValMerger +/// A builder for creating layers from unsorted update tuples. +#[derive(SizeOf)] +pub struct FallbackValBuilder where K: DataTrait + ?Sized, V: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } + #[size_of(skip)] + factories: FallbackValBatchFactories, + inner: BuilderInner, } -/// A builder for creating layers from unsorted update tuples. -pub struct FallbackValBuilder +#[derive(SizeOf)] +enum BuilderInner where K: DataTrait + ?Sized, V: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - factories: FallbackValBatchFactories, - inner: BuilderInner, + /// In-memory. + Vec(VecValBuilder), + + /// On-storage. + File(FileValBuilder), + + /// In-memory as long as we don't exceed a maximum threshold size. + Threshold { + vec: VecValBuilder, + + /// Bytes left to add until the threshold is exceeded. + remaining: usize, + }, } -enum BuilderInner +impl FallbackValBuilder where + Self: SizeOf, K: DataTrait + ?Sized, V: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { - File(FileValBuilder), - Vec(VecValBuilder), + /// We ran out of the bytes threshold for `BuilderInner::Threshold`. Spill + /// to storage as `BuilderInner::File`, writing `vec` as the initial + /// contents. + fn spill(&mut self, vec: OrdValBatch) { + let mut file = FileValBuilder::timed_with_capacity(&self.factories.file, 0); + copy_to_builder(&mut file, vec.cursor()); + self.inner = BuilderInner::File(file); + } } impl Builder> for FallbackValBuilder @@ -495,14 +526,17 @@ where ) -> Self { Self { factories: factories.clone(), - inner: if capacity < Runtime::min_storage_rows() { - BuilderInner::Vec(VecValBuilder::with_capacity(&factories.vec, time, capacity)) - } else { - BuilderInner::File(FileValBuilder::with_capacity( - &factories.file, - time, - capacity, - )) + inner: match BuildTo::for_capacity( + &factories.vec, + &factories.file, + time, + capacity, + VecValBuilder::with_capacity, + FileValBuilder::with_capacity, + ) { + BuildTo::Memory(vec) => BuilderInner::Vec(vec), + BuildTo::Storage(file) => BuilderInner::File(file), + BuildTo::Threshold(vec, remaining) => BuilderInner::Threshold { vec, remaining }, }, } } @@ -513,6 +547,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push(item), BuilderInner::Vec(vec) => vec.push(item), + BuilderInner::Threshold { vec, remaining } => { + let size = item.size_of().total_bytes(); + vec.push(item); + if size > *remaining { + let vec = replace( + vec, + VecValBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -520,6 +568,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_refs(key, val, weight), BuilderInner::Vec(vec) => vec.push_refs(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key, val, weight).size_of().total_bytes(); + vec.push_refs(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecValBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -527,6 +589,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_vals(key, val, weight), BuilderInner::Vec(vec) => vec.push_vals(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key as &K, val as &V, weight as &R).size_of().total_bytes(); + vec.push_vals(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecValBuilder::timed_with_capacity(&self.factories.vec, 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -535,36 +611,14 @@ where factories: self.factories, inner: match self.inner { BuilderInner::File(file) => Inner::File(file.done()), - BuilderInner::Vec(vec) => Inner::Vec(vec.done()), + BuilderInner::Vec(vec) | BuilderInner::Threshold { vec, .. } => { + Inner::Vec(vec.done()) + } }, } } } -impl SizeOf for FallbackValBuilder -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FallbackValBatch -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - impl Archive for FallbackValBatch where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/fallback/wset.rs b/crates/dbsp/src/trace/ord/fallback/wset.rs index 1c8574dd770..b45ced3cdaa 100644 --- a/crates/dbsp/src/trace/ord/fallback/wset.rs +++ b/crates/dbsp/src/trace/ord/fallback/wset.rs @@ -15,15 +15,18 @@ use crate::{ Batch, BatchFactories, BatchReader, BatchReaderFactories, Builder, FileWSet, FileWSetFactories, Filter, Merger, WeightedItem, }, - DBData, DBWeight, NumEntries, Runtime, + DBData, DBWeight, NumEntries, }; use rand::Rng; use rkyv::{ser::Serializer, Archive, Archived, Deserialize, Fallible, Serialize}; use size_of::SizeOf; -use std::fmt::{self, Debug}; +use std::{ + fmt::{self, Debug}, + mem::replace, +}; use std::{ops::Neg, path::PathBuf}; -use super::utils::GenericMerger; +use super::utils::{copy_to_builder, BuildTo, GenericMerger, MergeTo}; pub struct FallbackWSetFactories where @@ -101,11 +104,13 @@ where } } +#[derive(SizeOf)] pub struct FallbackWSet where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackWSetFactories, inner: Inner, } @@ -124,6 +129,7 @@ where } } +#[derive(SizeOf)] #[allow(clippy::large_enum_variant)] enum Inner where @@ -319,6 +325,14 @@ where } } + #[inline] + fn approximate_byte_size(&self) -> usize { + match &self.inner { + Inner::File(file) => file.approximate_byte_size(), + Inner::Vec(vec) => vec.approximate_byte_size(), + } + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) @@ -374,15 +388,18 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FallbackWSetMerger where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackWSetFactories, inner: MergerInner, } +#[derive(SizeOf)] enum MergerInner where K: DataTrait + ?Sized, @@ -404,27 +421,22 @@ where fn new_merger(batch1: &FallbackWSet, batch2: &FallbackWSet) -> Self { Self { factories: batch1.factories.clone(), - inner: if batch1.len() + batch2.len() < Runtime::min_storage_rows() { - match (&batch1.inner, &batch2.inner) { - (Inner::Vec(vec1), Inner::Vec(vec2)) => { - MergerInner::AllVec(VecWSetMerger::new_merger(vec1, vec2)) - } - _ => MergerInner::ToVec(GenericMerger::new( - &batch1.factories.vec, - batch1, - batch2, - )), + inner: match ( + MergeTo::from((batch1, batch2)), + &batch1.inner, + &batch2.inner, + ) { + (MergeTo::Memory, Inner::Vec(vec1), Inner::Vec(vec2)) => { + MergerInner::AllVec(VecWSetMerger::new_merger(vec1, vec2)) + } + (MergeTo::Memory, _, _) => { + MergerInner::ToVec(GenericMerger::new(&batch1.factories.vec, batch1, batch2)) } - } else { - match (&batch1.inner, &batch2.inner) { - (Inner::File(file1), Inner::File(file2)) => { - MergerInner::AllFile(FileWSetMerger::new_merger(file1, file2)) - } - _ => MergerInner::ToFile(GenericMerger::new( - &batch1.factories.file, - batch1, - batch2, - )), + (MergeTo::Storage, Inner::File(file1), Inner::File(file2)) => { + MergerInner::AllFile(FileWSetMerger::new_merger(file1, file2)) + } + (MergeTo::Storage, _, _) => { + MergerInner::ToFile(GenericMerger::new(&batch1.factories.file, batch1, batch2)) } }, } @@ -494,39 +506,54 @@ where } } -impl SizeOf for FallbackWSetMerger -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - MergerInner::AllFile(file) => file.size_of_children(context), - MergerInner::AllVec(vec) => vec.size_of_children(context), - MergerInner::ToFile(merger) => merger.size_of_children(context), - MergerInner::ToVec(merger) => merger.size_of_children(context), - } - } -} - /// A builder for batches from ordered update tuples. +#[derive(SizeOf)] pub struct FallbackWSetBuilder where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FallbackWSetFactories, inner: BuilderInner, } +#[derive(SizeOf)] #[allow(clippy::large_enum_variant)] enum BuilderInner where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { - File(FileWSetBuilder), + /// In-memory. Vec(VecWSetBuilder), + + /// On-storage. + File(FileWSetBuilder), + + /// In-memory as long as we don't exceed a maximum threshold size. + Threshold { + vec: VecWSetBuilder, + + /// Bytes left to add until the threshold is exceeded. + remaining: usize, + }, +} + +impl FallbackWSetBuilder +where + Self: SizeOf, + K: DataTrait + ?Sized, + R: WeightTrait + ?Sized, +{ + /// We ran out of the bytes threshold for `BuilderInner::Threshold`. Spill + /// to storage as `BuilderInner::File`, writing `vec` as the initial + /// contents. + fn spill(&mut self, vec: VecWSet) { + let mut file = FileWSetBuilder::with_capacity(&self.factories.file, (), 0); + copy_to_builder(&mut file, vec.cursor()); + self.inner = BuilderInner::File(file); + } } impl Builder> for FallbackWSetBuilder @@ -544,18 +571,17 @@ where fn with_capacity(factories: &FallbackWSetFactories, time: (), capacity: usize) -> Self { Self { factories: factories.clone(), - inner: if capacity < Runtime::min_storage_rows() { - BuilderInner::Vec(VecWSetBuilder::with_capacity( - &factories.vec, - time, - capacity, - )) - } else { - BuilderInner::File(FileWSetBuilder::with_capacity( - &factories.file, - time, - capacity, - )) + inner: match BuildTo::for_capacity( + &factories.vec, + &factories.file, + time, + capacity, + VecWSetBuilder::with_capacity, + FileWSetBuilder::with_capacity, + ) { + BuildTo::Memory(vec) => BuilderInner::Vec(vec), + BuildTo::Storage(file) => BuilderInner::File(file), + BuildTo::Threshold(vec, remaining) => BuilderInner::Threshold { vec, remaining }, }, } } @@ -568,6 +594,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push(item), BuilderInner::Vec(vec) => vec.push(item), + BuilderInner::Threshold { vec, remaining } => { + let size = item.size_of().total_bytes(); + vec.push(item); + if size > *remaining { + let vec = replace( + vec, + VecWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -576,6 +616,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_refs(key, val, weight), BuilderInner::Vec(vec) => vec.push_refs(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key, weight).size_of().total_bytes(); + vec.push_refs(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -584,6 +638,20 @@ where match &mut self.inner { BuilderInner::File(file) => file.push_vals(key, val, weight), BuilderInner::Vec(vec) => vec.push_vals(key, val, weight), + BuilderInner::Threshold { vec, remaining } => { + let size = (key as &K, weight as &R).size_of().total_bytes(); + vec.push_vals(key, val, weight); + if size > *remaining { + let vec = replace( + vec, + VecWSetBuilder::with_capacity(&self.factories.vec, (), 0), + ) + .done(); + self.spill(vec); + } else { + *remaining -= size; + } + } } } @@ -593,38 +661,14 @@ where factories: self.factories, inner: match self.inner { BuilderInner::File(file) => Inner::File(file.done()), - BuilderInner::Vec(vec) => Inner::Vec(vec.done()), + BuilderInner::Vec(vec) | BuilderInner::Threshold { vec, .. } => { + Inner::Vec(vec.done()) + } }, } } } -impl SizeOf for FallbackWSetBuilder -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - BuilderInner::File(file) => file.size_of_children(context), - BuilderInner::Vec(vec) => vec.size_of_children(context), - } - } -} - -impl SizeOf for FallbackWSet -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, context: &mut size_of::Context) { - match &self.inner { - Inner::Vec(vec) => vec.size_of_children(context), - Inner::File(file) => file.size_of_children(context), - } - } -} - impl Archive for FallbackWSet where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/file/indexed_wset_batch.rs b/crates/dbsp/src/trace/ord/file/indexed_wset_batch.rs index 44d5df387b9..c54bc996cec 100644 --- a/crates/dbsp/src/trace/ord/file/indexed_wset_batch.rs +++ b/crates/dbsp/src/trace/ord/file/indexed_wset_batch.rs @@ -137,14 +137,17 @@ where /// /// Each tuple in `FileIndexedWSet` has key type `K`, value type `V`, /// weight type `R`, and time `()`. +#[derive(SizeOf)] pub struct FileIndexedWSet where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileIndexedWSetFactories, #[allow(clippy::type_complexity)] + #[size_of(skip)] file: Reader<(&'static K, &'static DynUnit, (&'static V, &'static R, ()))>, lower_bound: usize, } @@ -355,6 +358,10 @@ where self.file.n_rows(1) as usize } + fn approximate_byte_size(&self) -> usize { + self.file.byte_size().unwrap() as usize + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) @@ -442,12 +449,14 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FileIndexedWSetMerger where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileIndexedWSetFactories, // Position in first batch. @@ -456,6 +465,7 @@ where lower2: usize, // Output so far. + #[size_of(skip)] writer: Writer2, } @@ -636,17 +646,6 @@ where } } -impl SizeOf for FileIndexedWSetMerger -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - type KeyCursor<'s, K, V, R> = FileCursor< 's, K, @@ -878,13 +877,16 @@ where } /// A builder for batches from ordered update tuples. +#[derive(SizeOf)] pub struct FileIndexedWSetBuilder where K: DataTrait + ?Sized, V: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileIndexedWSetFactories, + #[size_of(skip)] writer: Writer2, key: Box>, } @@ -981,28 +983,6 @@ where } } -impl SizeOf for FileIndexedWSetBuilder -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FileIndexedWSet -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - impl Archive for FileIndexedWSet where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/file/key_batch.rs b/crates/dbsp/src/trace/ord/file/key_batch.rs index 7c095142f13..2cb5c81a983 100644 --- a/crates/dbsp/src/trace/ord/file/key_batch.rs +++ b/crates/dbsp/src/trace/ord/file/key_batch.rs @@ -133,14 +133,17 @@ where /// /// Each tuple in `FileKeyBatch` has key type `K`, value type `()`, /// weight type `R`, and time type `R`. +#[derive(SizeOf)] pub struct FileKeyBatch where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileKeyBatchFactories, #[allow(clippy::type_complexity)] + #[size_of(skip)] file: Reader<( &'static K, &'static DynUnit, @@ -251,6 +254,10 @@ where self.file.n_rows(1) as usize } + fn approximate_byte_size(&self) -> usize { + self.file.byte_size().unwrap() as usize + } + fn lower(&self) -> AntichainRef<'_, T> { self.lower.as_ref() } @@ -348,12 +355,14 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FileKeyMerger where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileKeyBatchFactories, lower: Antichain, upper: Antichain, @@ -364,6 +373,7 @@ where lower2: usize, // Output so far. + #[size_of(skip)] writer: Writer2, R>, } @@ -539,17 +549,6 @@ where } } -impl SizeOf for FileKeyMerger -where - K: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - type RawKeyCursor<'s, K, T, R> = FileCursor< 's, K, @@ -828,14 +827,17 @@ where } /// A builder for creating layers from unsorted update tuples. +#[derive(SizeOf)] pub struct FileKeyBuilder where K: DataTrait + ?Sized, T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileKeyBatchFactories, time: T, + #[size_of(skip)] writer: Writer2, R>, key: Box>, } @@ -951,28 +953,6 @@ where } } -impl SizeOf for FileKeyBuilder -where - K: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FileKeyBatch -where - K: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - impl Archive for FileKeyBatch where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/file/val_batch.rs b/crates/dbsp/src/trace/ord/file/val_batch.rs index 7c7b77fa494..59ba1f97888 100644 --- a/crates/dbsp/src/trace/ord/file/val_batch.rs +++ b/crates/dbsp/src/trace/ord/file/val_batch.rs @@ -183,6 +183,7 @@ type RawValCursor<'s, K, V, T, R> = FileCursor< /// An immutable collection of update tuples, from a contiguous interval of /// logical times. +#[derive(SizeOf)] pub struct FileValBatch where K: DataTrait + ?Sized, @@ -190,7 +191,9 @@ where T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileValBatchFactories, + #[size_of(skip)] pub file: RawValBatch, pub lower_bound: usize, pub lower: Antichain, @@ -278,6 +281,10 @@ where self.file.n_rows(1) as usize } + fn approximate_byte_size(&self) -> usize { + self.file.byte_size().unwrap() as usize + } + fn lower(&self) -> AntichainRef<'_, T> { self.lower.as_ref() } @@ -407,6 +414,7 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FileValMerger where K: DataTrait + ?Sized, @@ -414,7 +422,9 @@ where T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileValBatchFactories, + #[size_of(skip)] result: Option>, lower: Antichain, upper: Antichain, @@ -702,18 +712,6 @@ where } } -impl SizeOf for FileValMerger -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - #[derive(SizeOf)] pub struct FileValCursor<'s, K, V, T, R> where @@ -1016,6 +1014,7 @@ where } /// A builder for creating layers from unsorted update tuples. +#[derive(SizeOf)] pub struct FileValBuilder where K: DataTrait + ?Sized, @@ -1023,8 +1022,10 @@ where T: Timestamp, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileValBatchFactories, time: T, + #[size_of(skip)] writer: Writer2, R>>, cur_key: Box>, } @@ -1133,18 +1134,6 @@ where } } -impl SizeOf for FileValBuilder -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - /* pub struct FileValConsumer { __type: PhantomData<(K, V, T, R)>, @@ -1194,18 +1183,6 @@ impl<'a, K, V, T, R> ValueConsumer<'a, V, R, T> for FileValValueConsumer<'a, K, } */ -impl SizeOf for FileValBatch -where - K: DataTrait + ?Sized, - V: DataTrait + ?Sized, - T: Timestamp, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - impl Archive for FileValBatch where K: DataTrait + ?Sized, diff --git a/crates/dbsp/src/trace/ord/file/wset_batch.rs b/crates/dbsp/src/trace/ord/file/wset_batch.rs index 872938b2dc2..7bd7d2cf5b2 100644 --- a/crates/dbsp/src/trace/ord/file/wset_batch.rs +++ b/crates/dbsp/src/trace/ord/file/wset_batch.rs @@ -127,12 +127,15 @@ where /// /// Each tuple in `FileWSet` has key type `K`, value type `()`, weight /// type `R`, and time type `()`. +#[derive(SizeOf)] pub struct FileWSet where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileWSetFactories, + #[size_of(skip)] file: Reader<(&'static K, &'static R, ())>, lower_bound: usize, } @@ -349,6 +352,10 @@ where self.key_count() } + fn approximate_byte_size(&self) -> usize { + self.file.byte_size().unwrap() as usize + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) @@ -435,11 +442,13 @@ where } /// State for an in-progress merge. +#[derive(SizeOf)] pub struct FileWSetMerger where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileWSetFactories, // Position in first batch. @@ -448,6 +457,7 @@ where lower2: usize, // Output so far. + #[size_of(skip)] writer: Writer1, } @@ -749,12 +759,15 @@ where } /// A builder for creating layers from unsorted update tuples. +#[derive(SizeOf)] pub struct FileWSetBuilder where K: DataTrait + ?Sized, R: WeightTrait + ?Sized, { + #[size_of(skip)] factories: FileWSetFactories, + #[size_of(skip)] writer: Writer1, } @@ -828,33 +841,3 @@ where self.done() } } - -impl SizeOf for FileWSetBuilder -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FileWSetMerger -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} - -impl SizeOf for FileWSet -where - K: DataTrait + ?Sized, - R: WeightTrait + ?Sized, -{ - fn size_of_children(&self, _context: &mut size_of::Context) { - // XXX - } -} diff --git a/crates/dbsp/src/trace/ord/vec/indexed_wset_batch.rs b/crates/dbsp/src/trace/ord/vec/indexed_wset_batch.rs index f10232ea4de..d43c5bd43b5 100644 --- a/crates/dbsp/src/trace/ord/vec/indexed_wset_batch.rs +++ b/crates/dbsp/src/trace/ord/vec/indexed_wset_batch.rs @@ -438,6 +438,11 @@ where self.layer.tuples() } + #[inline] + fn approximate_byte_size(&self) -> usize { + self.size_of().total_bytes() + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) diff --git a/crates/dbsp/src/trace/ord/vec/key_batch.rs b/crates/dbsp/src/trace/ord/vec/key_batch.rs index 9237809b393..dbab1644990 100644 --- a/crates/dbsp/src/trace/ord/vec/key_batch.rs +++ b/crates/dbsp/src/trace/ord/vec/key_batch.rs @@ -299,6 +299,10 @@ where as Trie>::tuples(&self.layer) } + fn approximate_byte_size(&self) -> usize { + self.size_of().total_bytes() + } + fn lower(&self) -> AntichainRef<'_, T> { self.lower.as_ref() } diff --git a/crates/dbsp/src/trace/ord/vec/val_batch.rs b/crates/dbsp/src/trace/ord/vec/val_batch.rs index 773763e0e6e..451c8680fab 100644 --- a/crates/dbsp/src/trace/ord/vec/val_batch.rs +++ b/crates/dbsp/src/trace/ord/vec/val_batch.rs @@ -368,6 +368,10 @@ where as Trie>::tuples(&self.layer) } + fn approximate_byte_size(&self) -> usize { + self.size_of().total_bytes() + } + fn lower(&self) -> AntichainRef<'_, T> { self.lower.as_ref() } diff --git a/crates/dbsp/src/trace/ord/vec/wset_batch.rs b/crates/dbsp/src/trace/ord/vec/wset_batch.rs index 6a022fda7e9..dbfd13fabce 100644 --- a/crates/dbsp/src/trace/ord/vec/wset_batch.rs +++ b/crates/dbsp/src/trace/ord/vec/wset_batch.rs @@ -332,6 +332,11 @@ impl BatchReader for VecWSet usize { + self.size_of().total_bytes() + } + #[inline] fn lower(&self) -> AntichainRef<'_, ()> { AntichainRef::new(&[()]) diff --git a/crates/dbsp/src/trace/spine_async/mod.rs b/crates/dbsp/src/trace/spine_async/mod.rs index 0d242b77b89..36a92c510d4 100644 --- a/crates/dbsp/src/trace/spine_async/mod.rs +++ b/crates/dbsp/src/trace/spine_async/mod.rs @@ -342,6 +342,10 @@ where self.fold_batches(0, |acc, batch| acc + batch.len()) } + fn approximate_byte_size(&self) -> usize { + self.fold_batches(0, |acc, batch| acc + batch.approximate_byte_size()) + } + fn lower(&self) -> AntichainRef<'_, Self::Time> { self.lower.as_ref() } diff --git a/crates/dbsp/src/trace/spine_fueled.rs b/crates/dbsp/src/trace/spine_fueled.rs index 4b4498d0557..293f73835f7 100644 --- a/crates/dbsp/src/trace/spine_fueled.rs +++ b/crates/dbsp/src/trace/spine_fueled.rs @@ -290,6 +290,10 @@ where self.fold_batches(0, |acc, batch| acc + batch.len()) } + fn approximate_byte_size(&self) -> usize { + self.fold_batches(0, |acc, batch| acc + batch.approximate_byte_size()) + } + fn lower(&self) -> AntichainRef<'_, Self::Time> { self.lower.as_ref() } diff --git a/crates/dbsp/src/trace/test/test_batch.rs b/crates/dbsp/src/trace/test/test_batch.rs index 160f4e99800..35240aacc72 100644 --- a/crates/dbsp/src/trace/test/test_batch.rs +++ b/crates/dbsp/src/trace/test/test_batch.rs @@ -1165,6 +1165,10 @@ where self.data.len() } + fn approximate_byte_size(&self) -> usize { + self.size_of().total_bytes() + } + fn lower(&self) -> AntichainRef<'_, Self::Time> { todo!() } diff --git a/crates/nexmark/benches/nexmark/main.rs b/crates/nexmark/benches/nexmark/main.rs index 115829c7412..b8195df5511 100644 --- a/crates/nexmark/benches/nexmark/main.rs +++ b/crates/nexmark/benches/nexmark/main.rs @@ -223,7 +223,7 @@ fn create_ascii_table(config: &NexmarkConfig) -> AsciiTable { ]; let mut max_width = 200; - if config.min_storage_rows != usize::MAX { + if config.min_storage_bytes != usize::MAX { result_columns.extend_from_slice(&[ "# Files", "Avg WrSz", @@ -266,7 +266,7 @@ fn run_query(config: &NexmarkConfig, snapshotter: &Snapshotter, query: Query) -> StorageCacheConfig::PageCache }, }), - min_storage_rows: config.min_storage_rows, + min_storage_bytes: config.min_storage_bytes, ..CircuitConfig::with_workers(num_cores) }; let (dbsp, input_handle) = @@ -491,7 +491,7 @@ fn main() -> Result<()> { format!("{}", HumanBytes::from(after.peak_rss)), format!("{}", after.page_faults - before.page_faults), ]; - if nexmark_config.min_storage_rows != usize::MAX { + if nexmark_config.min_storage_bytes != usize::MAX { row.extend_from_slice(&[ format!("{}", diff.n_created), format!("{}", HumanBytes::from(diff.avg_wblock)), diff --git a/crates/nexmark/src/config.rs b/crates/nexmark/src/config.rs index 37171aa116b..6a9313eb274 100644 --- a/crates/nexmark/src/config.rs +++ b/crates/nexmark/src/config.rs @@ -51,11 +51,11 @@ pub struct Config { pub progress: bool, /// Configure storage. Without this option, data is kept in memory. With - /// this option, if `ROWS` is 0 or omitted, all data is written to storage; + /// this option, if `BYTES` is 0 or omitted, all data is written to storage; /// otherwise, data batches are written to storage when they are estimated - /// to contain at least `ROWS` rows. - #[clap(long="storage", default_value_t = usize::MAX, default_missing_value = "0", hide_default_value(true), value_name("ROWS"))] - pub min_storage_rows: usize, + /// to contain at least `BYTES` bytes. + #[clap(long="storage", default_value_t = usize::MAX, default_missing_value = "0", hide_default_value(true), value_name("BYTES"))] + pub min_storage_bytes: usize, /// Storage cache configuration. By default, the benchmark uses the /// operating system page cache. This option enables the internal Feldera @@ -156,7 +156,7 @@ impl Default for Config { input_batch_size: 40_000, output_csv: None, progress: true, - min_storage_rows: usize::MAX, + min_storage_bytes: usize::MAX, feldera_cache: false, } } diff --git a/crates/pipeline-types/src/config.rs b/crates/pipeline-types/src/config.rs index 30c927962f2..d8b4eb723e7 100644 --- a/crates/pipeline-types/src/config.rs +++ b/crates/pipeline-types/src/config.rs @@ -158,14 +158,15 @@ pub struct RuntimeConfig { #[serde(default)] pub resources: ResourceConfig, - /// The minimum estimated number of rows in a batch to write it to storage. - /// This is provided for debugging and fine-tuning and should ordinarily be - /// left unset. It only has an effect when `storage` is set to true. + /// The minimum estimated number of bytes in a batch of data to write it to + /// storage. This is provided for debugging and fine-tuning and should + /// ordinarily be left unset. It only has an effect when `storage` is set to + /// true. /// /// A value of 0 will write even empty batches to storage, and nonzero /// values provide a threshold. `usize::MAX` would effectively disable /// storage. - pub min_storage_rows: Option, + pub min_storage_bytes: Option, } impl RuntimeConfig { diff --git a/crates/pipeline_manager/src/db/test.rs b/crates/pipeline_manager/src/db/test.rs index 05290a7d6c0..bde6b33114d 100644 --- a/crates/pipeline_manager/src/db/test.rs +++ b/crates/pipeline_manager/src/db/test.rs @@ -984,7 +984,7 @@ async fn versioning() { min_batch_size_records: 0, max_buffering_delay_usecs: 0, resources: ResourceConfig::default(), - min_storage_rows: None, + min_storage_bytes: None, }; handle .db @@ -1300,7 +1300,7 @@ pub(crate) fn runtime_config() -> impl Strategy { storage_mb_max: config.9, storage_class: None, }, - min_storage_rows: None, + min_storage_bytes: None, }) } @@ -1335,7 +1335,7 @@ pub(crate) fn option_runtime_config() -> impl Strategy