Skip to main content

faucet_core/
object_rollover.rs

1//! Cross-page accumulation and rollover for object-store sinks (#618).
2//!
3//! An object-store sink used to write **one object per `write_batch` page**
4//! with nothing retained between calls, so a small `batch_size` produced a
5//! swarm of tiny objects — the small-files problem that dominates read time on
6//! S3/Athena/Spark, where per-object overhead outweighs the bytes. The Parquet
7//! sink already accumulates across pages and rolls over on a row *or* byte
8//! threshold; this is that logic, extracted so every object-store sink shares
9//! one definition rather than four near-copies.
10//!
11//! ## What this owns and what it does not
12//!
13//! This is the **decision**, not the I/O: it says when the open object is full
14//! and hands back the bytes to upload. Each sink keeps its own upload — the
15//! vendor SDKs differ too much to unify here, and pushing the upload behind a
16//! trait would buy an abstraction whose only implementors are four call sites.
17//!
18//! ## Why a byte threshold and not only rows
19//!
20//! Rows are a poor proxy for object size: 10k rows of a wide table and 10k
21//! rows of `{"id":1}` differ by orders of magnitude, so a rows-only cap either
22//! writes tiny objects for narrow data or unbounded ones for wide data. The
23//! byte threshold is what actually bounds peak memory, since the open object's
24//! body is buffered until it rolls.
25
26use serde_json::Value;
27
28/// Accumulates encoded rows for one open object, and says when to roll.
29///
30/// The sink pushes each page's rows in, takes whatever completed objects come
31/// back, and calls [`finish`](Self::finish) at flush time for the partial
32/// remainder.
33#[derive(Debug)]
34pub struct ObjectAccumulator {
35    buf: Vec<u8>,
36    rows: usize,
37    /// Bytes already handed out as parts for the open object, so the byte cap
38    /// measures the whole object rather than just the unflushed tail.
39    parted_bytes: usize,
40    max_rows: Option<usize>,
41    max_bytes: Option<usize>,
42    part_bytes: Option<usize>,
43}
44
45/// One object's worth of accumulated bytes, ready to upload.
46#[derive(Debug, PartialEq, Eq)]
47pub struct CompletedObject {
48    /// Encoded body (for JSONL: one record per line, trailing newline).
49    pub body: Vec<u8>,
50    /// Records the body holds — for logging and metrics, and so a caller can
51    /// assert no row was lost across a rollover.
52    pub rows: usize,
53}
54
55/// What a push produced.
56///
57/// Three outcomes rather than an `Option`, because "a part is ready to upload"
58/// and "the object is finished" are different instructions to the sink: the
59/// first frees memory mid-object, the second closes it.
60#[derive(Debug, PartialEq, Eq)]
61pub enum Emit {
62    /// Keep accumulating.
63    Nothing,
64    /// A multipart part is full. Upload it and drop it — this is what keeps
65    /// peak memory at O(part size) instead of O(object size) when the object
66    /// cap is large or unset.
67    Part(Vec<u8>),
68    /// The object reached its record/byte cap.
69    Object(CompletedObject),
70}
71
72impl ObjectAccumulator {
73    /// Build an accumulator. `max_rows`/`max_bytes` of `None` or `0` mean "no
74    /// limit on this axis"; with neither set, nothing ever rolls and the whole
75    /// run lands in one object at `finish`.
76    pub fn new(max_rows: Option<usize>, max_bytes: Option<usize>) -> Self {
77        Self {
78            buf: Vec::new(),
79            rows: 0,
80            parted_bytes: 0,
81            max_rows: max_rows.filter(|n| *n > 0),
82            max_bytes: max_bytes.filter(|n| *n > 0),
83            part_bytes: None,
84        }
85    }
86
87    /// Emit [`Emit::Part`]s once the unflushed tail reaches `bytes`, so the
88    /// sink can stream them into a multipart upload and free the memory.
89    ///
90    /// Without this, an object with a large (or absent) byte cap is held whole
91    /// in RAM before its single-shot upload — which is what caps output size
92    /// by available memory. `0` disables parting.
93    pub fn with_part_size(mut self, bytes: usize) -> Self {
94        self.part_bytes = (bytes > 0).then_some(bytes);
95        self
96    }
97
98    /// Bytes of the open object already uploaded as parts.
99    pub fn parted_bytes(&self) -> usize {
100        self.parted_bytes
101    }
102
103    /// Whether any part of the open object has already been uploaded — the
104    /// sink needs this to know whether to complete a multipart upload or do a
105    /// single-shot put.
106    pub fn has_parts(&self) -> bool {
107        self.parted_bytes > 0
108    }
109
110    /// Rows currently held in the open object.
111    pub fn rows(&self) -> usize {
112        self.rows
113    }
114
115    /// Bytes currently held in the open object.
116    pub fn len(&self) -> usize {
117        self.buf.len()
118    }
119
120    /// Whether nothing is buffered.
121    pub fn is_empty(&self) -> bool {
122        self.rows == 0
123    }
124
125    /// Append one encoded record (caller-encoded, so the accumulator stays
126    /// format-agnostic) and return any object this completed.
127    ///
128    /// The threshold is checked **after** appending, so a single record larger
129    /// than `max_bytes` still lands in its own object rather than being
130    /// refused or split — splitting would corrupt it, and refusing would drop
131    /// data the source really produced.
132    pub fn push_encoded(&mut self, encoded: &[u8]) -> Emit {
133        self.buf.extend_from_slice(encoded);
134        self.rows += 1;
135        // The object cap is checked first: an object that is finished should
136        // be closed, not parted one push before closing.
137        if self.max_rows.is_some_and(|m| self.rows >= m)
138            || self
139                .max_bytes
140                .is_some_and(|m| self.parted_bytes + self.buf.len() >= m)
141        {
142            return Emit::Object(self.take());
143        }
144        if self.part_bytes.is_some_and(|m| self.buf.len() >= m) {
145            let part = std::mem::take(&mut self.buf);
146            self.parted_bytes += part.len();
147            return Emit::Part(part);
148        }
149        Emit::Nothing
150    }
151
152    /// Append one JSON record as an NDJSON line.
153    pub fn push_record(&mut self, record: &Value) -> Result<Emit, crate::FaucetError> {
154        let mut line = serde_json::to_vec(record)
155            .map_err(|e| crate::FaucetError::Sink(format!("JSON serialization failed: {e}")))?;
156        line.push(b'\n');
157        Ok(self.push_encoded(&line))
158    }
159
160    /// Close the open object, if any. Called at `flush`, and on drop-time
161    /// finalisation — an object left unfinished is data loss, so a sink must
162    /// never skip it.
163    ///
164    /// Returns `Some` whenever the object holds rows, **including when the
165    /// unflushed tail is empty but parts were already uploaded**: a multipart
166    /// upload still has to be completed.
167    pub fn finish(&mut self) -> Option<CompletedObject> {
168        (!self.is_empty()).then(|| self.take())
169    }
170
171    fn take(&mut self) -> CompletedObject {
172        self.parted_bytes = 0;
173        CompletedObject {
174            body: std::mem::take(&mut self.buf),
175            rows: std::mem::replace(&mut self.rows, 0),
176        }
177    }
178}
179
180#[cfg(test)]
181mod tests {
182    use super::*;
183    use serde_json::json;
184
185    fn rec(i: u64) -> Value {
186        json!({ "id": i })
187    }
188
189    /// Push `n` records and collect whatever the accumulator emitted.
190    fn push_all(acc: &mut ObjectAccumulator, n: u64) -> (Vec<CompletedObject>, Vec<Vec<u8>>) {
191        let (mut objects, mut parts) = (Vec::new(), Vec::new());
192        for i in 0..n {
193            match acc.push_record(&rec(i)).unwrap() {
194                Emit::Nothing => {}
195                Emit::Part(p) => parts.push(p),
196                Emit::Object(o) => objects.push(o),
197            }
198        }
199        (objects, parts)
200    }
201
202    #[test]
203    fn nothing_rolls_without_a_threshold() {
204        // The "one object for the whole run" configuration — what an operator
205        // who wants the fewest possible files asks for.
206        let mut acc = ObjectAccumulator::new(None, None);
207        let (objects, parts) = push_all(&mut acc, 1000);
208        assert!(objects.is_empty() && parts.is_empty());
209        let done = acc.finish().expect("the remainder must be written");
210        assert_eq!(done.rows, 1000);
211        assert!(acc.finish().is_none(), "finish must not double-emit");
212    }
213
214    #[test]
215    fn rows_across_several_pages_land_in_one_object() {
216        // The headline fix: a small `batch_size` used to mean one object per
217        // page. Three pages of 4 under a 10-row cap must produce one full
218        // object plus a remainder, not three objects.
219        let mut acc = ObjectAccumulator::new(Some(10), None);
220        let mut completed = Vec::new();
221        for page in 0..3 {
222            for i in 0..4 {
223                if let Emit::Object(o) = acc.push_record(&rec(page * 4 + i)).unwrap() {
224                    completed.push(o);
225                }
226            }
227        }
228        assert_eq!(completed.len(), 1, "one rollover at 10 rows");
229        assert_eq!(completed[0].rows, 10);
230        let rest = acc.finish().expect("2 rows remain");
231        assert_eq!(rest.rows, 2);
232        assert_eq!(
233            completed[0].rows + rest.rows,
234            12,
235            "no row may be lost across a rollover"
236        );
237    }
238
239    #[test]
240    fn the_byte_threshold_rolls_independently_of_rows() {
241        // Rows are a poor proxy for size; this is the axis that actually
242        // bounds peak memory.
243        let mut acc = ObjectAccumulator::new(None, Some(32));
244        let (objects, _) = push_all(&mut acc, 20);
245        assert!(!objects.is_empty(), "a byte cap must roll");
246        for obj in &objects {
247            assert!(
248                obj.body.len() >= 32,
249                "an object rolls at or past the threshold, not before: {}",
250                obj.body.len()
251            );
252            assert!(obj.body.len() < 32 + 64, "overshoot bounded by one record");
253        }
254    }
255
256    #[test]
257    fn whichever_threshold_hits_first_wins() {
258        let mut acc = ObjectAccumulator::new(Some(1000), Some(24));
259        let (objects, _) = push_all(&mut acc, 10);
260        assert!(
261            !objects.is_empty(),
262            "the byte cap must roll even though the row cap is far away"
263        );
264    }
265
266    #[test]
267    fn a_single_oversized_record_gets_its_own_object() {
268        // Splitting it would corrupt it and refusing it would drop data the
269        // source really produced, so the only correct answer is one object
270        // over the threshold.
271        let mut acc = ObjectAccumulator::new(None, Some(8));
272        let big = json!({ "blob": "x".repeat(500) });
273        let Emit::Object(obj) = acc.push_record(&big).unwrap() else {
274            panic!("an oversized record must complete an object immediately");
275        };
276        assert_eq!(obj.rows, 1);
277        assert!(obj.body.len() > 8);
278    }
279
280    #[test]
281    fn a_zero_threshold_means_no_limit_not_roll_every_record() {
282        // `0` is the house "no limit" sentinel; reading it as "roll always"
283        // would turn an existing config into one object per record.
284        let mut acc = ObjectAccumulator::new(Some(0), Some(0));
285        let (objects, parts) = push_all(&mut acc, 50);
286        assert!(objects.is_empty() && parts.is_empty());
287        assert_eq!(acc.finish().expect("remainder").rows, 50);
288    }
289
290    #[test]
291    fn bodies_are_ndjson_with_a_trailing_newline() {
292        let mut acc = ObjectAccumulator::new(Some(2), None);
293        acc.push_record(&rec(1)).unwrap();
294        let Emit::Object(obj) = acc.push_record(&rec(2)).unwrap() else {
295            panic!("rolled");
296        };
297        let text = String::from_utf8(obj.body).unwrap();
298        assert_eq!(text, "{\"id\":1}\n{\"id\":2}\n");
299    }
300
301    #[test]
302    fn counters_track_the_open_object() {
303        let mut acc = ObjectAccumulator::new(Some(10), None);
304        assert!(acc.is_empty());
305        acc.push_record(&rec(1)).unwrap();
306        assert_eq!(acc.rows(), 1);
307        assert!(!acc.is_empty());
308        acc.finish();
309        assert!(acc.is_empty(), "finish resets the open object");
310        assert_eq!(acc.rows(), 0);
311        assert_eq!(acc.len(), 0);
312    }
313
314    #[test]
315    fn parts_bound_peak_memory_for_an_uncapped_object() {
316        // The case the byte cap cannot help with: "one object for the whole
317        // run". Without parting, the entire object sits in RAM before a
318        // single-shot upload, which is what caps output size by memory.
319        let mut acc = ObjectAccumulator::new(None, None).with_part_size(64);
320        let (objects, parts) = push_all(&mut acc, 200);
321        assert!(objects.is_empty(), "no object cap, so nothing rolls");
322        assert!(!parts.is_empty(), "parts must be emitted");
323        for p in &parts {
324            assert!(p.len() >= 64, "a part fills before it is emitted");
325            assert!(p.len() < 64 + 64, "and is not held far past the threshold");
326        }
327        assert!(
328            acc.len() < 64,
329            "the unflushed tail stays under one part: {}",
330            acc.len()
331        );
332        assert!(acc.has_parts());
333        // Every byte is accounted for: parts + tail = the whole object.
334        let tail = acc.finish().expect("tail");
335        assert_eq!(tail.rows, 200, "rows count the whole object, not the tail");
336    }
337
338    #[test]
339    fn the_byte_cap_measures_the_whole_object_not_just_the_tail() {
340        // A part'd object must still roll at `max_bytes`; measuring only the
341        // unflushed tail would make the cap unreachable and produce one
342        // unbounded object.
343        let mut acc = ObjectAccumulator::new(None, Some(200)).with_part_size(64);
344        let (objects, parts) = push_all(&mut acc, 200);
345        assert!(!parts.is_empty(), "parts still stream");
346        assert!(
347            !objects.is_empty(),
348            "the object cap must still be reached once parts are counted"
349        );
350    }
351
352    #[test]
353    fn taking_an_object_resets_the_part_counter() {
354        // Otherwise the next object would inherit the previous one's parted
355        // bytes and roll early — and, worse, `has_parts` would tell the sink
356        // to complete a multipart upload that was never started.
357        let mut acc = ObjectAccumulator::new(Some(2), None).with_part_size(8);
358        acc.push_record(&rec(1)).unwrap();
359        let Emit::Object(_) = acc.push_record(&rec(2)).unwrap() else {
360            panic!("rolled at 2 rows");
361        };
362        assert_eq!(acc.parted_bytes(), 0);
363        assert!(!acc.has_parts());
364    }
365
366    #[test]
367    fn an_object_cap_closes_rather_than_parting_on_the_same_push() {
368        // A push that both fills a part and finishes the object must close it:
369        // emitting a part first would leave a completed object needing a
370        // second, empty finish.
371        let mut acc = ObjectAccumulator::new(Some(1), None).with_part_size(1);
372        assert!(matches!(acc.push_record(&rec(1)).unwrap(), Emit::Object(_)));
373    }
374}
375
376/// Cross-page record accumulation for **warehouse** sinks (#617).
377///
378/// The object-store sinks above accumulate encoded bytes; a warehouse sink
379/// needs the records themselves, because its commit is a statement it builds
380/// from them (`INSERT … SELECT FROM UNNEST`, `COPY`, `INSERT … FORMAT
381/// JSONEachRow`). The problem is the same one: the commit unit was the *page*
382/// unit, so a small `batch_size` meant one expensive warehouse operation per
383/// small page — slow and costly everywhere, and a hard "too many parts"
384/// failure on ClickHouse.
385///
386/// ## Why this is append-only
387///
388/// Accumulation is applied to the plain `write_batch` path and **not** to
389/// `write_batch_idempotent` or `write_batch_partial`:
390///
391/// - a commit token must land atomically with *its own* page, so deferring the
392///   write past the page the token names would make the watermark a lie;
393/// - a DLQ needs to know which rows of *this* page failed, which a deferred,
394///   merged commit cannot report.
395///
396/// This is the same carve-out the BigQuery sink makes for `write_batch_partial`
397/// (#612), for the same reasons.
398#[derive(Debug)]
399pub struct PageAccumulator {
400    rows: Vec<Value>,
401    bytes: usize,
402    max_rows: Option<usize>,
403    max_bytes: Option<usize>,
404}
405
406impl PageAccumulator {
407    /// `max_rows`/`max_bytes` of `None` or `0` mean "no limit on this axis".
408    /// With neither set, the whole run commits once at `flush`.
409    pub fn new(max_rows: Option<usize>, max_bytes: Option<usize>) -> Self {
410        Self {
411            rows: Vec::new(),
412            bytes: 0,
413            max_rows: max_rows.filter(|n| *n > 0),
414            max_bytes: max_bytes.filter(|n| *n > 0),
415        }
416    }
417
418    /// Records held for the next commit.
419    pub fn len(&self) -> usize {
420        self.rows.len()
421    }
422
423    /// Whether nothing is buffered.
424    pub fn is_empty(&self) -> bool {
425        self.rows.is_empty()
426    }
427
428    /// Estimated serialized size of the buffered records.
429    pub fn bytes(&self) -> usize {
430        self.bytes
431    }
432
433    /// Add a page and return the group to commit, if one is now full.
434    ///
435    /// The whole page is added before the check, so a group can overshoot by
436    /// at most one page — splitting a page across two commits would break the
437    /// DLQ's per-page row indices for any caller that later wants them, and
438    /// buys nothing when the threshold is a soft target.
439    pub fn push_page(&mut self, page: &[Value]) -> Option<Vec<Value>> {
440        for r in page {
441            // An estimate, not a measurement: serializing every record twice
442            // (once to size it, once to commit it) would cost more than the
443            // precision is worth for a rollover threshold.
444            self.bytes += estimate_size(r);
445            self.rows.push(r.clone());
446        }
447        let full = self.max_rows.is_some_and(|m| self.rows.len() >= m)
448            || self.max_bytes.is_some_and(|m| self.bytes >= m);
449        full.then(|| self.take())
450    }
451
452    /// Take whatever is buffered, for the `flush`-time commit.
453    pub fn finish(&mut self) -> Option<Vec<Value>> {
454        (!self.is_empty()).then(|| self.take())
455    }
456
457    fn take(&mut self) -> Vec<Value> {
458        self.bytes = 0;
459        std::mem::take(&mut self.rows)
460    }
461}
462
463/// Rough serialized size of a JSON value, without serializing it.
464///
465/// Used only to decide when a commit group is big enough, so it trades
466/// accuracy for not walking every value twice. Numbers and booleans are
467/// charged a flat width; strings their length plus quoting.
468fn estimate_size(v: &Value) -> usize {
469    match v {
470        Value::Null => 4,
471        Value::Bool(_) => 5,
472        Value::Number(_) => 8,
473        Value::String(s) => s.len() + 2,
474        Value::Array(a) => 2 + a.iter().map(estimate_size).sum::<usize>() + a.len(),
475        Value::Object(m) => {
476            2 + m
477                .iter()
478                .map(|(k, v)| k.len() + 3 + estimate_size(v))
479                .sum::<usize>()
480        }
481    }
482}
483
484#[cfg(test)]
485mod page_accumulator_tests {
486    use super::*;
487    use serde_json::json;
488
489    fn page(n: usize) -> Vec<Value> {
490        (0..n).map(|i| json!({ "id": i })).collect()
491    }
492
493    #[test]
494    fn small_pages_merge_into_one_commit_group() {
495        // The headline fix: `batch_size` could only ever *split* a page, so
496        // ten 10-row pages meant ten warehouse operations.
497        let mut acc = PageAccumulator::new(Some(100), None);
498        let mut commits = Vec::new();
499        for _ in 0..10 {
500            if let Some(g) = acc.push_page(&page(10)) {
501                commits.push(g);
502            }
503        }
504        assert_eq!(commits.len(), 1, "ten small pages → one commit");
505        assert_eq!(commits[0].len(), 100);
506        assert!(acc.finish().is_none(), "nothing left over");
507    }
508
509    #[test]
510    fn a_group_overshoots_by_at_most_one_page() {
511        // Pages are never split across commits, so the threshold is a soft
512        // target — pinned so a future "exact" rewrite has to justify itself.
513        let mut acc = PageAccumulator::new(Some(10), None);
514        let g = acc
515            .push_page(&page(25))
516            .expect("one page over the cap commits");
517        assert_eq!(g.len(), 25);
518    }
519
520    #[test]
521    fn the_byte_cap_rolls_independently_of_rows() {
522        let mut acc = PageAccumulator::new(None, Some(200));
523        let mut commits = 0;
524        for _ in 0..20 {
525            if acc.push_page(&page(5)).is_some() {
526                commits += 1;
527            }
528        }
529        assert!(commits > 0, "a byte cap must roll");
530    }
531
532    #[test]
533    fn no_cap_means_one_commit_for_the_whole_run() {
534        let mut acc = PageAccumulator::new(None, None);
535        for _ in 0..50 {
536            assert!(acc.push_page(&page(10)).is_none());
537        }
538        assert_eq!(acc.finish().expect("flush commits").len(), 500);
539    }
540
541    #[test]
542    fn a_zero_cap_means_no_limit() {
543        // `0` is the house "no limit" sentinel; reading it as "commit always"
544        // would restore the per-page behaviour this exists to remove.
545        let mut acc = PageAccumulator::new(Some(0), Some(0));
546        assert!(acc.push_page(&page(10)).is_none());
547        assert_eq!(acc.finish().expect("remainder").len(), 10);
548    }
549
550    #[test]
551    fn counters_track_the_open_group() {
552        let mut acc = PageAccumulator::new(Some(100), None);
553        assert!(acc.is_empty());
554        acc.push_page(&page(3));
555        assert_eq!(acc.len(), 3);
556        assert!(acc.bytes() > 0);
557        acc.finish();
558        assert!(acc.is_empty());
559        assert_eq!(acc.bytes(), 0, "finish resets the size estimate too");
560    }
561
562    #[test]
563    fn size_estimate_grows_with_real_content() {
564        // It only has to be monotonic in the data to be useful as a threshold.
565        let small = estimate_size(&json!({ "a": 1 }));
566        let big = estimate_size(&json!({ "a": 1, "b": "x".repeat(1000) }));
567        assert!(big > small + 900, "{small} vs {big}");
568        assert!(estimate_size(&json!([1, 2, 3])) > estimate_size(&json!([])));
569    }
570}