| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
|
@ntkathole hopefully this helps! |
Sorry, something went wrong.
|
⚠️ Please install the Codecov Report❌ Patch coverage is 86.88699% with 123 lines in your changes missing coverage. Please review.
@@ 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
... and 44 files with indirect coverage changes Continue to review full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
Sorry, something went wrong.
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
There was a problem hiding this comment.
Thanks for this contribution, please resolve the inline comments
Sorry, something went wrong.
| 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.
Sorry, something went wrong.
|
|
||
| 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:
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.
Sorry, something went wrong.
| 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:
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.
Sorry, something went wrong.
| 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.
Sorry, something went wrong.
| 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."
)
Sorry, something went wrong.
| @@ -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:
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:
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.
Sorry, something went wrong.
| 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.
Sorry, something went wrong.
|
|
||
| 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.*.
Sorry, something went wrong.
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. |
There was a problem hiding this comment.
@ntkathole I implemented 3 methods here, append, merge & overwrite
Sorry, something went wrong.
|
@ntkathole ready for review! |
Sorry, something went wrong.
| Back | FazBrowse Home | New Git URL |
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:
Which issue(s) this PR fixes:
Checks
Testing Strategy