Skip to content

Commit 9550f22

Browse files
Marcus-Rostijyejare
authored andcommitted
feat: Add quality monitoring features for Trino offline store
Signed-off-by: Marcus Rosti <rostimarcus@gmail.com>
1 parent 8ab92e8 commit 9550f22

5 files changed

Lines changed: 1267 additions & 7 deletions

File tree

docs/how-to-guides/feature-monitoring.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -359,6 +359,7 @@ Monitoring works natively with all offline stores that serve as compute engines
359359
| BigQuery | SQL push-down | `MERGE` into BQ tables |
360360
| Redshift | SQL push-down | `MERGE` via Data API |
361361
| Spark | SparkSQL push-down | Parquet tables |
362+
| Trino | SQL push-down | Trino tables |
362363
| Oracle | SQL via Ibis | `MERGE` from `DUAL` |
363364
| DuckDB | In-memory SQL | Parquet files |
364365
| Dask | PyArrow compute | Parquet files |

sdk/python/feast/infra/offline_stores/contrib/trino_offline_store/connectors/upload.py

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
```
1919
"""
2020

21-
from datetime import datetime, timezone
21+
from datetime import date, datetime, timezone
2222
from typing import Any, Dict, Iterator, Optional, Set
2323

2424
import numpy as np
@@ -117,20 +117,27 @@ def _is_nan(value: Any) -> bool:
117117
def _format_value(row: pd.Series, schema: Dict[str, Any]) -> str:
118118
formated_values = []
119119
for row_name, row_value in row.items():
120-
if schema[row_name].startswith("timestamp"):
120+
if _is_nan(row_value):
121+
formated_values.append("NULL")
122+
elif schema[row_name].startswith("timestamp"):
121123
if isinstance(row_value, datetime):
122124
row_value = format_datetime(row_value)
123125
formated_values.append(f"TIMESTAMP '{row_value}'")
126+
elif schema[row_name].startswith("date"):
127+
if isinstance(row_value, (datetime, date)):
128+
row_value = row_value.strftime("%Y-%m-%d")
129+
formated_values.append(f"DATE '{row_value}'")
130+
elif isinstance(row_value, (bool, np.bool_)):
131+
formated_values.append("TRUE" if row_value else "FALSE")
124132
elif isinstance(row_value, list):
125133
formated_values.append(f"ARRAY{row_value}")
126134
elif isinstance(row_value, np.ndarray):
127135
formated_values.append(f"ARRAY{row_value.tolist()}")
128136
elif isinstance(row_value, tuple):
129137
formated_values.append(f"ARRAY{list(row_value)}")
130138
elif isinstance(row_value, str):
131-
formated_values.append(f"'{row_value}'")
132-
elif _is_nan(row_value):
133-
formated_values.append("NULL")
139+
escaped = row_value.replace("'", "''")
140+
formated_values.append(f"'{escaped}'")
134141
else:
135142
formated_values.append(f"{row_value}")
136143

0 commit comments

Comments
 (0)