Repository navigation
fix(engine): stream unranked scans and hydrate return-only columns after the limit - #887
Conversation
…ter the limit A read that projected a wide column failed with `query_memory_bytes` once the type's projected columns exceeded the 150 MiB query pool, even for `limit 1`: an unranked table `ScanExec` streamed only when a contains join marked it, and otherwise collected every Lance batch of the type and concatenated them before the limit or top-k above it ran. Every unranked table scan now takes the pipelined path: byte-capped Lance batches, bounded read-ahead and an explicit 64 MiB Lance I/O buffer, each batch charged until its consumer takes it. Ranked scans keep their breaker for the overfetch ladder; only a marked scan records the runtime-filter counters. Pass 6 (late_materialization) now also runs on query plans. Under a root `limit K`, a binding whose scan estimates more than 4K rows reads its deciding columns plus `_rowaddr`, and a new `HydrateColumns` node above the limit fetches the columns only the output reads by row address for the kept rows (`HydrateExec`, chunked under an eighth of the query pool). One query reads one pinned snapshot, so the address is valid for the take. Tests: a slow GQT case at scale (red on v0.12.0), a 16 MiB-pool engine test, a fast GQT case over hash-join and id-lookup destinations compared with engine v1, planner rule, bound-plan, lowering, report and replay tests, and a `hydrate $var: columns [..]` plan claim.
…y-memory-failure-bdbef8
ragnorc
left a comment
There was a problem hiding this comment.
Recommendation: request changes for the hydration memory bound described inline. Reviewed baec750d8c1e1d4c959ab38e884fbc63bfdf3963. This is a COMMENT review because the authenticated account owns the PR.
A query can ask for one document but scan a table full of large text or vectors. Previously, an unranked scan held the table before the limit could stop it. This PR streams those scans. It also carries row addresses through filters, joins, and sorting. After the limit chooses the result rows, one new operator reads their return-only columns. This lets small answers avoid reading and retaining large amounts of payload data.
The planner pass addresses the cause. It keeps row selection separate from payload retrieval. The same accepted snapshot supplies both reads. Snapshot::open_lance_dataset opens the manifest-pinned version, and Lance 11's address take uses physical addresses within that dataset. Compaction on a later version does not justify reading a different version here. I found no new publication or policy path.
The target workload is a large table with wide values and a small result. On file, S3, or Azure storage, it trades a second, selective read for less payload I/O and decoded memory. The four-to-one row threshold does not estimate value width, filter selectivity, or remote request latency. It is a heuristic, not evidence that every qualifying query is faster. The 64 MiB scan I/O buffer and two decode tasks reduce read-ahead, but can reduce throughput on high-latency storage. This buffer is separate from the query pool. Concurrent scan buffers also cost headroom: the PR raises the single-hop memory fixture from 16 to 20 MiB. The final result still must fit the existing collection and concatenation budget.
Liability is mixed. Reusing the scan producer and pinned snapshot removes a whole-table retention requirement. One hydration operator composes better than separate fixes for each query shape. However, this adds 1,750 net lines, a serialized plan variant, hidden address columns, positional output reconstruction, and chunk-sizing rules. Those rules need one memory-admission authority. Five similar operators with private size checks would increase maintenance and failure risk. Reuse the existing checked Arrow take path and account for repeated output rows before allocation.
Validation:
- Local, using compiled artifacts from a clean checkout at this exact head: 49 planner query tests passed. The fast hydration logic case and the 200 MB wide-column logic case passed. The two focused scan/hydration memory tests and the hydration replay test passed.
- Fresh isolated build with a temporary assertion in
engine_v2_memory.rs: the fanout chunk-bound assertion failed. A 16 MiB pool returned 108 copies of a 40 KiB payload in three hydration chunks, including the four-row seed. The remaining 104 rows need at least three more chunks to stay below 2 MiB each. All returned payloads had the expected size. The original test passed before this extension. I restored the temporary source edit. - Exact-head CI: main checks and GQ logic tests passed. The DST suite failed only
dst_bench_cost_count_golden: 50 passed, one failed, 28 ignored._auditGETs changed from 706 to 734, and_verifyGETs from 1,349 to 1,416. Explain those extra reads before accepting a new golden count. - I checked the pinned Lance 11 take/scanner code and Arrow 58.3 string-take allocation, plus the relevant upstream guides. Local checks used filesystem storage. Cloud throughput, peak process RSS, and performance across workload distributions remain unmeasured.
The inline memory finding blocks my recommendation. The cost-golden failure also needs resolution before merge. No other blocking defect found in this scope.
…y-memory-failure-bdbef8
Lance reads one batch per fragment, and the pipelined scan sizes its reads for 4 KiB rows, so a streamed scan over many small fragments or narrow rows sent many tiny batches. A per-batch consumer (an indexed hop's edge lookup) then paid once per fragment; the DST cost golden caught the extra reads. Kept rows now gather under one `v2 scan batch` charge until they reach the session's batch rows (or the rows of a `Limit` directly over the scan's projection) or the byte target, and leave as one batch. Two contains-join tests counted batches on tables smaller than one gathered batch; their fixtures now exceed a scan batch so the assertions still prove the scan streams.
ragnorc
left a comment
There was a problem hiding this comment.
Recommendation: request changes. The existing P2 hydration finding still applies. This follow-up reviews 4d49aa813379501865efe43e5dc809863d09bbca. I submit a COMMENT review because the authenticated account owns the PR.
This PR lets a query select a few rows without first holding a table of large payloads. It streams unranked scans, carries row addresses through row selection, and reads return-only values after the limit. The new commit fixes a cost of that streaming design: small stored fragments produced small batches, and each batch could trigger another graph lookup. The scan now gathers kept rows into larger batches before it sends them downstream.
The new gathering path fits fragmented tables with many small writes. It keeps the input order, retains memory charges, and uses the existing checked WorkMemory::concat before allocation. It flushes the final partial batch. A limit hint reduces the row target for a scan directly below a projection and limit. It changes batch size, not row selection. The accepted snapshot and pinned Lance dataset remain the read authority. I found no new publication or policy path.
The tradeoff is fewer downstream lookups for an extra copy when several batches combine. The first result can also arrive later. The row and byte targets are flush thresholds: one incoming batch can take the total above either target. The shared query pool still governs memory admission for this gathering path. Wide values tend to hit the byte target sooner, so they gain less from row gathering. The limit hint covers one plan shape; it is not a general limit pushdown. These choices apply to file, S3, and Azure reads, but the benefit depends on fragment size, row width, and storage latency.
Liability improves in the batching part: one scan boundary controls batch shape, so downstream operators do not each need a small-fragment workaround. The new commit adds 170 lines and removes 12, including tests. That cost is reasonable for a local helper that reuses the existing memory and producer mechanisms. Five copies of this policy in different operators would be harder to maintain. Keep this shared ownership.
The remaining hydration path breaks that same design principle. Its size check counts distinct rows fetched from Lance. Its Arrow take copies payloads for repeated result rows before the pool admits those copies. This file is unchanged from the first review. Size the expanded output and reserve its memory before allocation. Reuse the existing checked take mechanism, with the required null-index behavior, rather than adding another private memory rule.
Validation:
- Local, fresh build at this exact head with Rust 1.97.1: the new small-fragment gathering test passed. The existing wide-column memory test passed before and after the temporary probe. Documentation checks passed for 218 Markdown files. Formatting and repository instruction-link checks passed.
- The temporary fanout assertion failed again. A 16 MiB pool returned 108 copies of a 40 KiB payload in three hydration chunks, including the four-row seed. The remaining 104 rows need at least three more chunks to keep each below 2 MiB. Returned values were correct. I restored all source edits and rebuilt the original test. The isolated checkout is clean.
- Exact-head DST CI passed, including
dst_bench_cost_count_golden. The main DST target reported 51 passed and 28 ignored. This resolves the cost-golden failure from the first review. GQT CI also passed. In main CI, the workspace tests, lint, and format checks passed. The S3 failpoint and Azurite jobs were still running at the final check. - Code inspection covered the merge from main, the new batching delta, callers, memory leases, and cancellation. The merge tree matches Git's automatic merge. I checked the pinned Lance 11.0.0 scanner implementation and the relevant upstream documentation. Lance's strict row batching cannot be combined with its byte target. These local tests use filesystem storage. I did not measure cloud latency, peak process RSS, or throughput across workload distributions.
No new blocking defect found in the batching delta. The existing hydration finding still blocks approval. I did not add a duplicate inline comment.
… them `HydrateExec` sized a chunk by the distinct rows Lance returned, then copied their values to every output row with a bare Arrow `take`. A row a join repeats is copied once per repetition, so the copies could pass the chunk bound, and the pool charged them only after they existed. The chunk now takes its copies through `WorkMemory::take_once`, which admits them before Arrow allocates, and keeps a prefix, halved until the pool's estimate of its copies fits beside the fetched rows. The checked take counts a null index as a null row rather than the row its raw value names. A `peak_chunk_bytes` gauge records the most one chunk held.
…y-memory-failure-bdbef8 # Conflicts: # crates/omnigraph-gqt-core/src/plan.rs # crates/omnigraph-planner/src/lib.rs # crates/omnigraph-planner/tests/bound_plan.rs # crates/omnigraph-planner/tests/lower_walk.rs # crates/omnigraph/src/engine/lower.rs # crates/omnigraph/src/engine/scan.rs
…y-memory-failure-bdbef8
Why
On 0.12.0 a read that projects a wide column fails with
resource limit exceeded for query_memory_bytesonce the type's projected columns exceed the150 MiB query pool, even when it returns one row. A list view over a node type
of about 85,000 rows whose text column holds message and transcript bodies
fails this way for
limit 1,limit 10andorder … limit 101alike.The cause is the leaf. An unranked table
ScanExecstreamed only when acontains join marked it; every other table scan was a breaker that collected
every Lance batch of the type (Lance's 8,192-row default, no byte cap) and then
concatenated them under a 2x admission.
LimitExecand the top-kSortExecalready stream, but they received a table-sized batch. Engine v1 had the same
collect-and-concatenate scan with no budget, which is why 0.11 answered.
What changes
The property this establishes: a read's resident memory follows its result
window, not its tables.
ScanExectakes the existingpipelined path for every table read under no search mode: byte-capped Lance
batches, two decoded ahead, a two-slot producer queue, each batch charged
until its consumer takes it. The Lance I/O buffer is set explicitly to
64 MiB (Lance's default is 32 MiB per I/O thread, 2 GiB on a cloud store,
uncharged and read ahead of a consumer that may stop after one batch). A
limitnow stops the read after its first batches; a top-k holds its keptrows and one batch. Kept batches gather until they reach the session's
batch rows (or the rows of a
Limitdirectly over the scan's projection)or the byte target, then leave as one batch: Lance reads a batch per
fragment, and without gathering a per-batch consumer such as an indexed
hop's edge lookup would pay once per fragment. Ranked scans keep their
breaker, which the overfetch ladder needs. Only a scan a contains join
marked records the runtime-filter counters.
(
late_materialization) now also runs on query plans(
materialize_returns_late). A bare$b.propertythat nothing but theoutput reads leaves its binding's physical scan, which reads
_rowaddrinstead; a new
HydrateColumnsnode above the rootLimitfetches thosecolumns by row address for the kept rows (
HydrateExec, built on Lance'sTakeBuilder::try_new_from_addresses, chunked under an eighth of the querypool and declared as the node's
retained_limit). Lance returns eachdistinct row once; the values are then copied to every output row that
names them through the pool's checked take (
WorkMemory::take_once), whichadmits one copy per output row before Arrow builds it. A row a join repeats
is copied once per repetition, so a chunk whose copies would pass the bound
keeps a prefix, halved until they fit. The checked take now counts a null
index as a null row instead of the row its raw value names. It fires under a root
limit Kfor every binding whose scan estimates more than4Krows (themanifest count, one row under a key equality), because a scan reads more
rows than the limit keeps whenever a sort, filter or traversal sits below
it, and Lance reads ahead of a consumer that stops early even when nothing
does. Keys, sort keys, filtered, joined or aggregated columns, sorted
aliases and
Blobcolumns stay on the scan; fused ranking and aggregatesare left alone. One query reads one pinned snapshot, so the address the scan
read is the row's location in the same table version.
Merged with the typed IR (#901):
HydrateColumnsis checked as a returnoperator like
Projection. Its deferred columns take their fields from theplan's declared result schema, not the catalog. The projection under it is
checked against every declared column except the deferred ones, with its
row-address columns hidden like the
~sort columns.For that list view the scan now reads the key and the identity for every row
and the text column for the 101 rows returned, where it used to read and hold
all of it.
Not changed: the 150 MiB budget. Raising it would not fix anything, because the
demand grew with the table.
Evidence
--measure), streaming alone vs this change: the unorderedlimit 1reads3,874,592 → 94,724 bytes, the top-k 906,020 → 141,468 bytes. The fixture's
bodies are identical and compress well, so distinct text widens the gap.
cases_slow/v2/wide_column_scan_limit_answers.gqt(5,000 rowsof a 40,000-byte body, 200 MB decoded): on v0.12.0 the unordered
limit 1is refused at
actual 158237376, limit 157286400and the top-k atactual 158354916; with this change both answer, and engine v1 returns thesame rows.
a_wide_column_read_under_a_limit_follows_its_result(
engine_v2_memory.rs, 64 MB of payload under a 16 MiB pool): refused onv0.12.0 at
actual 67298352; now the scan emits fewer rows than the typeholds and
HydrateExechydrates exactly one row.cases/v2/planner/hydrate_return_columns_above_the_limit.gqt:a top-k with a null winner, a traversal whose two rows share one destination
(through a hash join) with a vector column, a limit over a traversal (through
id lookup), and three shapes that must not hydrate; every rows step is
compared with engine v1.
query_plan/cost.rs), the bound-plan round tripand the lowering walk; engine report and plan-replay tests for the new node.
hydrated_copies_of_a_joined_row_stay_within_the_chunk_bound(
engine_v2_memory.rs): sixteen narrow sources, then a hub with 512 edgesand a 40 KiB payload,
limit 108under a 16 MiB pool. Every payload comesback intact and the new
peak_chunk_bytesgauge stays under the 2 MiB chunkbound. With the prefix cut disabled, one chunk holds 3,810,196 bytes.
a_take_admits_a_copy_per_index_and_a_null_row_per_null_index(in-source,memory.rs) covers the checked take: 64 null indices whose raw value namesa 512 KiB row admit a few KiB, and four repetitions of that row are refused
under a 1 MiB pool before Arrow allocates them.
a_streamed_scan_gathers_small_fragments_into_one_batch(engine_v2.rs):twenty one-row commits reach the hop as one batch per scan; without the
gathering each scan sends twenty. The DST cost golden
(
dst_bench_cost_count_golden) caught the per-fragment cost first: withoutgathering its
_auditand_verifyreads rose from 706 and 1,349 to 734and 1,416; with it the golden matches unchanged.
a_text_contains_join_answers_where_the_unfiltered_scan_refusesnow provesthe needle-free product answers too (renamed), and
expand_streams_a_single_hop_under_a_pool_the_pairs_would_not_fitgets a20 MiB pool (still below the 24 MB of pairs it guards), because the source
scan's read-ahead batches now stream beside the hop's neighbour map.
one gathered batch; they had seen several only because Lance's reads are
sized for 4 KiB rows. Their fixtures now exceed one scan batch, so the same
assertions still prove the scan streams:
the_scan_and_the_join_share_one_matcherhas 1 KiB passages (2 MB under its 14.5 MiB pool, the matcher and a full
channel still fit), and the generated corpus of
a_contains_join_and_its_cross_join_agree_on_generated_texts_over_several_batcheshas 20,000 passages, so even its sieved rows fill more than 8,192.
Run locally before merging
main:cargo test -p omnigraph-gqt --test gq_logic_tests): 256/256;cases_slow/: 2/2; GQT lib 192/192.12, unit 21).
--features failpoints):engine_v2,engine_v2_memory,engine_v2_plan_replay,engine_v2_scrubbed_replay,traversal_indexed,traversal_adaptive,search,ordering,aggregation,literal_filters,end_to_end,scalar_indexes,rrf_prefilter_gate,forbidden_apis,lance_surface_guards,proptest_equivalence,consistency,composite_flow,point_in_time,branching,changes,export,writes, and the lib'sengine::tests: all pass.data_routesandstored_queries; CLIcli_dataandparity_matrix: all pass.cargo fmt --all --check, workspace and GQT clippy with-D warnings,scripts/check-docs.py,typos: clean.runner_dispatch::concurrent_block_refusalsfailed under heavy machine loadand passes otherwise; its
starved_wallfixture needs more than 3.0 s andless than 3.5 s of its 4 s budget on both v0.12.0 and this change, so the
flake is that budget, not this change.
After merging
mainand adding the gathering:cargo fmt --all --checkandworkspace clippy with
-D warningsare clean; every engine suite listed aboveand the lib's
engine::tests pass; the DST suite (scenarios51, which holdsthe cost golden,
lane_b,torn_init, lib 53) passes; GQT corpus 265/265,cases_slow/2/2, GQT lib,runner_dispatchandomnigraph-gqt-corepass;server
data_routesandstored_queriesand CLIcli_dataandparity_matrixpass.Follow-ups
retained_limitandwork_bytesinderive_propertiesand reshape or refuse at plan time a planwhose bound is unknown or above the pool. That changes the query contract,
so it needs an RFC.
Expandsingle hop builds its neighbour map per input batch, so its workfollows the batch's fan-out, not its output slice.