feat: Add Trino compute engine for batch retrieval and materialization - #6898
Marcus-Rosti wants to merge 1 commit into
Conversation
|
@ntkathole hopefully this helps! |
|
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ 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
... and 4 files with indirect coverage changes Continue to review full report in Codecov by Harness.
🚀 New features to boost your workflow:
|
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
6854b37 to
2a0b1c8
Compare
| logger = logging.getLogger(__name__) | ||
|
|
||
|
|
||
| class TrinoFeatureBuilder(FeatureBuilder): |
There was a problem hiding this comment.
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" |
There was a problem hiding this comment.
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: UsesSELECT {explicit_column_list}(excludesDEDUP_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}") |
There was a problem hiding this comment.
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})") |
There was a problem hiding this comment.
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] |
There was a problem hiding this comment.
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 | |||
There was a problem hiding this comment.
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:
- SQL compilation correctness — generated CTEs are never executed against a Trino parser/planner
- entity_df upload + PIT join —
upload_pandas_dataframe_to_trinois mocked - Streaming batch correctness —
stream_trino_arrow_batchesusescursor._query(private API), never tested against the real client - Offline staging swap atomicity — the DROP/RENAME sequence is never tested for failure modes
- Type mapping round-trips —
from_feast_to_trino_typeandtrino_to_pa_value_typeend-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: |
There was a problem hiding this comment.
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" |
There was a problem hiding this comment.
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.*.
What this PR does / why we need it:
Adds Trino as a native Feast compute engine (
type: trino.engineortype: trino), enabling DAG-based historical feature retrieval and materialization directly via Trino SQL without requiring Spark or Ray:TrinoTransformationsupporting SQL templates ({}) and callable UDFs.Which issue(s) this PR fixes:
Checks
git commit -s)Testing Strategy