Skip to content

feat: Add Trino compute engine for batch retrieval and materialization - #6898

Open
Marcus-Rosti wants to merge 1 commit into
feast-dev:masterfrom
Marcus-Rosti:mrosti/trino-compute
Open

Marcus-Rosti wants to merge 1 commit into
feast-dev:masterfrom
Marcus-Rosti:mrosti/trino-compute

Conversation

@Marcus-Rosti

Copy link
Copy Markdown
Contributor

What this PR does / why we need it:

Adds Trino as a native Feast compute engine (type: trino.engine or type: trino), enabling DAG-based historical feature retrieval and materialization directly via Trino SQL without requiring Spark or Ray:

  • Historical Retrieval: Compiles DAG operations into an ANSI Trino query with point-in-time joins, windowed aggregations, and deduplication.
  • Materialization: Supports streaming Arrow batches into online stores and atomic staging/rename-swap into offline tables.
  • Transformations: Introduces TrinoTransformation supporting SQL templates ({}) and callable UDFs.
  • Docs & Coverage: Includes reference documentation and unit test suite (>93% coverage).

Which issue(s) this PR fixes:

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
  • Manual tests
  • Testing is not required for this change

@Marcus-Rosti
Marcus-Rosti requested a review from a team as a code owner September 29, 2026 21:27
@Marcus-Rosti

Copy link
Copy Markdown
Contributor Author

@ntkathole hopefully this helps!

@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

❌ Patch coverage is 87.51592% with 98 lines in your changes missing coverage. Please review.
✅ Project coverage is 48.65%. Comparing base (9c2da18) to head (2a0b1c8).

Files with missing lines Patch % Lines
.../python/feast/infra/compute_engines/trino/utils.py 79.72% 18 Missing and 11 partials ⚠️
.../python/feast/infra/compute_engines/trino/nodes.py 91.22% 16 Missing and 9 partials ⚠️
...ast/infra/compute_engines/trino/feature_builder.py 68.33% 18 Missing and 1 partial ⚠️
...dk/python/feast/infra/compute_engines/trino/job.py 87.03% 9 Missing and 5 partials ⚠️
...ython/feast/infra/compute_engines/trino/compute.py 94.30% 5 Missing and 2 partials ⚠️
...n/feast/infra/compute_engines/trino/sql_builder.py 88.23% 2 Missing and 2 partials ⚠️
❗ 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    #6898      +/-   ##
==========================================
+ Coverage   48.04%   48.65%   +0.60%     
==========================================
  Files         427      435       +8     
  Lines       53591    54376     +785     
  Branches     7800     7897      +97     
==========================================
+ Hits        25749    26454     +705     
- Misses      25986    26038      +52     
- Partials     1856     1884      +28     
Flag Coverage Δ
go-feature-server 30.58% <ø> (ø)
python-unit 50.01% <87.51%> (+0.62%) ⬆️
Files with missing lines Coverage Δ
...dk/python/feast/infra/compute_engines/dag/model.py 100.00% <100.00%> (ø)
...thon/feast/infra/compute_engines/trino/__init__.py 100.00% <100.00%> (ø)
sdk/python/feast/repo_config.py 79.52% <ø> (ø)
sdk/python/feast/transformation/mode.py 100.00% <100.00%> (ø)
...ython/feast/transformation/trino_transformation.py 100.00% <100.00%> (ø)
...n/feast/infra/compute_engines/trino/sql_builder.py 88.23% <88.23%> (ø)
...ython/feast/infra/compute_engines/trino/compute.py 94.30% <94.30%> (ø)
...dk/python/feast/infra/compute_engines/trino/job.py 87.03% <87.03%> (ø)
...ast/infra/compute_engines/trino/feature_builder.py 68.33% <68.33%> (ø)
.../python/feast/infra/compute_engines/trino/nodes.py 91.22% <91.22%> (ø)
... and 1 more

... and 4 files with indirect coverage changes


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 9c2da18...2a0b1c8. 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.

Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>

@ntkathole ntkathole left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this contribution, please resolve the inline comments

logger = logging.getLogger(__name__)


class TrinoFeatureBuilder(FeatureBuilder):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Blocker: entity_df is silently ignored for historical retrieval.

This class inherits the base _build() which only creates a JoinNode when there are upstream input_nodes (multi-FeatureView DAGs). For the common single-FeatureView case, the DAG is:

ReadNode -> FilterNode -> Agg/Dedup -> WriteNode — no JoinNode is ever created.

This means get_historical_features(entity_df=..., features=[...]) returns the full feature table instead of the PIT-correct subset matching the entity_df.

How Flink solves this: FlinkFeatureBuilder overrides _build() and adds:

if self._should_join_entity_df():
    last_node = self.build_join_node(view, [last_node])

How Spark solves this: SparkReadNode delegates to create_offline_store_retrieval_job() which handles entity_df at the offline store level.

Trino does neither. Please override _build() here (like Flink) to inject the entity_df join step for HistoricalRetrievalTask.

