From 394ba1e81145f05e6d6ffb1a8ccf0fd1db9f57f8 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Mon, 16 May 2022 14:52:17 -0700 Subject: [PATCH 1/4] fix: Feature Logging test & python server ports Signed-off-by: Oleksii Moskalenko --- .../integration/e2e/test_go_feature_server.py | 23 ++++++++++++++----- .../feature_repos/repo_configuration.py | 2 +- 2 files changed, 18 insertions(+), 7 deletions(-) diff --git a/sdk/python/tests/integration/e2e/test_go_feature_server.py b/sdk/python/tests/integration/e2e/test_go_feature_server.py index e469c90c11f..61588b6d7a0 100644 --- a/sdk/python/tests/integration/e2e/test_go_feature_server.py +++ b/sdk/python/tests/integration/e2e/test_go_feature_server.py @@ -10,7 +10,7 @@ import pytest import pytz -from feast import FeatureService, ValueType +from feast import FeatureService, FeatureView, ValueType from feast.embedded_go.lib.embedded import LoggingOptions from feast.embedded_go.online_features_service import EmbeddedOnlineFeatureServer from feast.feast_object import FeastObject @@ -162,13 +162,14 @@ def test_feature_logging( _, datasets, _ = universal_data_sources latest_rows = get_latest_rows(datasets.driver_df, "driver_id", driver_ids) + feature_view = fs.get_feature_view("driver_stats") features = [ feature.name for proj in feature_service.feature_view_projections for feature in proj.features ] expected_logs = generate_expected_logs( - latest_rows, "driver_stats", features, ["driver_id"], "event_timestamp" + latest_rows, feature_view, features, ["driver_id"], "event_timestamp" ) def retrieve(): @@ -213,15 +214,25 @@ def get_latest_rows(df, join_key, entity_values): def generate_expected_logs( - df, feature_view_name, features, join_keys, timestamp_column + df: pd.DataFrame, + feature_view: FeatureView, + features: List[str], + join_keys: List[str], + timestamp_column: str, ): logs = pd.DataFrame() for join_key in join_keys: logs[join_key] = df[join_key] for feature in features: - logs[f"{feature_view_name}__{feature}"] = df[feature] - logs[f"{feature_view_name}__{feature}__timestamp"] = df[timestamp_column] - logs[f"{feature_view_name}__{feature}__status"] = FieldStatus.PRESENT + col = f"{feature_view.name}__{feature}" + logs[col] = df[feature] + logs[f"{col}__timestamp"] = df[timestamp_column] + logs[f"{col}__status"] = FieldStatus.PRESENT + logs[f"{col}__status"] = logs[f"{col}__status"].mask( + df[timestamp_column] + < datetime.utcnow().replace(tzinfo=pytz.UTC) - feature_view.ttl, + FieldStatus.OUTSIDE_MAX_AGE, + ) return logs.sort_values(by=join_keys).reset_index(drop=True) diff --git a/sdk/python/tests/integration/feature_repos/repo_configuration.py b/sdk/python/tests/integration/feature_repos/repo_configuration.py index 27cf1a52e9d..ce469777c89 100644 --- a/sdk/python/tests/integration/feature_repos/repo_configuration.py +++ b/sdk/python/tests/integration/feature_repos/repo_configuration.py @@ -349,7 +349,7 @@ def get_local_server_port(self) -> int: worker_id_num = int(parsed_worker_id[0]) else: worker_id_num = 0 - return 6000 + 100 * worker_id_num + self.id + return 16000 + 100 * worker_id_num + self.id def table_name_from_data_source(ds: DataSource) -> Optional[str]: From ac4f573e7522dfca48f49610cc708f06eeaaa8e0 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Mon, 16 May 2022 15:02:02 -0700 Subject: [PATCH 2/4] optional timedelta Signed-off-by: Oleksii Moskalenko --- .../tests/integration/e2e/test_go_feature_server.py | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/sdk/python/tests/integration/e2e/test_go_feature_server.py b/sdk/python/tests/integration/e2e/test_go_feature_server.py index 61588b6d7a0..4e4cfc1fb8f 100644 --- a/sdk/python/tests/integration/e2e/test_go_feature_server.py +++ b/sdk/python/tests/integration/e2e/test_go_feature_server.py @@ -229,10 +229,11 @@ def generate_expected_logs( logs[col] = df[feature] logs[f"{col}__timestamp"] = df[timestamp_column] logs[f"{col}__status"] = FieldStatus.PRESENT - logs[f"{col}__status"] = logs[f"{col}__status"].mask( - df[timestamp_column] - < datetime.utcnow().replace(tzinfo=pytz.UTC) - feature_view.ttl, - FieldStatus.OUTSIDE_MAX_AGE, - ) + if feature_view.ttl: + logs[f"{col}__status"] = logs[f"{col}__status"].mask( + df[timestamp_column] + < datetime.utcnow().replace(tzinfo=pytz.UTC) - feature_view.ttl, + FieldStatus.OUTSIDE_MAX_AGE, + ) return logs.sort_values(by=join_keys).reset_index(drop=True) From fd8bba98cb2f8dd9533d683ebc70c9b7b29a6999 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Mon, 16 May 2022 15:28:32 -0700 Subject: [PATCH 3/4] revert Signed-off-by: Oleksii Moskalenko --- .../tests/integration/feature_repos/repo_configuration.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/python/tests/integration/feature_repos/repo_configuration.py b/sdk/python/tests/integration/feature_repos/repo_configuration.py index ce469777c89..c82063e464f 100644 --- a/sdk/python/tests/integration/feature_repos/repo_configuration.py +++ b/sdk/python/tests/integration/feature_repos/repo_configuration.py @@ -349,7 +349,7 @@ def get_local_server_port(self) -> int: worker_id_num = int(parsed_worker_id[0]) else: worker_id_num = 0 - return 16000 + 100 * worker_id_num + self.id + return g6000 + 100 * worker_id_num + self.id def table_name_from_data_source(ds: DataSource) -> Optional[str]: From 090a3db366ed95be949064d9b7ad5c80d763d316 Mon Sep 17 00:00:00 2001 From: Oleksii Moskalenko Date: Mon, 16 May 2022 15:29:44 -0700 Subject: [PATCH 4/4] typo Signed-off-by: Oleksii Moskalenko --- .../tests/integration/feature_repos/repo_configuration.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/python/tests/integration/feature_repos/repo_configuration.py b/sdk/python/tests/integration/feature_repos/repo_configuration.py index c82063e464f..27cf1a52e9d 100644 --- a/sdk/python/tests/integration/feature_repos/repo_configuration.py +++ b/sdk/python/tests/integration/feature_repos/repo_configuration.py @@ -349,7 +349,7 @@ def get_local_server_port(self) -> int: worker_id_num = int(parsed_worker_id[0]) else: worker_id_num = 0 - return g6000 + 100 * worker_id_num + self.id + return 6000 + 100 * worker_id_num + self.id def table_name_from_data_source(ds: DataSource) -> Optional[str]: