Commit f0a2d99
feat: streaming column-major merge engine with page-bounded body cols (PR-6b.2)
Rebuilds PR-6b on top of PR-6a.2's per-page Arrow decoder. The
streaming merge engine now keeps body-col memory bounded by output
page size (not column-chunk size) while preserving caller-specified
M:N output splitting at sorted_series boundaries.
Architecture (Husky multi-input → multi-output sorted merge):
Phase 0 (async) — drain sort cols from each input. With Husky
column ordering, sort cols + sorted_series are the prefix of each
row group's body bytes, so the decoder can stop after they are
fully decoded; the remaining body col pages stay un-read in the
input stream, ready for phase 3.
Phase 1 — compute_merge_order over the per-input sort-col
RecordBatches using the existing k-way (sorted_series,
timestamp_secs) heap.
Phase 2 — compute_output_boundaries with the caller's
num_outputs, splitting at sorted_series transitions.
Phase 3 (blocking + block_on bridges) — streaming write. All M
output writers are alive for the duration. For each column in
Husky order, every output's col K is written in turn:
- Sort col / sorted_series: applied via arrow::interleave from
the already-buffered phase-0 data.
- Body col: each output page is assembled via arrow::interleave
from input page slices, with decoders advanced page-by-page via
handle.block_on from inside the sync iterator passed to
write_next_column_arrays. Pages flush to the writer's sink as
SerializedColumnWriter's page-size threshold trips — memory
stays bounded by the in-flight output page plus a small number
of in-flight input pages.
After all M outputs' col K is done, every input decoder is at the
start of col K+1 in its single row group. Move to col K+1.
PR-6b.2 only handles single-row-group inputs (real or PR-5-
adapter-presented). Multi-RG metric-aligned inputs are rejected
with a clear error message; supporting them requires either
consuming + discarding body cols of RG[i-1] from the stream to
reach RG[i]'s sort cols, or a second body GET — both are larger
scope changes that land in a follow-up.
Page-bounded contract verified by
test_body_col_streams_many_pages_per_column_chunk: with
data_page_row_count_limit=1000 on an 8000-row merge, the output
value column spans ≥ 2 pages, demonstrating that body col writes
respect data_page_size and do not materialise whole column chunks.
Tests (9, all passing): two-input merge, single-RG output for
single-metric_name input, total-rows-preserved across M:N,
sort-schema mismatch rejection, KV metadata propagation,
all-empty-inputs no-output, output drainable by StreamDecoder,
multi-RG input rejection, page-bounded body col streaming.
Also exposes existing helpers in merge/writer.rs as pub(super)
(apply_merge_permutation, build_merge_kv_metadata,
build_sorting_columns, resolve_sort_field_names, verify_sort_order)
so streaming.rs can reuse the same MC-3 / KV / sorting-columns
construction the non-streaming engine uses. PR-7 will fold the
non-streaming engine away.
PR-6c.2 will add file-size monitoring on top: close the current
output at the next sorted_series transition when an in-progress
file approaches the size cap, producing additional splits beyond
the caller's N.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>1 parent 898a709 commit f0a2d99
5 files changed
Lines changed: 2258 additions & 10 deletions
File tree
- quickwit/quickwit-parquet-engine/src
- merge
- storage
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
24 | 24 | | |
25 | 25 | | |
26 | 26 | | |
| 27 | + | |
27 | 28 | | |
28 | 29 | | |
29 | 30 | | |
| |||
0 commit comments