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}