Conversation
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
requested review from
ejscribner,
robhowley and
tokoko
and removed request for
a team
September 29, 2026 20:53
|
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ 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
Continue to review full report in Codecov by Harness.
🚀 New features to boost your workflow:
|
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What this PR does / why we need it:
Adds a default
RetrievalJob.to_arrow_reader()(wrapsto_arrow(), so every existing offline store keeps working unchanged), a Trino-specific override that pages results viacursor.fetchmany()instead offetchall()(newstreaming_batch_sizeconfig, default 200000), and a shared_emit_offline_store_request_metrics()helper extracted fromto_arrow()'s metrics/audit code so the new method doesn't duplicate it.OfflineServer._execute_read_apianddo_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_readerre-chunks itfor 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 zonetopa.timestamp("us", tz="UTC")instead of atz-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 schemacomparison.
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
git commit -s)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:
offline_store.pyoffline_server.py(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.