FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

feat: Add Trino compute engine for batch retrieval and materialization by Marcus-Rosti · Pull Request #6898 · feast-dev/feast · GitHub

Repository navigation

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

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

Marcus-Rosti wants to merge 10 commits into
feast-dev:masterfrom
Marcus-Rosti:mrosti/trino-compute

Conversation

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 requested a review from a team as a code owner September 29, 2026 21:27

Copy link
Copy Markdown
Contributor Author

@ntkathole hopefully this helps!

codecov-commenter commented Sep 29, 2026 •
edited
Loading

Copy link
Copy Markdown

⚠️ Please install the to ensure uploads and comments are reliably processed by Codecov.

Codecov Report

❌ Patch coverage is 86.88699% with 123 lines in your changes missing coverage. Please review.
✅ Project coverage is 49.74%. Comparing base (9c2da18) to head (fa6d4e2).
⚠️ Report is 29 commits behind head on master.

Files with missing lines Patch % Lines
.../python/feast/infra/compute_engines/trino/nodes.py 91.12% 19 Missing and 15 partials ⚠️
...ast/infra/compute_engines/trino/feature_builder.py 60.24% 26 Missing and 7 partials ⚠️
.../python/feast/infra/compute_engines/trino/utils.py 81.13% 18 Missing and 12 partials ⚠️
...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.57% 5 Missing and 2 partials ⚠️
...n/feast/infra/compute_engines/trino/sql_builder.py 90.24% 2 Missing and 2 partials ⚠️
...tores/contrib/trino_offline_store/trino_queries.py 66.66% 1 Missing ⚠️
❗ Your organization needs to install the Codecov GitHub app to enable full functionality.
Additional details and impacted files

@@            Coverage Diff             @@
##           master    #6898      +/-   ##
==========================================
+ Coverage   48.04%   49.74%   +1.70%     
==========================================
  Files         427      441      +14     
  Lines       53591    55267    +1676     
  Branches     7800     8047     +247     
==========================================
+ Hits        25749    27495    +1746     
+ Misses      25986    25852     -134     
- Partials     1856     1920      +64     
Flag Coverage Δ
go-feature-server 30.58% <ø> (ø)
python-unit 51.16% <86.88%> (+1.78%) ⬆️
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%> (ø)
...tores/contrib/trino_offline_store/trino_queries.py 55.78% <66.66%> (+10.62%) ⬆️
...n/feast/infra/compute_engines/trino/sql_builder.py 90.24% <90.24%> (ø)
...ython/feast/infra/compute_engines/trino/compute.py 94.57% <94.57%> (ø)
...dk/python/feast/infra/compute_engines/trino/job.py 87.03% <87.03%> (ø)
.../python/feast/infra/compute_engines/trino/utils.py 81.13% <81.13%> (ø)
... and 2 more

... and 44 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...fa6d4e2. 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 left a comment •
edited
Loading

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

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

Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
def _execute_offline_write(
self, plan: TrinoQueryPlan, context: ExecutionContext
) -> None:
"""Execute materialization to an offline table.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

@ntkathole I implemented 3 methods here, append, merge & overwrite

Copy link
Copy Markdown
Contributor Author

@ntkathole ready for review!

This branch has not been deployed

No deployments
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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants


Back | FazBrowse Home | New Git URL