Skip to content

feat: Stream Trino offline results via to_arrow_reader - #6897

Open
dbbvitor wants to merge 12 commits into
feast-dev:masterfrom
dbbvitor:feat/trino-streaming-arrow-reader
Open

dbbvitor wants to merge 12 commits into
feast-dev:masterfrom
dbbvitor:feat/trino-streaming-arrow-reader

Conversation

@dbbvitor

Copy link
Copy Markdown
Contributor

What this PR does / why we need it:

Adds a default RetrievalJob.to_arrow_reader() (wraps to_arrow(), so every existing offline store keeps working unchanged), a Trino-specific override that pages results via cursor.fetchmany() instead of fetchall() (new streaming_batch_size config, default 200000), and a shared _emit_offline_store_request_metrics() helper extracted from to_arrow()'s metrics/audit code so the new method doesn't duplicate it. OfflineServer._execute_read_api and do_exchange's historical-features branch now call .to_arrow_reader() instead of .to_arrow(), so Trino queries stream through the offline server instead of materializing the full result before PR A's _bounded_reader re-chunks it
for Flight.

Design note: the fork patch this replaces used hasattr(job, "to_arrow_reader") duck-typing in the server; this PR uses a real base-class method instead, so the on-demand-feature-view fallback logic lives in the Trino override (where it belongs) rather than in the server.

Also fixes trino_to_pa_value_type() to map ... with time zone to pa.timestamp("us", tz="UTC") instead of a
tz-naive timestamp -- needed because the new streaming path trusts the declared type-map schema (unlike the old materialized path, which got tz-awareness for free from pandas' per-row inference); without the fix, streaming would have silently dropped timezone info that the existing to_arrow() path preserves. See Misc for the full old-vs-new schema
comparison.

Retires fork patches patch_trino_queries_streaming.py, patch_trino_streaming_arrow_reader.py, patch_offline_store_metrics_helper.py.

Which issue(s) this PR fixes:

Closes #6863 (merge after #6896 )

Checks

  • I've made sure the tests are passing.
  • My commits are signed off (git commit -s)
  • My PR title follows conventional commits format

Testing Strategy

  • Unit tests

  • Integration tests

  • 275 unit tests across test_metrics.py, test_offline_store.py, test_trino_streaming.py (new), test_trino_type_map.py, test_offline_server.py.

  • 100% diff coverage (line + branch) against upstream/master (113 changed lines, 0 missing -- independently re-verified).

  • Mutation-tested with cosmic-ray in 3 scoped sessions:

    Module Mutants Killed Survivors Timeouts
    offline_store.py 39 39 0 0
    Trino contrib package 115 115 0 0
    offline_server.py 35 35 0 0

    (Trino package independently re-verified at the same 115/115. offline_server.py's 35 pending mutants all sit in fix: Bound Arrow Flight batch size in the offline server #6896 own lines, this PR's own hunks there generated no mutants, since method-name swaps and a type annotation aren't mutable by cosmic-ray's operators.)

  • Also ran the real integration test (tests/integration/offline_server/test_offline_server.py --integration), a genuine Flight client/server round trip through the new reader path, all 4 tests pass.

The offline server's do_get and do_exchange handlers each returned a
single unbounded Arrow Table/RecordBatch to the Flight client. Once a
result approaches the IPC layer's 2GiB-per-batch limit, Flight raises
"Cannot send record batches exceeding 2GiB yet" and the request fails
outright instead of degrading gracefully.

Add module-level _split_oversized/_bounded_reader helpers that
re-batch any Table or RecordBatchReader so no single RecordBatch
exceeds 1.5GB, and route both do_get's RecordBatchStream and
do_exchange's writer.write_batch loop through them. do_exchange is
the newer HPA-safe read path (bae07fc) and was still writing one
unbounded table via write_table, so it needed the same fix as do_get.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
cosmic-ray's mutation run over the diff found one surviving mutant:
`half = batch.num_rows // 2` mutated to `/`. A real pyarrow
RecordBatch.slice() silently truncates a float offset/length to the
same int, so no test built on real batches can observe the
difference. Add a duck-typed fake batch that records the exact
slice() call arguments and asserts they stay int, which does
distinguish the two operators.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
Extract to_arrow()'s metrics/audit finally-block into a reusable
_emit_offline_store_request_metrics helper, and give every RetrievalJob a
default to_arrow_reader() that wraps to_arrow(). This keeps every existing
offline store working unchanged while giving stores that can page results
natively (Trino, next) a hook to override without duplicating metrics code.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
Add Query.start()/iterate_pages() to page through cursor.fetchmany()
instead of fetchall(), and override TrinoRetrievalJob.to_arrow_reader()
to consume them lazily. Jobs carrying on-demand feature views still use
the materialized to_arrow() path, since ODFVs need the full table.

Also fix trino_to_pa_value_type() to map "... with time zone" to a
tz-aware pa.timestamp("us", tz="UTC") instead of a tz-naive one, since
the cursor returns tz-aware datetimes for those columns.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
Switch _execute_read_api's three read APIs, and do_exchange's
get_historical_features branch, from to_arrow() to to_arrow_reader().
The default to_arrow_reader() makes this a no-op for every non-Trino
store; Trino now streams instead of materializing the full result
before the existing _bounded_reader (PR A) re-batches it for Flight.

SQL errors still surface before the stream starts, since Query.start()
executes eagerly. Cursor/temp-table cleanup on an early client
disconnect is best-effort: RecordBatchReader.close() does not close
the underlying generator.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
Cosmic-ray found real gaps: the elapsed-time subtraction in both the
metrics helper and to_arrow_reader() had no test pinning its exact
value, so add/mul/mod/floordiv mutants survived; _stringify_complex's
depth==0 check and its depth-1 recursion had no negative-depth or
multi-level test to distinguish them from <=0 or >>1/%1/^1; and
_complex_column_depth's trailing slice needed three levels of array
nesting before an off-by-one became observable.

Also drop the dead preserve_index=False kwarg from
pyarrow.Table.from_pandas() in the Trino reader: since we always pass
an explicit schema with no index field, pyarrow ignores preserve_index
entirely for the table's data and public schema (verified empirically),
so the argument was an untestable no-op rather than a missing test case.

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
Update the RetrievalJob functionality matrix in overview.md to match
the per-store trino.md page: Trino now streams via to_arrow_reader().

Signed-off-by: dbbvitor <vitor.diniz@gympass.com>
@dbbvitor
dbbvitor requested review from a team as code owners September 29, 2026 20:53
@dbbvitor
dbbvitor requested review from ejscribner, robhowley and tokoko and removed request for a team September 29, 2026 20:53
@codecov-commenter

codecov-commenter commented Sep 29, 2026 •

Copy link
Copy Markdown

⚠️ Please install the 'codecov app svg image' to ensure uploads and comments are reliably processed by Codecov.

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 48.70%. Comparing base (fe27230) to head (f299854).
⚠️ Report is 3 commits behind head on master.
❗ Your organization needs to install the Codecov GitHub app to enable full functionality.

Additional details and impacted files

Impacted file tree graph

@@            Coverage Diff             @@
##           master    #6897      +/-   ##
==========================================
+ Coverage   48.50%   48.70%   +0.19%     
==========================================
  Files         427      427              
  Lines       53755    53845      +90     
  Branches     7827     7840      +13     
==========================================
+ Hits        26076    26227     +151     
+ Misses      25813    25752      -61     
  Partials     1866     1866              
Flag Coverage Δ
go-feature-server 30.58% <ø> (ø)
python-unit 50.08% <100.00%> (+0.21%) ⬆️
Files with missing lines Coverage Δ
...ffline_stores/contrib/trino_offline_store/trino.py 44.55% <100.00%> (+5.99%) ⬆️
...tores/contrib/trino_offline_store/trino_queries.py 90.35% <100.00%> (+45.18%) ⬆️
...ores/contrib/trino_offline_store/trino_type_map.py 97.53% <100.00%> (+0.06%) ⬆️
...python/feast/infra/offline_stores/offline_store.py 82.58% <100.00%> (+0.40%) ⬆️
sdk/python/feast/offline_server.py 42.33% <100.00%> (+6.58%) ⬆️

Continue to review full report in Codecov by Harness.

Legend - Click here to learn more
Δ = absolute <relative> (impact), ø = not affected, ? = missing data
Powered by Codecov. Last update fe27230...f299854. Read the comment docs.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Large Arrow results are fully materialized before use in three places (Flight, Trino, Spark)

3 participants