Skip to content

fix(engine): bound hydration's Lance fetch before it allocates - #935

Open
ragnorc wants to merge 2 commits into
ModernRelay:mainfrom
ragnorc:fix/hydrate-fetch-bound
Open

ragnorc wants to merge 2 commits into
ModernRelay:mainfrom
ragnorc:fix/hydrate-fetch-bound

Conversation

@ragnorc

@ragnorc ragnorc commented Oct 10, 2026

Copy link
Copy Markdown
Contributor

Why

HydrateExec (from #887) fetches the return-only columns of the rows a limit kept. It planned each fetch window from the widest row seen so far, then took the window's distinct rows with one TakeBuilder::execute(). That call allocates the whole result, and Lance reads ahead up to its default I/O buffer (32 MiB per I/O thread, 2 GiB on a cloud store), all before the query pool sees a byte. When rows widen after a narrow run, a window planned from the narrow rows overshoots the chunk bound by the whole window. The new test below reproduced this: a top-k of 150 under a 16 MiB pool held 3,278,868 bytes in one fetch against a 2 MiB bound. The #887 review thread named this as the remaining limit.

What changes

The fetch is bounded by design, not by the widths already seen.

  • Row-by-row decode, charged as it goes. Hydration reads through FileFragment::open and FragmentReader::take(offsets, 1, ..), one row per decode task. The reader merges overlays and applies the deletion vector, as TakeBuilder does. Each decoded row is copied at its own size through the pool's checked take, then charged before the fetch waits for the next row. The copy matters because a Lance row can slice a buffer of the whole read, and the pool charges a buffer by its allocation's capacity. At most 8 rows decode at once.
  • Stop and replan. When the decoded rows pass half of the chunk's share, the fetch stops and drops the read. The window is then planned again from the widest row decoded.
  • Read-ahead inside the bound. An engine ScanScheduler caps Lance's read-ahead at a quarter of hydrate_chunk_bytes. That quarter is reserved in the pool for the operator's life as hydrate read-ahead.
  • The bound. The operator holds at most hydrate_chunk_bytes (the declared retained_limit) plus the row that crosses it. Lance treats its I/O buffer as a soft limit, so it can read one head request past it.
  • Checked interleave. Copies to output rows go through a new checked WorkMemory::interleave_once, which admits one copy per pick before Arrow builds it. take_bytes is gone; hydration was its only caller.
  • Fallback for inherited files. Lance 11 ignores a passed scheduler for a data file reached through a base_id, opening its own max-bandwidth scheduler instead (open_current_file_reader). A Lance branch reaches its parent's files this way. In OmniGraph only table forks created before v11 hold such files; graph branches stage detached commits in the table's own base. A fragment with such a file is read one row per request and counted in a new single_row_reads operator metric.
  • lance-io becomes a regular engine dependency. lance already links it.

Evidence

  • engine_v2_memory.rs::a_hydrated_fetch_stops_at_the_chunk_bound_when_rows_widen: red on main (one fetch held 3,278,868 bytes against a 2,097,152-byte bound), green here. The narrow rows sit in their own fragment. In a shared one, Lance's shared-buffer measurement makes the narrow seed look wide, which hides the overshoot.

  • lance_surface_guards.rs::fragment_reads_bypass_the_passed_scheduler_for_inherited_files pins the Lance behaviour behind the fallback. It goes red at the Lance release that honours the passed scheduler for every base, and then the fallback can be removed.

  • In-source operators::hydrate::tests::fragment_rows_reads_inherited_and_own_fragments_by_address exercises both read shapes against a real Lance branch.

  • engine_v2_memory.rs::hydration_on_a_graph_branch_reads_in_coalesced_requests: on a graph branch, single_row_reads stays 0.

  • The memory.rs unit test now also covers the checked interleave: per-pick admission, and refusal before Arrow allocates.

  • The GQT case hydrate_return_columns_above_the_limit.gqt gains a step that hydrates on a branch.

  • Request counts are unchanged. On the slow wide-column case, run on the in-memory test store with --measure: limit 1 takes 202 store requests and the top-101 takes 204, on main and on this branch alike.

  • Local, before merging main:

    • fmt, workspace clippy with -D warnings, check-docs.py and typos are clean.
    • All 62 Lance guards, the engine suites, the 461 lib tests and the DST suite (51 scenarios, including the cost golden) pass.
    • The planner and GQT-core suites pass, as do both slow cases.
    • The full GQT corpus didn't run clean on the shared machine (load average 30–300). Every failure was a 10 s timeout in a case unrelated to hydration, with a different set each run. Each re-ran alone in 2–6 s.
  • After merging main (52beecfe):

    • fmt, clippy, check-docs.py and typos are clean.
    • All 64 Lance guards, forbidden_apis, engine_v2, engine_v2_memory, engine_v2_plan_replay, end_to_end, search, branching and the 481 lib tests pass.
    • The GQT hydration case and both slow cases pass.
    • The full GQT corpus at default parallelism passed 299 of 308; all 9 failures were timeouts, and each passes alone in about 5.5 s.
    • runner_dispatch::concurrent_block_refusals failed under load. Its wall-clock budget is known to be load-sensitive.

Follow-ups

  • Declared residency for every query node (plan-time bounds) still needs an RFC.

HydrateExec planned each window from the widest row seen so far and took
the window's distinct rows with one `TakeBuilder` call, which allocates the
whole result, and lets Lance read ahead up to its default I/O buffer (32 MiB
per I/O thread), before the pool sees a byte. Rows that widen after a narrow
run overshot the chunk bound by the whole window: a top-k of 150 under a
16 MiB pool held 3.28 MB in one fetch against a 2 MiB bound.

The fetch now reads through `FragmentReader::take` with one row per decode
task. Each decoded row is copied at its own size through the pool's checked
take (a Lance row can slice a buffer of the whole read) and charged before
the fetch waits for the next. Once the decoded rows pass half of the chunk's
share, the fetch stops and the window is planned again from the widest row.
An engine `ScanScheduler` caps Lance's read-ahead at a quarter of the bound,
reserved in the pool, so the operator holds at most `hydrate_chunk_bytes`
and the row that crosses it, with at most eight rows decoded and not yet
charged. Copies to output rows go through a new checked
`WorkMemory::interleave_once`.

Lance 11 ignores a passed scheduler for a data file reached through a
`base_id` (a table fork created before v11); such a fragment is read one row
per request and counted as `single_row_reads`. A Lance guard pins that
behaviour so the fallback goes when Lance honours the scheduler. Graph
branches stage writes in the table's own base and keep coalesced reads.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant