Skip to content

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

Description

@dbbvitor

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:

  1. 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
    
  2. 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.).

  3. 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

  • A synthetic result whose single RecordBatch exceeds 2GiB serializes over Arrow Flight
    without raising ArrowInvalid.
  • TrinoRetrievalJob.to_arrow_reader() never calls cursor.fetchall(), and its output
    matches to_arrow() for a query with row(...)/map(...)/array(...) columns and NULLs.
  • SparkReadNode.execute() uses native pyarrow.Table ingestion on Spark 4.0+, and output
    rows/schema match the .to_pandas() path for a Table with nulls and multiple chunks.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions