Is your feature request related to a problem? Please describe.
Yes, in three independent places — Feast fully materializes an Arrow result as one giant,
unchunked in-memory object before anything downstream can act on it, and none of them scale:
-
OfflineServer.do_get()
wraps every result straight into fl.RecordBatchStream(table), where table comes from each
offline store's _to_arrow_internal — one giant, unchunked RecordBatch. Arrow's Flight/IPC
writer refuses to serialize any single RecordBatch above 2GiB, so a job that completed
successfully server-side crashes the moment the client tries to read the response:
pyarrow.lib.ArrowInvalid: Cannot send record batches exceeding 2GiB yet
-
TrinoRetrievalJob._to_arrow_internal()
does cursor.fetchall() → pandas DataFrame → pyarrow.Table.from_pandas() — two full
in-memory copies of the result before a single row reaches the caller, and exactly what
produces the oversized batch in (1.).
-
SparkReadNode.execute()
converts the retrieved Arrow Table to Spark via arrow_table.to_pandas() then
createDataFrame(...) — to_pandas() on a Flight/IPC-sourced, multi-chunk Table always
makes a full pandas copy (pyarrow can't zero-copy multi-chunk/nullable columns, and
self_destruct can't free the source either), doubling peak driver memory for every
materialization job and driving real driver OOMs on large batch sources.
Describe the solution you'd like
offline_server.py — before wrapping in RecordBatchStream, re-slice the outgoing table by
row count (to_batches(max_chunksize=N)) and re-check each resulting batch's actual encoded
byte size, bisecting further by row count for any batch still over a safety threshold.
- Trino offline store — an opt-in streaming retrieval path:
# trino_queries.py
query.start() # run the query, return column metadata
query.iterate_pages(batch_size) # page via cursor.fetchmany()
# trino.py
TrinoRetrievalJob.to_arrow_reader() # -> pyarrow.RecordBatchReader, built from those pages
do_get() prefers this reader over to_arrow() whenever a job carries no on-demand feature
views.
compute_engines/spark/nodes.py — on Spark 4.0+, pass the pyarrow.Table directly to
createDataFrame() (native ingestion via _create_from_arrow_table) instead of
.to_pandas(), falling back to the pandas round-trip on older PySpark.
Describe alternatives you've considered
- (1.): just raising the offline server's chunking limits — doesn't address the memory doubling
in (2.), only the point where it becomes fatal.
- (2.): making
to_arrow() itself lazy for every offline store — too large a behavior change for
callers (on-demand transforms, dataset validation) that need the full Table; scoping to an
opt-in method on Trino avoids that.
- (3.): no real alternative once Spark 4.0+ supports native ingestion — the pandas round-trip is
pure overhead.
Additional context
Test plan
Is your feature request related to a problem? Please describe.
Yes, in three independent places — Feast fully materializes an Arrow result as one giant,
unchunked in-memory object before anything downstream can act on it, and none of them scale:
OfflineServer.do_get()wraps every result straight into
fl.RecordBatchStream(table), wheretablecomes from eachoffline store's
_to_arrow_internal— one giant, unchunkedRecordBatch. Arrow's Flight/IPCwriter refuses to serialize any single
RecordBatchabove 2GiB, so a job that completedsuccessfully server-side crashes the moment the client tries to read the response:
TrinoRetrievalJob._to_arrow_internal()does
cursor.fetchall()→ pandasDataFrame→pyarrow.Table.from_pandas()— two fullin-memory copies of the result before a single row reaches the caller, and exactly what
produces the oversized batch in (1.).
SparkReadNode.execute()converts the retrieved Arrow
Tableto Spark viaarrow_table.to_pandas()thencreateDataFrame(...)—to_pandas()on a Flight/IPC-sourced, multi-chunkTablealwaysmakes a full pandas copy (pyarrow can't zero-copy multi-chunk/nullable columns, and
self_destructcan't free the source either), doubling peak driver memory for everymaterialization job and driving real driver OOMs on large batch sources.
Describe the solution you'd like
offline_server.py— before wrapping inRecordBatchStream, re-slice the outgoing table byrow count (
to_batches(max_chunksize=N)) and re-check each resulting batch's actual encodedbyte size, bisecting further by row count for any batch still over a safety threshold.
do_get()prefers this reader overto_arrow()whenever a job carries no on-demand featureviews.
compute_engines/spark/nodes.py— on Spark 4.0+, pass thepyarrow.Tabledirectly tocreateDataFrame()(native ingestion via_create_from_arrow_table) instead of.to_pandas(), falling back to the pandas round-trip on older PySpark.Describe alternatives you've considered
in (2.), only the point where it becomes fatal.
to_arrow()itself lazy for every offline store — too large a behavior change forcallers (on-demand transforms, dataset validation) that need the full
Table; scoping to anopt-in method on Trino avoids that.
pure overhead.
Additional context
Test plan
RecordBatchexceeds 2GiB serializes over Arrow Flightwithout raising
ArrowInvalid.TrinoRetrievalJob.to_arrow_reader()never callscursor.fetchall(), and its outputmatches
to_arrow()for a query withrow(...)/map(...)/array(...)columns and NULLs.SparkReadNode.execute()uses nativepyarrow.Tableingestion on Spark 4.0+, and outputrows/schema match the
.to_pandas()path for aTablewith nulls and multiple chunks.