From a9cd6e6a054ddc4c665a096f962012989ba80e5e Mon Sep 17 00:00:00 2001 From: Jeff Date: Thu, 4 Nov 2021 17:32:14 -0700 Subject: [PATCH 1/7] ODFV UDFs should handle list types ODFV UDFs handling list types (e.g., embeddings/vectors) should be registered without error. Signed-off-by: Jeff --- sdk/python/feast/driver_test_data.py | 1 + sdk/python/feast/type_map.py | 4 ++ .../feature_repos/universal/feature_views.py | 37 +++++++++++++++++++ .../test_universal_odfv_feature_inference.py | 15 +++++++- 4 files changed, 56 insertions(+), 1 deletion(-) diff --git a/sdk/python/feast/driver_test_data.py b/sdk/python/feast/driver_test_data.py index 1c9a1dd20bc..b760f3e996a 100644 --- a/sdk/python/feast/driver_test_data.py +++ b/sdk/python/feast/driver_test_data.py @@ -124,6 +124,7 @@ def create_driver_hourly_stats_df(drivers, start_date, end_date) -> pd.DataFrame df_all_drivers["avg_daily_trips"] = np.random.randint(0, 1000, size=rows).astype( np.int32 ) + df_all_drivers["embedding"] = [[1.0, 0.0]] * rows df_all_drivers["created"] = pd.to_datetime(pd.Timestamp.now(tz=None).round("ms")) # Create duplicate rows that should be filtered by created timestamp diff --git a/sdk/python/feast/type_map.py b/sdk/python/feast/type_map.py index c615d2f50e9..cbf7391dc7b 100644 --- a/sdk/python/feast/type_map.py +++ b/sdk/python/feast/type_map.py @@ -89,6 +89,10 @@ def feast_value_type_to_pandas_type(value_type: ValueType) -> Any: ValueType.BYTES: "bytes", ValueType.BOOL: "bool", ValueType.UNIX_TIMESTAMP: "datetime", + ValueType.DOUBLE_LIST: "object", + ValueType.FLOAT_LIST: "object", + ValueType.INT32_LIST: "object", + ValueType.INT64_LIST: "object", } if value_type in value_type_to_pandas_type: return value_type_to_pandas_type[value_type] diff --git a/sdk/python/tests/integration/feature_repos/universal/feature_views.py b/sdk/python/tests/integration/feature_repos/universal/feature_views.py index 23ce1b4cb19..515772c8733 100644 --- a/sdk/python/tests/integration/feature_repos/universal/feature_views.py +++ b/sdk/python/tests/integration/feature_repos/universal/feature_views.py @@ -1,6 +1,7 @@ from datetime import timedelta from typing import Dict, List, Optional, Union +import numpy as np import pandas as pd from feast import Feature, FeatureView, OnDemandFeatureView, ValueType @@ -68,6 +69,35 @@ def conv_rate_plus_100_feature_view( ) +def similarity(features_df: pd.DataFrame) -> pd.DataFrame: + if features_df.size == 0: + return pd.DataFrame({"cos": [0.0]}) # give hint to Feast about return type + vectors_a = features_df["embedding"].apply(np.array) + vectors_b = features_df["vector"].apply(np.array) + dot_products = vectors_a.mul(vectors_b).apply(sum) + norms_q = vectors_a.apply(np.linalg.norm) + norms_doc = vectors_b.apply(np.linalg.norm) + df = pd.DataFrame() + df["cos"] = dot_products / (norms_q * norms_doc) + return df + + +def similarity_feature_view( + inputs: Dict[str, Union[RequestDataSource, FeatureView]], + infer_features: bool = False, + features: Optional[List[Feature]] = None, +) -> OnDemandFeatureView: + _features = features or [ + Feature("cos", ValueType.DOUBLE), + ] + return OnDemandFeatureView( + name=similarity.__name__, + inputs=inputs, + features=[] if infer_features else _features, + udf=similarity, + ) + + def create_driver_age_request_feature_view(): return RequestFeatureView( name="driver_age", @@ -83,6 +113,12 @@ def create_conv_rate_request_data_source(): ) +def create_similarity_request_data_source(): + return RequestDataSource( + name="similarity_input", schema={"vector": ValueType.DOUBLE_LIST} + ) + + def create_driver_hourly_stats_feature_view(source, infer_features: bool = False): driver_stats_feature_view = FeatureView( name="driver_stats", @@ -93,6 +129,7 @@ def create_driver_hourly_stats_feature_view(source, infer_features: bool = False Feature(name="conv_rate", dtype=ValueType.FLOAT), Feature(name="acc_rate", dtype=ValueType.FLOAT), Feature(name="avg_daily_trips", dtype=ValueType.INT32), + Feature(name="embedding", dtype=ValueType.DOUBLE_LIST), ], batch_source=source, ttl=timedelta(hours=2), diff --git a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py index 6d6750081bc..1ca47a47a58 100644 --- a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py +++ b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py @@ -7,6 +7,8 @@ conv_rate_plus_100_feature_view, create_conv_rate_request_data_source, create_driver_hourly_stats_feature_view, + create_similarity_request_data_source, + similarity_feature_view, ) @@ -27,11 +29,22 @@ def test_infer_odfv_features(environment, universal_data_sources, infer_features infer_features=infer_features, ) - feast_objects = [driver_hourly_stats, driver_odfv, driver(), customer()] + sim_odfv = similarity_feature_view( + { + "driver": driver_hourly_stats, + "input_request": create_similarity_request_data_source(), + }, + infer_features=infer_features, + ) + + feast_objects = [driver_hourly_stats, driver_odfv, sim_odfv, driver(), customer()] store.apply(feast_objects) odfv = store.get_on_demand_feature_view("conv_rate_plus_100") assert len(odfv.features) == 3 + odfv = store.get_on_demand_feature_view("similarity") + assert len(odfv.features) == 1 + @pytest.mark.integration @pytest.mark.universal From 07ec0978072aec54aa6d92d905aa409ecc68ea06 Mon Sep 17 00:00:00 2001 From: Jeff Date: Fri, 5 Nov 2021 08:21:10 -0700 Subject: [PATCH 2/7] handle all value type names that end in _LIST Signed-off-by: Jeff --- sdk/python/feast/type_map.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/sdk/python/feast/type_map.py b/sdk/python/feast/type_map.py index cbf7391dc7b..4fc3d944da8 100644 --- a/sdk/python/feast/type_map.py +++ b/sdk/python/feast/type_map.py @@ -89,11 +89,9 @@ def feast_value_type_to_pandas_type(value_type: ValueType) -> Any: ValueType.BYTES: "bytes", ValueType.BOOL: "bool", ValueType.UNIX_TIMESTAMP: "datetime", - ValueType.DOUBLE_LIST: "object", - ValueType.FLOAT_LIST: "object", - ValueType.INT32_LIST: "object", - ValueType.INT64_LIST: "object", } + if value_type.name.endswith("_LIST"): + return "object" if value_type in value_type_to_pandas_type: return value_type_to_pandas_type[value_type] raise TypeError( From befc9569df3d674e61198f7476a2998eab068ead Mon Sep 17 00:00:00 2001 From: Jeff Date: Fri, 5 Nov 2021 18:44:27 -0700 Subject: [PATCH 3/7] clearly define dummy vector for driver embedding test data Signed-off-by: Jeff --- sdk/python/feast/driver_test_data.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdk/python/feast/driver_test_data.py b/sdk/python/feast/driver_test_data.py index b760f3e996a..84095046f83 100644 --- a/sdk/python/feast/driver_test_data.py +++ b/sdk/python/feast/driver_test_data.py @@ -124,7 +124,8 @@ def create_driver_hourly_stats_df(drivers, start_date, end_date) -> pd.DataFrame df_all_drivers["avg_daily_trips"] = np.random.randint(0, 1000, size=rows).astype( np.int32 ) - df_all_drivers["embedding"] = [[1.0, 0.0]] * rows + dummy_vector = [1.0, 0.0] + df_all_drivers["embedding"] = [dummy_vector] * rows df_all_drivers["created"] = pd.to_datetime(pd.Timestamp.now(tz=None).round("ms")) # Create duplicate rows that should be filtered by created timestamp From a80a9b8e32d9068ed38215679c7bf1c5ad3d40ed Mon Sep 17 00:00:00 2001 From: Jeff Date: Fri, 5 Nov 2021 19:14:28 -0700 Subject: [PATCH 4/7] example embedding in test_write_to_online_store() Signed-off-by: Jeff --- .../tests/integration/online_store/test_universal_online.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/sdk/python/tests/integration/online_store/test_universal_online.py b/sdk/python/tests/integration/online_store/test_universal_online.py index c90021f9ce9..349cfdb9358 100644 --- a/sdk/python/tests/integration/online_store/test_universal_online.py +++ b/sdk/python/tests/integration/online_store/test_universal_online.py @@ -143,11 +143,13 @@ def test_write_to_online_store(environment, universal_data_sources): fs.apply([driver_hourly_stats, driver_entity]) # fake data to ingest into Online Store + dummy_vector = [1.0, 0.0] data = { "driver_id": [123], "conv_rate": [0.85], "acc_rate": [0.91], "avg_daily_trips": [14], + "embedding": [dummy_vector], "event_timestamp": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], "created": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], } From af3dc6b40b715da40517734573b6ecdf22a95743 Mon Sep 17 00:00:00 2001 From: Jeff Date: Fri, 5 Nov 2021 19:15:15 -0700 Subject: [PATCH 5/7] map Arrow list types to Redshift super type Signed-off-by: Jeff --- sdk/python/feast/type_map.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/sdk/python/feast/type_map.py b/sdk/python/feast/type_map.py index 4fc3d944da8..1de54800b95 100644 --- a/sdk/python/feast/type_map.py +++ b/sdk/python/feast/type_map.py @@ -453,6 +453,9 @@ def pa_to_redshift_value_type(pa_type: pyarrow.DataType) -> str: # PyArrow decimal types (e.g. "decimal(38,37)") luckily directly map to the Redshift type. return pa_type_as_str + if pa_type_as_str.startswith("list"): + return "super" + # We have to take into account how arrow types map to parquet types as well. # For example, null type maps to int32 in parquet, so we have to use int4 in Redshift. # Other mappings have also been adjusted accordingly. From 464e717e3a4fe988ce2a78cdd8538357557bafdd Mon Sep 17 00:00:00 2001 From: Jeff Date: Sat, 6 Nov 2021 00:42:26 -0700 Subject: [PATCH 6/7] ensure float list types in ODFV UDFs can be appied Signed-off-by: Jeff --- sdk/python/feast/driver_test_data.py | 3 ++- .../feature_repos/universal/feature_views.py | 24 +++++++++++++------ .../online_store/test_universal_online.py | 3 ++- .../test_universal_odfv_feature_inference.py | 2 +- 4 files changed, 22 insertions(+), 10 deletions(-) diff --git a/sdk/python/feast/driver_test_data.py b/sdk/python/feast/driver_test_data.py index 84095046f83..35c79c3cfe4 100644 --- a/sdk/python/feast/driver_test_data.py +++ b/sdk/python/feast/driver_test_data.py @@ -125,7 +125,8 @@ def create_driver_hourly_stats_df(drivers, start_date, end_date) -> pd.DataFrame np.int32 ) dummy_vector = [1.0, 0.0] - df_all_drivers["embedding"] = [dummy_vector] * rows + df_all_drivers["embedding_double"] = [dummy_vector] * rows + df_all_drivers["embedding_float"] = df_all_drivers["embedding_double"] df_all_drivers["created"] = pd.to_datetime(pd.Timestamp.now(tz=None).round("ms")) # Create duplicate rows that should be filtered by created timestamp diff --git a/sdk/python/tests/integration/feature_repos/universal/feature_views.py b/sdk/python/tests/integration/feature_repos/universal/feature_views.py index 515772c8733..f809ae2fbde 100644 --- a/sdk/python/tests/integration/feature_repos/universal/feature_views.py +++ b/sdk/python/tests/integration/feature_repos/universal/feature_views.py @@ -71,14 +71,18 @@ def conv_rate_plus_100_feature_view( def similarity(features_df: pd.DataFrame) -> pd.DataFrame: if features_df.size == 0: - return pd.DataFrame({"cos": [0.0]}) # give hint to Feast about return type - vectors_a = features_df["embedding"].apply(np.array) - vectors_b = features_df["vector"].apply(np.array) + # give hint to Feast about return type + df = pd.DataFrame({"cos_double": [0.0]}) + df["cos_float"] = df["cos_double"].astype(np.float32) + return df + vectors_a = features_df["embedding_double"].apply(np.array) + vectors_b = features_df["vector_double"].apply(np.array) dot_products = vectors_a.mul(vectors_b).apply(sum) norms_q = vectors_a.apply(np.linalg.norm) norms_doc = vectors_b.apply(np.linalg.norm) df = pd.DataFrame() - df["cos"] = dot_products / (norms_q * norms_doc) + df["cos_double"] = dot_products / (norms_q * norms_doc) + df["cos_float"] = df["cos_double"].astype(np.float32) return df @@ -88,7 +92,8 @@ def similarity_feature_view( features: Optional[List[Feature]] = None, ) -> OnDemandFeatureView: _features = features or [ - Feature("cos", ValueType.DOUBLE), + Feature("cos_double", ValueType.DOUBLE), + Feature("cos_float", ValueType.FLOAT), ] return OnDemandFeatureView( name=similarity.__name__, @@ -115,7 +120,11 @@ def create_conv_rate_request_data_source(): def create_similarity_request_data_source(): return RequestDataSource( - name="similarity_input", schema={"vector": ValueType.DOUBLE_LIST} + name="similarity_input", + schema={ + "vector_double": ValueType.DOUBLE_LIST, + "vector_float": ValueType.FLOAT_LIST, + }, ) @@ -129,7 +138,8 @@ def create_driver_hourly_stats_feature_view(source, infer_features: bool = False Feature(name="conv_rate", dtype=ValueType.FLOAT), Feature(name="acc_rate", dtype=ValueType.FLOAT), Feature(name="avg_daily_trips", dtype=ValueType.INT32), - Feature(name="embedding", dtype=ValueType.DOUBLE_LIST), + Feature(name="embedding_double", dtype=ValueType.DOUBLE_LIST), + Feature(name="embedding_float", dtype=ValueType.FLOAT_LIST), ], batch_source=source, ttl=timedelta(hours=2), diff --git a/sdk/python/tests/integration/online_store/test_universal_online.py b/sdk/python/tests/integration/online_store/test_universal_online.py index 349cfdb9358..6acf900c8b8 100644 --- a/sdk/python/tests/integration/online_store/test_universal_online.py +++ b/sdk/python/tests/integration/online_store/test_universal_online.py @@ -149,7 +149,8 @@ def test_write_to_online_store(environment, universal_data_sources): "conv_rate": [0.85], "acc_rate": [0.91], "avg_daily_trips": [14], - "embedding": [dummy_vector], + "embedding_double": [dummy_vector], + "embedding_float": [dummy_vector], "event_timestamp": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], "created": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], } diff --git a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py index 1ca47a47a58..03bf8757c94 100644 --- a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py +++ b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py @@ -43,7 +43,7 @@ def test_infer_odfv_features(environment, universal_data_sources, infer_features assert len(odfv.features) == 3 odfv = store.get_on_demand_feature_view("similarity") - assert len(odfv.features) == 1 + assert len(odfv.features) == 2 @pytest.mark.integration From c655d1be8115903f808a0d3bd9f38f6f6f20894b Mon Sep 17 00:00:00 2001 From: Jeff Date: Thu, 11 Nov 2021 21:09:08 -0800 Subject: [PATCH 7/7] isolate ODFV list type feature test to smaller code changes Signed-off-by: Jeff --- sdk/python/feast/driver_test_data.py | 3 -- .../feature_repos/universal/entities.py | 4 ++ .../feature_repos/universal/feature_views.py | 18 +++++++- .../online_store/test_universal_online.py | 3 -- .../test_universal_odfv_feature_inference.py | 45 ++++++++++++++----- 5 files changed, 55 insertions(+), 18 deletions(-) diff --git a/sdk/python/feast/driver_test_data.py b/sdk/python/feast/driver_test_data.py index 35c79c3cfe4..1c9a1dd20bc 100644 --- a/sdk/python/feast/driver_test_data.py +++ b/sdk/python/feast/driver_test_data.py @@ -124,9 +124,6 @@ def create_driver_hourly_stats_df(drivers, start_date, end_date) -> pd.DataFrame df_all_drivers["avg_daily_trips"] = np.random.randint(0, 1000, size=rows).astype( np.int32 ) - dummy_vector = [1.0, 0.0] - df_all_drivers["embedding_double"] = [dummy_vector] * rows - df_all_drivers["embedding_float"] = df_all_drivers["embedding_double"] df_all_drivers["created"] = pd.to_datetime(pd.Timestamp.now(tz=None).round("ms")) # Create duplicate rows that should be filtered by created timestamp diff --git a/sdk/python/tests/integration/feature_repos/universal/entities.py b/sdk/python/tests/integration/feature_repos/universal/entities.py index 3b4ec05f5d9..e8e90a6af62 100644 --- a/sdk/python/tests/integration/feature_repos/universal/entities.py +++ b/sdk/python/tests/integration/feature_repos/universal/entities.py @@ -16,3 +16,7 @@ def customer(): def location(): return Entity(name="location_id", value_type=ValueType.INT64) + + +def item(): + return Entity(name="item_id", value_type=ValueType.INT64) diff --git a/sdk/python/tests/integration/feature_repos/universal/feature_views.py b/sdk/python/tests/integration/feature_repos/universal/feature_views.py index f809ae2fbde..3d19212f485 100644 --- a/sdk/python/tests/integration/feature_repos/universal/feature_views.py +++ b/sdk/python/tests/integration/feature_repos/universal/feature_views.py @@ -128,6 +128,22 @@ def create_similarity_request_data_source(): ) +def create_item_embeddings_feature_view(source, infer_features: bool = False): + item_embeddings_feature_view = FeatureView( + name="item_embeddings", + entities=["item"], + features=None + if infer_features + else [ + Feature(name="embedding_double", dtype=ValueType.DOUBLE_LIST), + Feature(name="embedding_float", dtype=ValueType.FLOAT_LIST), + ], + batch_source=source, + ttl=timedelta(hours=2), + ) + return item_embeddings_feature_view + + def create_driver_hourly_stats_feature_view(source, infer_features: bool = False): driver_stats_feature_view = FeatureView( name="driver_stats", @@ -138,8 +154,6 @@ def create_driver_hourly_stats_feature_view(source, infer_features: bool = False Feature(name="conv_rate", dtype=ValueType.FLOAT), Feature(name="acc_rate", dtype=ValueType.FLOAT), Feature(name="avg_daily_trips", dtype=ValueType.INT32), - Feature(name="embedding_double", dtype=ValueType.DOUBLE_LIST), - Feature(name="embedding_float", dtype=ValueType.FLOAT_LIST), ], batch_source=source, ttl=timedelta(hours=2), diff --git a/sdk/python/tests/integration/online_store/test_universal_online.py b/sdk/python/tests/integration/online_store/test_universal_online.py index 6acf900c8b8..c90021f9ce9 100644 --- a/sdk/python/tests/integration/online_store/test_universal_online.py +++ b/sdk/python/tests/integration/online_store/test_universal_online.py @@ -143,14 +143,11 @@ def test_write_to_online_store(environment, universal_data_sources): fs.apply([driver_hourly_stats, driver_entity]) # fake data to ingest into Online Store - dummy_vector = [1.0, 0.0] data = { "driver_id": [123], "conv_rate": [0.85], "acc_rate": [0.91], "avg_daily_trips": [14], - "embedding_double": [dummy_vector], - "embedding_float": [dummy_vector], "event_timestamp": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], "created": [pd.Timestamp(datetime.datetime.utcnow()).round("ms")], } diff --git a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py index 03bf8757c94..ee3180f863e 100644 --- a/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py +++ b/sdk/python/tests/integration/registration/test_universal_odfv_feature_inference.py @@ -1,12 +1,17 @@ +from datetime import datetime + +import pandas as pd import pytest from feast import Feature, ValueType from feast.errors import SpecifiedFeaturesNotPresentError -from tests.integration.feature_repos.universal.entities import customer, driver +from feast.infra.offline_stores.file_source import FileSource +from tests.integration.feature_repos.universal.entities import customer, driver, item from tests.integration.feature_repos.universal.feature_views import ( conv_rate_plus_100_feature_view, create_conv_rate_request_data_source, create_driver_hourly_stats_feature_view, + create_item_embeddings_feature_view, create_similarity_request_data_source, similarity_feature_view, ) @@ -29,19 +34,39 @@ def test_infer_odfv_features(environment, universal_data_sources, infer_features infer_features=infer_features, ) - sim_odfv = similarity_feature_view( - { - "driver": driver_hourly_stats, - "input_request": create_similarity_request_data_source(), - }, - infer_features=infer_features, - ) - - feast_objects = [driver_hourly_stats, driver_odfv, sim_odfv, driver(), customer()] + feast_objects = [driver_hourly_stats, driver_odfv, driver(), customer()] store.apply(feast_objects) odfv = store.get_on_demand_feature_view("conv_rate_plus_100") assert len(odfv.features) == 3 + +@pytest.mark.integration +@pytest.mark.parametrize("infer_features", [True, False], ids=lambda v: str(v)) +def test_infer_odfv_list_features(environment, infer_features, tmp_path): + fake_embedding = [1.0, 1.0] + items_df = pd.DataFrame( + data={ + "item_id": [0], + "embedding_float": [fake_embedding], + "embedding_double": [fake_embedding], + "event_timestamp": [pd.Timestamp(datetime.utcnow())], + "created": [pd.Timestamp(datetime.utcnow())], + } + ) + output_path = f"{tmp_path}/items.parquet" + items_df.to_parquet(output_path) + fake_items_src = FileSource( + path=output_path, + event_timestamp_column="event_timestamp", + created_timestamp_column="created", + ) + items = create_item_embeddings_feature_view(fake_items_src) + sim_odfv = similarity_feature_view( + {"items": items, "input_request": create_similarity_request_data_source()}, + infer_features=infer_features, + ) + store = environment.feature_store + store.apply([item(), items, sim_odfv]) odfv = store.get_on_demand_feature_view("similarity") assert len(odfv.features) == 2