Repository navigation
Expand file tree
/
Copy pathoffline_store.py
More file actions
799 lines (686 loc) · 30.6 KB
/
Copy pathoffline_store.py
File metadata and controls
799 lines (686 loc) · 30.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
# Copyright 2019 The Feast Authors
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# https://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import logging
import time
import warnings
from abc import ABC
from datetime import date, datetime, timedelta, timezone
from pathlib import Path
from typing import (
TYPE_CHECKING,
Any,
Callable,
Dict,
Iterable,
List,
Optional,
Tuple,
Union,
)
import pandas as pd
import pyarrow
from feast import flags_helper
from feast.data_source import DataSource
from feast.dataframe import DataFrameEngine, FeastDataFrame
from feast.dqm.errors import ValidationFailed
from feast.feature_logging import LoggingConfig, LoggingSource
from feast.feature_view import FeatureView
from feast.infra.registry.base_registry import BaseRegistry
from feast.on_demand_feature_view import OnDemandFeatureView
from feast.repo_config import RepoConfig
from feast.saved_dataset import SavedDatasetStorage
from feast.torch_wrapper import get_torch
if TYPE_CHECKING:
from feast.saved_dataset import ValidationReference
warnings.simplefilter("once", RuntimeWarning)
class RetrievalMetadata:
min_event_timestamp: Optional[datetime]
max_event_timestamp: Optional[datetime]
# List of feature references
features: List[str]
# List of entity keys + ODFV inputs
keys: List[str]
def __init__(
self,
features: List[str],
keys: List[str],
min_event_timestamp: Optional[datetime] = None,
max_event_timestamp: Optional[datetime] = None,
):
self.features = features
self.keys = keys
self.min_event_timestamp = min_event_timestamp
self.max_event_timestamp = max_event_timestamp
def _extract_retrieval_metadata(job: "RetrievalJob") -> tuple:
"""Return ``(feature_view_names, feature_count)`` from a RetrievalJob's metadata."""
from feast.utils import _parse_feature_ref
try:
meta = job.metadata
if meta:
feature_count = len(meta.features)
feature_views = list(
{_parse_feature_ref(ref)[0] for ref in meta.features if ":" in ref}
)
return feature_views, feature_count
except (NotImplementedError, AttributeError):
pass
return [], 0
def _emit_offline_store_request_metrics(
job: "RetrievalJob",
method: str,
status_label: str,
row_count: int,
elapsed: float,
) -> None:
"""Record offline-store request metrics and the audit log entry. Never raises."""
try:
from feast import metrics as feast_metrics
if feast_metrics._config.offline_features:
feast_metrics.offline_store_request_total.labels(
method=method, status=status_label
).inc()
feast_metrics.offline_store_request_latency_seconds.labels(
method=method
).observe(elapsed)
feast_metrics.offline_store_row_count.labels(method=method).observe(
row_count
)
if feast_metrics._config.audit_logging:
feature_views, feature_count = _extract_retrieval_metadata(job)
end_dt = datetime.now(tz=timezone.utc)
start_dt = end_dt - timedelta(seconds=elapsed)
feast_metrics.emit_offline_audit_log(
method=method,
feature_views=feature_views,
feature_count=feature_count,
row_count=row_count,
status=status_label,
start_time=start_dt.isoformat(),
end_time=end_dt.isoformat(),
duration_ms=elapsed * 1000,
)
except Exception:
logging.getLogger(__name__).debug(
"Failed to record offline store metrics", exc_info=True
)
class RetrievalJob(ABC):
"""A RetrievalJob manages the execution of a query to retrieve data from the offline store."""
def to_df(
self,
validation_reference: Optional["ValidationReference"] = None,
timeout: Optional[int] = None,
) -> pd.DataFrame:
"""
Synchronously executes the underlying query and returns the result as a pandas dataframe.
On demand transformations will be executed. If a validation reference is provided, the dataframe
will be validated.
Args:
validation_reference (optional): The validation to apply against the retrieved dataframe.
timeout (optional): The query timeout if applicable.
"""
return (
self.to_arrow(validation_reference=validation_reference, timeout=timeout)
.to_pandas()
.reset_index(drop=True)
)
def to_feast_df(
self,
validation_reference: Optional["ValidationReference"] = None,
timeout: Optional[int] = None,
) -> FeastDataFrame:
"""
Synchronously executes the underlying query and returns the result as a FeastDataFrame.
This is the new primary method that returns FeastDataFrame with proper engine detection.
On demand transformations will be executed. If a validation reference is provided, the dataframe
will be validated.
Args:
validation_reference (optional): The validation to apply against the retrieved dataframe.
timeout (optional): The query timeout if applicable.
"""
# Get Arrow table as before
arrow_table = self.to_arrow(
validation_reference=validation_reference, timeout=timeout
)
# Prepare metadata
metadata = {}
# Add features to metadata if available
if hasattr(self, "features"):
metadata["features"] = self.features
else:
metadata["features"] = []
# Add on-demand feature views to metadata
if hasattr(self, "on_demand_feature_views") and self.on_demand_feature_views:
metadata["on_demand_feature_views"] = [
odfv.name for odfv in self.on_demand_feature_views
]
else:
metadata["on_demand_feature_views"] = []
# Wrap in FeastDataFrame with Arrow engine and metadata
return FeastDataFrame(
data=arrow_table, engine=DataFrameEngine.ARROW, metadata=metadata
)
def to_arrow(
self,
validation_reference: Optional["ValidationReference"] = None,
timeout: Optional[int] = None,
) -> pyarrow.Table:
"""
Synchronously executes the underlying query and returns the result as an arrow table.
On demand transformations will be executed. If a validation reference is provided, the dataframe
will be validated.
Args:
validation_reference (optional): The validation to apply against the retrieved dataframe.
timeout (optional): The query timeout if applicable.
"""
start_wall = time.monotonic()
status_label = "success"
row_count = 0
try:
features_table = self._to_arrow_internal(timeout=timeout)
row_count = features_table.num_rows
except Exception:
status_label = "error"
raise
finally:
_emit_offline_store_request_metrics(
job=self,
method="to_arrow",
status_label=status_label,
row_count=row_count,
elapsed=time.monotonic() - start_wall,
)
if self.on_demand_feature_views:
# Build a mapping of ODFV name to requested feature names
# This ensures we only return the features that were explicitly requested
odfv_feature_refs: Dict[str, set[str]] = {}
try:
metadata = self.metadata
except NotImplementedError:
metadata = None
if metadata and metadata.features:
for feature_ref in metadata.features:
if ":" in feature_ref:
view_name, feature_name = feature_ref.split(":", 1)
# Check if this view_name matches any of the ODFVs
for odfv in self.on_demand_feature_views:
if (
odfv.name == view_name
or odfv.projection.name_to_use() == view_name
):
if view_name not in odfv_feature_refs:
odfv_feature_refs[view_name] = set()
# Store the feature name in the format that will appear in transformed_arrow
expected_col_name = (
f"{odfv.projection.name_to_use()}__{feature_name}"
if self.full_feature_names
else feature_name
)
odfv_feature_refs[view_name].add(expected_col_name)
for odfv in self.on_demand_feature_views:
transformed_arrow = odfv.transform_arrow(
features_table, self.full_feature_names
)
# Determine which columns to include from this ODFV
# If we have metadata with requested features, filter to only those
# Otherwise, include all columns (backward compatibility)
requested_features_for_odfv = (
odfv_feature_refs.get(odfv.name)
if odfv.name in odfv_feature_refs
else odfv_feature_refs.get(odfv.projection.name_to_use())
)
for col in transformed_arrow.column_names:
if col.startswith("__index"):
continue
# Only append the column if it was requested, or if we don't have feature metadata
if (
requested_features_for_odfv is None
or col in requested_features_for_odfv
):
features_table = features_table.append_column(
col, transformed_arrow[col]
)
if validation_reference:
if not flags_helper.is_test():
warnings.warn(
"Dataset validation is an experimental feature. "
"This API is unstable and it could and most probably will be changed in the future. "
"We do not guarantee that future changes will maintain backward compatibility.",
RuntimeWarning,
)
validation_result = validation_reference.profile.validate(
features_table.to_pandas()
)
if not validation_result.is_success:
raise ValidationFailed(validation_result)
return features_table
def to_arrow_reader(
self, timeout: Optional[int] = None
) -> pyarrow.RecordBatchReader:
"""
Returns the result as a stream of record batches.
The default materializes the full result via ``to_arrow()``; offline stores that
can page results natively override this to bound memory.
"""
table = self.to_arrow(timeout=timeout)
return pyarrow.RecordBatchReader.from_batches(table.schema, table.to_batches())
def to_tensor(
self,
kind: str = "torch",
default_value: Any = float("nan"),
timeout: Optional[int] = None,
) -> Dict[str, Any]:
"""
Converts historical features into a dictionary of 1D torch tensors or lists (for non-numeric types).
Args:
kind: "torch" (default and only supported kind).
default_value: Value to replace missing (None or NaN) entries.
timeout: Optional timeout for query execution.
Returns:
Dict[str, Union[torch.Tensor, List]]: Feature column name -> tensor or list.
"""
if kind != "torch":
raise ValueError(
f"Unsupported tensor kind: {kind}. Only 'torch' is supported."
)
torch = get_torch()
device = "cuda" if torch.cuda.is_available() else "cpu"
df = self.to_df(timeout=timeout)
tensor_dict = {}
for column in df.columns:
values = df[column].fillna(default_value).tolist()
first_non_null = next((v for v in values if v is not None), None)
if isinstance(first_non_null, (int, float, bool)):
tensor_dict[column] = torch.tensor(values, device=device)
else:
tensor_dict[column] = values
return tensor_dict
def to_ray_dataset(self) -> Any:
"""Convert the retrieval result to a Ray Dataset.
This is a first-class method on every ``RetrievalJob`` regardless of
the configured offline store backend:
* **Ray offline store** – ``RayRetrievalJob`` overrides this method and
returns the underlying ``ray.data.Dataset`` *without* materialising the
result to the driver, keeping the full computation on the cluster.
* **All other offline stores** – the default implementation here pulls
the result to an Arrow table on the driver and converts it via
``ray.data.from_arrow()``. This is less efficient than the native
Ray path but lets every backend participate in Ray pipelines.
Example::
from feast import FeatureStore
store = FeatureStore(".")
ds = store.get_historical_features(
entity_df=entity_df,
features=["driver_stats:conv_rate"],
).to_ray_dataset()
# Chain Ray Data transforms directly
predictions = ds.map_batches(MyModel, num_gpus=1)
Returns:
ray.data.Dataset: Dataset containing the retrieved feature rows.
Raises:
ImportError: If ``ray`` is not installed. Install via
``pip install 'feast[ray]'``.
"""
try:
import ray.data
except ImportError as e:
raise ImportError(
"Ray is required to call to_ray_dataset(). "
"Install it with: pip install 'feast[ray]'"
) from e
return ray.data.from_arrow(self.to_arrow())
def to_sql(self) -> str:
"""
Return RetrievalJob generated SQL statement if applicable.
"""
raise NotImplementedError
def _to_df_internal(self, timeout: Optional[int] = None) -> pd.DataFrame:
"""
Synchronously executes the underlying query and returns the result as a pandas dataframe.
timeout: RetreivalJob implementations may implement a timeout.
Does not handle on demand transformations or dataset validation. For either of those,
`to_df` should be used.
"""
raise NotImplementedError
def _to_arrow_internal(self, timeout: Optional[int] = None) -> pyarrow.Table:
"""
Synchronously executes the underlying query and returns the result as an arrow table.
timeout: RetreivalJob implementations may implement a timeout.
Does not handle on demand transformations or dataset validation. For either of those,
`to_arrow` should be used.
"""
raise NotImplementedError
@property
def full_feature_names(self) -> bool:
"""Returns True if full feature names should be applied to the results of the query."""
raise NotImplementedError
@property
def on_demand_feature_views(self) -> List[OnDemandFeatureView]:
"""Returns a list containing all the on demand feature views to be handled."""
raise NotImplementedError
def persist(
self,
storage: SavedDatasetStorage,
allow_overwrite: bool = False,
timeout: Optional[int] = None,
):
"""
Synchronously executes the underlying query and persists the result in the same offline store
at the specified destination.
Args:
storage: The saved dataset storage object specifying where the result should be persisted.
allow_overwrite: If True, a pre-existing location (e.g. table or file) can be overwritten.
Currently not all individual offline store implementations make use of this parameter.
"""
raise NotImplementedError
@property
def metadata(self) -> Optional[RetrievalMetadata]:
"""Returns metadata about the retrieval job."""
raise NotImplementedError
def supports_remote_storage_export(self) -> bool:
"""Returns True if the RetrievalJob supports `to_remote_storage`."""
return False
def to_remote_storage(self) -> List[str]:
"""
Synchronously executes the underlying query and exports the results to remote storage (e.g. S3 or GCS).
Implementations of this method should export the results as multiple parquet files, each file sized
appropriately depending on how much data is being returned by the retrieval job.
Returns:
A list of parquet file paths in remote storage.
"""
raise NotImplementedError
class OfflineStore(ABC):
"""
An offline store defines the interface that Feast uses to interact with the storage and compute system that
handles offline features.
Each offline store implementation is designed to work only with the corresponding data source. For example,
the SnowflakeOfflineStore can handle SnowflakeSources but not FileSources.
"""
supports_filter_by_created_timestamp: bool = False
"""Whether get_historical_features supports the filter_by_created_timestamp flag."""
def ensure_filter_by_created_timestamp_supported(self) -> None:
"""Raises NotImplementedError if this store does not support filter_by_created_timestamp."""
if not self.supports_filter_by_created_timestamp:
raise NotImplementedError(
f"filter_by_created_timestamp is not supported by {type(self).__name__}"
)
@staticmethod
def pull_latest_from_table_or_query(
config: RepoConfig,
data_source: DataSource,
join_key_columns: List[str],
feature_name_columns: List[str],
timestamp_field: str,
created_timestamp_column: Optional[str],
start_date: datetime,
end_date: datetime,
) -> RetrievalJob:
"""
Extracts the latest entity rows (i.e. the combination of join key columns, feature columns, and
timestamp columns) from the specified data source that lie within the specified time range.
All of the column names should refer to columns that exist in the data source. In particular,
any mapping of column names must have already happened.
Args:
config: The config for the current feature store.
data_source: The data source from which the entity rows will be extracted.
join_key_columns: The columns of the join keys.
feature_name_columns: The columns of the features.
timestamp_field: The timestamp column, used to determine which rows are the most recent.
created_timestamp_column: The column indicating when the row was created, used to break ties.
start_date: The start of the time range.
end_date: The end of the time range.
Returns:
A RetrievalJob that can be executed to get the entity rows.
"""
raise NotImplementedError
@staticmethod
def get_historical_features(
config: RepoConfig,
feature_views: List[FeatureView],
feature_refs: List[str],
entity_df: Optional[Union[pd.DataFrame, str]],
registry: BaseRegistry,
project: str,
full_feature_names: bool = False,
) -> RetrievalJob:
"""
Retrieves the point-in-time correct historical feature values for the specified entity rows.
Args:
config: The config for the current feature store.
feature_views: A list containing all feature views that are referenced in the entity rows.
feature_refs: The features to be retrieved.
entity_df: A collection of rows containing all entity columns on which features need to be joined,
as well as the timestamp column used for point-in-time joins. Either a pandas dataframe can be
provided or a SQL query. If None, features will be retrieved for the specified timestamp range.
registry: The registry for the current feature store.
project: Feast project to which the feature views belong.
full_feature_names: If True, feature names will be prefixed with the corresponding feature view name,
changing them from the format "feature" to "feature_view__feature" (e.g. "daily_transactions"
changes to "customer_fv__daily_transactions").
Keyword Args:
start_date: Start date for the timestamp range when retrieving features without entity_df.
end_date: End date for the timestamp range when retrieving features without entity_df. By default, the current time is used.
filter_by_created_timestamp: If True, only feature values whose created timestamp (the
``created_timestamp_column`` of the batch source) is at or before the entity row's event
timestamp are considered. Only passed through when a store declares
``supports_filter_by_created_timestamp``.
Returns:
A RetrievalJob that can be executed to get the features.
"""
raise NotImplementedError
@staticmethod
def pull_all_from_table_or_query(
config: RepoConfig,
data_source: DataSource,
join_key_columns: List[str],
feature_name_columns: List[str],
timestamp_field: str,
created_timestamp_column: Optional[str] = None,
start_date: Optional[datetime] = None,
end_date: Optional[datetime] = None,
) -> RetrievalJob:
"""
Extracts all the entity rows (i.e. the combination of join key columns, feature columns, and
timestamp columns) from the specified data source that lie within the specified time range.
All of the column names should refer to columns that exist in the data source. In particular,
any mapping of column names must have already happened.
Args:
config: The config for the current feature store.
data_source: The data source from which the entity rows will be extracted.
join_key_columns: The columns of the join keys.
feature_name_columns: The columns of the features.
timestamp_field: The timestamp column, used to determine which rows are the most recent.
created_timestamp_column (Optional): The column indicating when the row was created, used to break ties.
start_date (Optional): The start of the time range.
end_date (Optional): The end of the time range.
Returns:
A RetrievalJob that can be executed to get the entity rows.
"""
raise NotImplementedError
@staticmethod
def write_logged_features(
config: RepoConfig,
data: Union[pyarrow.Table, Path],
source: LoggingSource,
logging_config: LoggingConfig,
registry: BaseRegistry,
):
"""
Writes logged features to a specified destination in the offline store.
If the specified destination exists, data will be appended; otherwise, the destination will be
created and data will be added. Thus this function can be called repeatedly with the same
destination to flush logs in chunks.
Args:
config: The config for the current feature store.
data: An arrow table or a path to parquet directory that contains the logs to write.
source: The logging source that provides a schema and some additional metadata.
logging_config: A LoggingConfig object that determines where the logs will be written.
registry: The registry for the current feature store.
"""
raise NotImplementedError
@staticmethod
def offline_write_batch(
config: RepoConfig,
feature_view: FeatureView,
table: pyarrow.Table,
progress: Optional[Callable[[int], Any]],
):
"""
Writes the specified arrow table to the data source underlying the specified feature view.
Args:
config: The config for the current feature store.
feature_view: The feature view whose batch source should be written.
table: The arrow table to write.
progress: Function to be called once a portion of the data has been written, used
to show progress.
"""
raise NotImplementedError
def validate_data_source(
self,
config: RepoConfig,
data_source: DataSource,
):
"""
Validates the underlying data source.
Args:
config: Configuration object used to configure a feature store.
data_source: DataSource object that needs to be validated
"""
data_source.validate(config=config)
def get_table_column_names_and_types_from_data_source(
self,
config: RepoConfig,
data_source: DataSource,
) -> Iterable[Tuple[str, str]]:
"""
Returns the list of column names and raw column types for a DataSource.
Args:
config: Configuration object used to configure a feature store.
data_source: DataSource object
"""
return data_source.get_table_column_names_and_types(config=config)
@staticmethod
def compute_monitoring_metrics(
config: RepoConfig,
data_source: DataSource,
feature_columns: List[Tuple[str, str]],
timestamp_field: str,
start_date: Optional[datetime] = None,
end_date: Optional[datetime] = None,
histogram_bins: int = 20,
top_n: int = 10,
) -> List[Dict[str, Any]]:
"""
Compute monitoring metrics (stats, percentiles, histograms) directly
in the offline store using its native compute engine.
Backends that don't support this should leave it unimplemented;
the monitoring service will fall back to Python-based computation.
Args:
config: The config for the current feature store.
data_source: The data source to compute metrics from.
feature_columns: List of (feature_name, feature_type) where
feature_type is "numeric" or "categorical".
timestamp_field: Column used for time-range filtering.
start_date: Start of the time range.
end_date: End of the time range.
histogram_bins: Number of bins for numeric histograms.
top_n: Number of top values for categorical histograms.
Returns:
A list of metric dicts, one per feature, matching the format
produced by MetricsCalculator.compute_all.
"""
raise NotImplementedError
@staticmethod
def get_monitoring_max_timestamp(
config: RepoConfig,
data_source: DataSource,
timestamp_field: str,
) -> Optional[datetime]:
"""
Return the maximum event timestamp from the data source.
Used by the monitoring service to determine date ranges for
auto-compute. Backends that don't support this should leave it
unimplemented; the caller will fall back to a full-table scan.
Args:
config: The config for the current feature store.
data_source: The data source to query.
timestamp_field: The timestamp column name.
Returns:
The maximum timestamp, or None if no data exists.
"""
raise NotImplementedError
# ------------------------------------------------------------------ #
# Monitoring metrics storage (native)
# ------------------------------------------------------------------ #
MONITORING_VALID_GRANULARITIES = (
"daily",
"weekly",
"biweekly",
"monthly",
"quarterly",
)
@staticmethod
def ensure_monitoring_tables(config: RepoConfig) -> None:
"""Create the monitoring metrics tables if they do not exist.
Backends that don't support native monitoring storage should
leave this unimplemented; the monitoring service will raise an
error indicating the backend lacks storage support.
"""
raise NotImplementedError
@staticmethod
def save_monitoring_metrics(
config: RepoConfig,
metric_type: str,
metrics: List[Dict[str, Any]],
) -> None:
"""Persist monitoring metrics (upsert semantics).
Args:
config: The config for the current feature store.
metric_type: One of "feature", "feature_view", "feature_service".
metrics: List of metric dicts to upsert.
"""
raise NotImplementedError
@staticmethod
def query_monitoring_metrics(
config: RepoConfig,
project: str,
metric_type: str,
filters: Optional[Dict[str, Any]] = None,
start_date: Optional[date] = None,
end_date: Optional[date] = None,
) -> List[Dict[str, Any]]:
"""Read monitoring metrics with optional filtering.
Args:
config: The config for the current feature store.
project: Feast project name.
metric_type: One of "feature", "feature_view", "feature_service".
filters: Column-value pairs for WHERE clauses.
start_date: Inclusive lower bound on metric_date.
end_date: Inclusive upper bound on metric_date.
Returns:
List of metric dicts ordered by metric_date ascending.
"""
raise NotImplementedError
@staticmethod
def clear_monitoring_baseline(
config: RepoConfig,
project: str,
feature_view_name: Optional[str] = None,
feature_name: Optional[str] = None,
data_source_type: Optional[str] = None,
) -> None:
"""Set is_baseline=FALSE for matching feature metric rows.
Used to ensure only one baseline exists per feature before
writing a new baseline.
"""
raise NotImplementedError