Note: the component test test_trino_compute_engine_get_historical_features hides this because mock_client.execute_query returns canned data — it never validates the generated SQL includes entity_df columns.


cte_name = f"_dedup_{self.name.replace(':', '_')}"
query = (
f"SELECT * FROM (\n"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Regression: _feast_rn column leaks into output.

This outer SELECT * includes the _feast_rn ROW_NUMBER column in the CTE output. Downstream nodes and the final result schema will contain this spurious column.

How peer engines handle this:

  • Spark SparkDedupNode: .drop("row_num") — explicitly removes the internal column
  • Flink FlinkDedupNode: Uses SELECT {explicit_column_list} (excludes DEDUP_ROW_NUMBER) + _drop_internal_columns() in the output node

Fix: Either wrap in an outer CTE that lists all columns except _feast_rn, or track the upstream column list in TrinoQueryPlan.columns and use it here.

try:
self.client.execute_query(create_sql)
# Step 2: Atomic swap
self.client.execute_query(f"DROP TABLE IF EXISTS {target_table}")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Data loss risk: non-atomic swap + destructive offline write.

Two issues here:

1. Non-atomic swap loses data if RENAME fails:
This line DROPs the target table before the RENAME. If the RENAME on line 617 fails (e.g., permission error, catalog issue), the original target is already gone. The except block then also drops the staging table — all data is permanently lost.

Fix: Reverse the order: RENAME target to backup first, then RENAME staging to target, then DROP backup.

2. Regression vs peers: offline write replaces entire table instead of appending:

  • Spark uses write.mode("append").save(path) — preserves existing data
  • Flink uses offline_store.offline_write_batch() — delegates to store's append logic
  • Trino uses DROP + CREATE TABLE AS — destroys ALL existing offline data

Every materialize() call for a feature view with offline=True wipes all previously materialized data.

Fix: Use INSERT INTO target_table (SELECT ... FROM ...) for appending, or implement a time-partitioned merge. At minimum, document this as a destructive operation.

f"{quote_identifier(ts_col)} >= {quote_identifier(ENTITY_TS_ALIAS)} - INTERVAL '{ttl_seconds}' SECOND"
)
if self.filter_condition:
conditions.append(f"({self.filter_condition})")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SQL injection surface: filter_condition is injected directly into SQL with f"({self.filter_condition})" — no sanitization or parameterization. While typically developer-provided, a malformed string will execute arbitrary Trino SQL.

At minimum, add a docstring noting that filter_condition must be trusted input.

try:
builder = TrinoFeatureBuilder(
registry=registry,
client=self.client, # type: ignore[arg-type]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

self.client can be None here (if host/catalog/user aren't configured and offline store isn't TrinoOfflineStoreConfig). The # type: ignore suppresses the type error but None will cause a confusing runtime crash inside TrinoFeatureBuilder.

Please add an explicit guard:

if self.client is None:
    raise RuntimeError(
        "Trino client is not configured. Set host, catalog, and user "
        "in batch_engine config or use a TrinoOfflineStoreConfig."
    )

@@ -0,0 +1,236 @@
from datetime import datetime, timedelta, timezone

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Missing: integration tests against a real Trino cluster.

The component and unit tests use mocked Trino clients throughout, which is good for unit-level coverage but misses:

  1. SQL compilation correctness — generated CTEs are never executed against a Trino parser/planner
  2. entity_df upload + PIT join — upload_pandas_dataframe_to_trino is mocked
  3. Streaming batch correctness — stream_trino_arrow_batches uses cursor._query (private API), never tested against the real client
  4. Offline staging swap atomicity — the DROP/RENAME sequence is never tested for failure modes
  5. Type mapping round-trips — from_feast_to_trino_type and trino_to_pa_value_type end-to-end

Other engines have integration suites under tests/integration/compute_engines/. Consider adding a Trino integration test (e.g., using Testcontainers with trinodb/trino Docker image) that validates at least:

  • Historical retrieval with entity_df produces PIT-correct results
  • Materialization writes correct data to online + offline stores
  • Streaming batch size respects the config

Also: test_trino_compute_engine_get_historical_features doesn't verify the generated SQL includes entity_df join conditions — it should assert entity columns appear in the compiled query.

return True
if isinstance(value, (list, tuple, np.ndarray)):
return False
try:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

stream_trino_arrow_batches uses private APIs: client._get_cursor() and cursor._query. These are internal to trino-python-client and can break on library upgrades.

Consider using cursor.description (public API) for column metadata instead of cursor._query.columns, and adding a version pin on the trino dependency.


cte_name = f"_pit_join_{self.name.replace(':', '_')}"
query = (
f"SELECT _entity.*, {latest_cte}.*\n"

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor: SELECT _entity.*, {latest_cte}.* will produce duplicate columns when entity_df and the feature table share join key columns (e.g., driver_id). Downstream Arrow/Pandas conversion may fail or silently rename them.

Consider qualifying the feature-side SELECT to exclude join keys already present from _entity.*.

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.

3 participants