From 304e307c8f75953835ca8da74cb8b82497377761 Mon Sep 17 00:00:00 2001 From: Jacob Klegar Date: Tue, 23 Mar 2021 18:05:24 -0400 Subject: [PATCH 1/3] Add materialize_incremental Signed-off-by: Jacob Klegar --- protos/feast/core/FeatureView.proto | 8 ++ sdk/python/feast/feature_store.py | 111 +++++++++++++----- sdk/python/feast/feature_view.py | 25 +++- ..._materialize_from_bigquery_to_datastore.py | 35 ++++-- 4 files changed, 138 insertions(+), 41 deletions(-) diff --git a/protos/feast/core/FeatureView.proto b/protos/feast/core/FeatureView.proto index 64c319034f9..d98f54825a6 100644 --- a/protos/feast/core/FeatureView.proto +++ b/protos/feast/core/FeatureView.proto @@ -71,4 +71,12 @@ message FeatureViewMeta { // Time where this Feature View is last updated google.protobuf.Timestamp last_updated_timestamp = 2; + + // List of pairs (start_time, end_time) for which this feature view has been materialized. + repeated MaterializationInterval materialization_intervals = 3; +} + +message MaterializationInterval { + google.protobuf.Timestamp start_time = 1; + google.protobuf.Timestamp end_time = 2; } diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index bff787dfaa9..6006190ad77 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -181,6 +181,53 @@ def get_historical_features( ) return job + def materialize_incremental( + self, + feature_views: Optional[List[str]], + end_date: datetime, + ) -> None: + """ + Materialize incremental new data from the offline store into the online store. + + This method loads incremental new feature data up to the specified end time from either + the specified feature views, or all feature views if none are specified, + into the online store where it is available for online serving. The start time of + the interval materialized is either the most recent end time of a prior materialization or + (now - ttl) if no such prior materialization exists. + + Args: + feature_views (List[str]): Optional list of feature view names. If selected, will only run + materialization for the specified feature views. + end_date (datetime): End date for time range of data to materialize into the online store + + Examples: + Materialize all features into the online store up to 5 minutes ago. + >>> from datetime import datetime, timedelta + >>> from feast.feature_store import FeatureStore + >>> + >>> fs = FeatureStore(config=RepoConfig(provider="gcp")) + >>> fs.materialize_incremental( + >>> end_date=datetime.utcnow() - timedelta(minutes=5) + >>> ) + """ + feature_views_to_materialize = [] + registry = self._get_registry() + if feature_views is None: + feature_views_to_materialize = registry.list_feature_views( + self.config.project + ) + else: + for name in feature_views: + feature_view = registry.get_feature_view(name, self.config.project) + feature_views_to_materialize.append(feature_view) + + # TODO paging large loads + for feature_view in feature_views_to_materialize: + start_date = feature_view.most_recent_end_time + if start_date is None: + start_date = datetime.utcnow() - feature_view.ttl + self._materialize_single_feature_view(feature_view, start_date, end_date) + def materialize( self, feature_views: Optional[List[str]], @@ -225,39 +272,45 @@ def materialize( # TODO paging large loads for feature_view in feature_views_to_materialize: - if isinstance(feature_view.input, FileSource): - raise NotImplementedError( - "This function is not yet implemented for File data sources" - ) - ( - entity_names, - feature_names, - event_timestamp_column, - created_timestamp_column, - ) = _run_reverse_field_mapping(feature_view) - - offline_store = get_offline_store(self.config) - table = offline_store.pull_latest_from_table_or_query( - feature_view.input, - entity_names, - feature_names, - event_timestamp_column, - created_timestamp_column, - start_date, - end_date, + self._materialize_single_feature_view(feature_view, start_date, end_date) + + def _materialize_single_feature_view(self, feature_view: FeatureView, start_date: datetime, end_date: datetime) -> None: + if isinstance(feature_view.input, FileSource): + raise NotImplementedError( + "This function is not yet implemented for File data sources" ) + ( + entity_names, + feature_names, + event_timestamp_column, + created_timestamp_column, + ) = _run_reverse_field_mapping(feature_view) + + offline_store = get_offline_store(self.config) + table = offline_store.pull_latest_from_table_or_query( + feature_view.input, + entity_names, + feature_names, + event_timestamp_column, + created_timestamp_column, + start_date, + end_date, + ) - if feature_view.input.field_mapping is not None: - table = _run_forward_field_mapping( - table, feature_view.input.field_mapping - ) + if feature_view.input.field_mapping is not None: + table = _run_forward_field_mapping( + table, feature_view.input.field_mapping + ) - rows_to_write = _convert_arrow_to_proto(table, feature_view) + rows_to_write = _convert_arrow_to_proto(table, feature_view) - provider = self._get_provider() - provider.online_write_batch( - self.config.project, feature_view, rows_to_write - ) + provider = self._get_provider() + provider.online_write_batch( + self.config.project, feature_view, rows_to_write + ) + + feature_view.materialization_intervals.append((start_date, end_date)) + self.apply([feature_view]) def get_online_features( self, feature_refs: List[str], entity_rows: List[Dict[str, Any]], diff --git a/sdk/python/feast/feature_view.py b/sdk/python/feast/feature_view.py index 4bcfff2db9e..19712e8acc4 100644 --- a/sdk/python/feast/feature_view.py +++ b/sdk/python/feast/feature_view.py @@ -11,20 +11,19 @@ # 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. -from datetime import timedelta -from typing import Dict, List, Optional, Union +from datetime import datetime, timedelta +from typing import Dict, List, Optional, Tuple, Union from google.protobuf.duration_pb2 import Duration from google.protobuf.timestamp_pb2 import Timestamp from feast.data_source import BigQuerySource, DataSource, FileSource from feast.feature import Feature -from feast.protos.feast.core.FeatureView_pb2 import FeatureView as FeatureViewProto from feast.protos.feast.core.FeatureView_pb2 import ( + FeatureView as FeatureViewProto, FeatureViewMeta as FeatureViewMetaProto, -) -from feast.protos.feast.core.FeatureView_pb2 import ( FeatureViewSpec as FeatureViewSpecProto, + MaterializationInterval as MaterializationIntervalProto, ) from feast.value_type import ValueType @@ -44,6 +43,7 @@ class FeatureView: created_timestamp: Optional[Timestamp] = None last_updated_timestamp: Optional[Timestamp] = None + materialization_intervals: List[Tuple[datetime, datetime]] = [] def __init__( self, @@ -98,7 +98,13 @@ def to_proto(self) -> FeatureViewProto: meta = FeatureViewMetaProto( created_timestamp=self.created_timestamp, last_updated_timestamp=self.last_updated_timestamp, + materialization_intervals=[], ) + for interval in self.materialization_intervals: + interval_proto = MaterializationIntervalProto() + interval_proto.start_time.FromDatetime(interval[0]) + interval_proto.end_time.FromDatetime(interval[1]) + meta.materialization_intervals.append(interval_proto) if self.ttl is not None: ttl_duration = Duration() @@ -152,4 +158,13 @@ def from_proto(cls, feature_view_proto: FeatureViewProto): feature_view.created_timestamp = feature_view_proto.meta.created_timestamp + for interval in feature_view_proto.meta.materialization_intervals: + feature_view.materialization_intervals.append((interval.start_time.ToDatetime(), interval.end_time.ToDatetime())) + return feature_view + + @property + def most_recent_end_time(self) -> Optional[datetime]: + if len(self.materialization_intervals) == 0: + return None + return max([interval[1] for interval in self.materialization_intervals]) diff --git a/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py b/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py index 3a6788130bc..1e140feebcf 100644 --- a/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py +++ b/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py @@ -28,12 +28,13 @@ def setup_method(self): def test_bigquery_table_to_datastore_correctness(self): # create dataset - ts = pd.Timestamp.now(tz="UTC").round("ms") + now = datetime.utcnow() + ts = pd.Timestamp(now).round("ms") data = { - "id": [1, 2, 1], - "value": [0.1, 0.2, 0.3], - "ts_1": [ts - timedelta(minutes=2), ts, ts], - "created_ts": [ts, ts, ts], + "id": [1, 2, 1, 3, 3], + "value": [0.1, 0.2, 0.3, 4, 5], + "ts_1": [ts - timedelta(minutes=4), ts, ts - timedelta(minutes=3), ts - timedelta(minutes=4), ts - timedelta(minutes=1)], + "created_ts": [ts, ts, ts, ts, ts], } df = pd.DataFrame.from_dict(data) @@ -68,8 +69,8 @@ def test_bigquery_table_to_datastore_correctness(self): # run materialize() fs.materialize( [fv.name], - datetime.utcnow() - timedelta(minutes=5), - datetime.utcnow() - timedelta(minutes=0), + now - timedelta(minutes=5), + now - timedelta(minutes=2), ) # check result of materialize() @@ -78,6 +79,26 @@ def test_bigquery_table_to_datastore_correctness(self): ).to_dict() assert abs(response_dict[f"{fv.name}:value"][0] - 0.3) < 1e-6 + # check prior value for materialize_incremental() + entity_key = EntityKeyProto( + entity_names=["driver_id"], entity_values=[ValueProto(int64_val=3)] + ) + t, val = fs._get_provider().online_read("default", fv, entity_key) + assert abs(val["value"].double_val - 4) < 1e-6 + + # run materialize_incremental() + fs.materialize_incremental( + [fv.name], + now - timedelta(minutes=0), + ) + + # check result of materialize_incremental() + entity_key = EntityKeyProto( + entity_names=["driver_id"], entity_values=[ValueProto(int64_val=3)] + ) + t, val = fs._get_provider().online_read("default", fv, entity_key) + assert abs(val["value"].double_val - 5) < 1e-6 + def test_bigquery_query_to_datastore_correctness(self): # create dataset ts = pd.Timestamp.now(tz="UTC").round("ms") From d6c2ec89c313f4ae5e1a2359dd02e43ef94248d2 Mon Sep 17 00:00:00 2001 From: Jacob Klegar Date: Tue, 23 Mar 2021 19:21:25 -0400 Subject: [PATCH 2/3] Rebase Signed-off-by: Jacob Klegar --- sdk/python/feast/feature_store.py | 20 +++++------ sdk/python/feast/feature_view.py | 10 ++++-- ..._materialize_from_bigquery_to_datastore.py | 33 ++++++++++--------- 3 files changed, 35 insertions(+), 28 deletions(-) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 6006190ad77..8cafe2eec5d 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -182,9 +182,7 @@ def get_historical_features( return job def materialize_incremental( - self, - feature_views: Optional[List[str]], - end_date: datetime, + self, feature_views: Optional[List[str]], end_date: datetime, ) -> None: """ Materialize incremental new data from the offline store into the online store. @@ -225,6 +223,10 @@ def materialize_incremental( for feature_view in feature_views_to_materialize: start_date = feature_view.most_recent_end_time if start_date is None: + if feature_view.ttl is None: + raise Exception( + f"No start time found for feature view {feature_view.name}. materialize_incremental() requires either a ttl to be set or for materialize() to have been run." + ) start_date = datetime.utcnow() - feature_view.ttl self._materialize_single_feature_view(feature_view, start_date, end_date) @@ -274,7 +276,9 @@ def materialize( for feature_view in feature_views_to_materialize: self._materialize_single_feature_view(feature_view, start_date, end_date) - def _materialize_single_feature_view(self, feature_view: FeatureView, start_date: datetime, end_date: datetime) -> None: + def _materialize_single_feature_view( + self, feature_view: FeatureView, start_date: datetime, end_date: datetime + ) -> None: if isinstance(feature_view.input, FileSource): raise NotImplementedError( "This function is not yet implemented for File data sources" @@ -298,16 +302,12 @@ def _materialize_single_feature_view(self, feature_view: FeatureView, start_date ) if feature_view.input.field_mapping is not None: - table = _run_forward_field_mapping( - table, feature_view.input.field_mapping - ) + table = _run_forward_field_mapping(table, feature_view.input.field_mapping) rows_to_write = _convert_arrow_to_proto(table, feature_view) provider = self._get_provider() - provider.online_write_batch( - self.config.project, feature_view, rows_to_write - ) + provider.online_write_batch(self.config.project, feature_view, rows_to_write) feature_view.materialization_intervals.append((start_date, end_date)) self.apply([feature_view]) diff --git a/sdk/python/feast/feature_view.py b/sdk/python/feast/feature_view.py index 19712e8acc4..d31394ad6e0 100644 --- a/sdk/python/feast/feature_view.py +++ b/sdk/python/feast/feature_view.py @@ -19,10 +19,14 @@ from feast.data_source import BigQuerySource, DataSource, FileSource from feast.feature import Feature +from feast.protos.feast.core.FeatureView_pb2 import FeatureView as FeatureViewProto from feast.protos.feast.core.FeatureView_pb2 import ( - FeatureView as FeatureViewProto, FeatureViewMeta as FeatureViewMetaProto, +) +from feast.protos.feast.core.FeatureView_pb2 import ( FeatureViewSpec as FeatureViewSpecProto, +) +from feast.protos.feast.core.FeatureView_pb2 import ( MaterializationInterval as MaterializationIntervalProto, ) from feast.value_type import ValueType @@ -159,7 +163,9 @@ def from_proto(cls, feature_view_proto: FeatureViewProto): feature_view.created_timestamp = feature_view_proto.meta.created_timestamp for interval in feature_view_proto.meta.materialization_intervals: - feature_view.materialization_intervals.append((interval.start_time.ToDatetime(), interval.end_time.ToDatetime())) + feature_view.materialization_intervals.append( + (interval.start_time.ToDatetime(), interval.end_time.ToDatetime()) + ) return feature_view diff --git a/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py b/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py index 1e140feebcf..f0095a5adda 100644 --- a/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py +++ b/sdk/python/tests/test_materialize_from_bigquery_to_datastore.py @@ -33,7 +33,13 @@ def test_bigquery_table_to_datastore_correctness(self): data = { "id": [1, 2, 1, 3, 3], "value": [0.1, 0.2, 0.3, 4, 5], - "ts_1": [ts - timedelta(minutes=4), ts, ts - timedelta(minutes=3), ts - timedelta(minutes=4), ts - timedelta(minutes=1)], + "ts_1": [ + ts - timedelta(seconds=4), + ts, + ts - timedelta(seconds=3), + ts - timedelta(seconds=4), + ts - timedelta(seconds=1), + ], "created_ts": [ts, ts, ts, ts, ts], } df = pd.DataFrame.from_dict(data) @@ -68,9 +74,7 @@ def test_bigquery_table_to_datastore_correctness(self): # run materialize() fs.materialize( - [fv.name], - now - timedelta(minutes=5), - now - timedelta(minutes=2), + [fv.name], now - timedelta(seconds=5), now - timedelta(seconds=2), ) # check result of materialize() @@ -80,24 +84,21 @@ def test_bigquery_table_to_datastore_correctness(self): assert abs(response_dict[f"{fv.name}:value"][0] - 0.3) < 1e-6 # check prior value for materialize_incremental() - entity_key = EntityKeyProto( - entity_names=["driver_id"], entity_values=[ValueProto(int64_val=3)] - ) - t, val = fs._get_provider().online_read("default", fv, entity_key) - assert abs(val["value"].double_val - 4) < 1e-6 + response_dict = fs.get_online_features( + [f"{fv.name}:value"], [{"driver_id": 3}] + ).to_dict() + assert abs(response_dict[f"{fv.name}:value"][0] - 4) < 1e-6 # run materialize_incremental() fs.materialize_incremental( - [fv.name], - now - timedelta(minutes=0), + [fv.name], now - timedelta(seconds=0), ) # check result of materialize_incremental() - entity_key = EntityKeyProto( - entity_names=["driver_id"], entity_values=[ValueProto(int64_val=3)] - ) - t, val = fs._get_provider().online_read("default", fv, entity_key) - assert abs(val["value"].double_val - 5) < 1e-6 + response_dict = fs.get_online_features( + [f"{fv.name}:value"], [{"driver_id": 3}] + ).to_dict() + assert abs(response_dict[f"{fv.name}:value"][0] - 5) < 1e-6 def test_bigquery_query_to_datastore_correctness(self): # create dataset From 93c1677ae4169d46e236c7d673c99669325b4ba5 Mon Sep 17 00:00:00 2001 From: Jacob Klegar Date: Thu, 25 Mar 2021 11:33:27 -0400 Subject: [PATCH 3/3] Address comment Signed-off-by: Jacob Klegar --- sdk/python/feast/feature_store.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 8cafe2eec5d..76850ef7609 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -225,7 +225,7 @@ def materialize_incremental( if start_date is None: if feature_view.ttl is None: raise Exception( - f"No start time found for feature view {feature_view.name}. materialize_incremental() requires either a ttl to be set or for materialize() to have been run." + f"No start time found for feature view {feature_view.name}. materialize_incremental() requires either a ttl to be set or for materialize() to have been run at least once." ) start_date = datetime.utcnow() - feature_view.ttl self._materialize_single_feature_view(feature_view, start_date, end_date)