From 0c8e6f2de36b154c1305d5baf0b6c07f09bdc1d9 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Thu, 8 Sep 2022 09:44:28 -0700 Subject: [PATCH 1/3] Clean up entity key serialization Signed-off-by: Felix Wang --- .../adding-support-for-a-new-online-store.md | 10 ++++++++-- .../cassandra_online_store/cassandra_online_store.py | 4 ++-- sdk/python/feast/infra/online_stores/snowflake.py | 8 ++++---- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md b/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md index 52f0897138d..35ad98ed82a 100644 --- a/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md +++ b/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md @@ -154,7 +154,10 @@ def online_write_batch( project = config.project for entity_key, values, timestamp, created_ts in data: - entity_key_bin = serialize_entity_key(entity_key).hex() + entity_key_bin = serialize_entity_key( + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, + ).hex() timestamp = _to_naive_utc(timestamp) if created_ts is not None: created_ts = _to_naive_utc(created_ts) @@ -184,7 +187,10 @@ def online_read( project = config.project for entity_key in entity_keys: - entity_key_bin = serialize_entity_key(entity_key).hex() + entity_key_bin = serialize_entity_key( + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, + ).hex() print(f"entity_key_bin: {entity_key_bin}") cur.execute( diff --git a/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py b/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py index ee0cb19fef5..5aedbfa40d6 100644 --- a/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py +++ b/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py @@ -314,7 +314,7 @@ def online_write_batch( project = config.project for entity_key, values, timestamp, created_ts in data: entity_key_bin = serialize_entity_key( - entity_key, entity_key_serialization_version=2 + entity_key, entity_key_serialization_version=config.entity_key_serialization_version, ).hex() with tracing_span(name="remote_call"): self._write_rows( @@ -353,7 +353,7 @@ def online_read( for entity_key in entity_keys: entity_key_bin = serialize_entity_key( - entity_key, entity_key_serialization_version=2 + entity_key, entity_key_serialization_version=config.entity_key_serialization_version, ).hex() with tracing_span(name="remote_call"): diff --git a/sdk/python/feast/infra/online_stores/snowflake.py b/sdk/python/feast/infra/online_stores/snowflake.py index a52beb73f76..7bd2092109f 100644 --- a/sdk/python/feast/infra/online_stores/snowflake.py +++ b/sdk/python/feast/infra/online_stores/snowflake.py @@ -97,9 +97,9 @@ def online_write_batch( for j, (feature_name, val) in enumerate(values.items()): df.loc[j, "entity_feature_key"] = serialize_entity_key( - entity_key, 2 + entity_key, entity_key_serialization_version=config.entity_key_serialization_version ) + bytes(feature_name, encoding="utf-8") - df.loc[j, "entity_key"] = serialize_entity_key(entity_key, 2) + df.loc[j, "entity_key"] = serialize_entity_key(entity_key, entity_key_serialization_version=config.entity_key_serialization_version) df.loc[j, "feature_name"] = feature_name df.loc[j, "value"] = val.SerializeToString() df.loc[j, "event_ts"] = timestamp @@ -165,7 +165,7 @@ def online_read( ( "TO_BINARY(" + hexlify( - serialize_entity_key(combo[0], 2) + serialize_entity_key(combo[0], entity_key_serialization_version=config.entity_key_serialization_version) + bytes(combo[1], encoding="utf-8") ).__str__()[1:] + ")" @@ -187,7 +187,7 @@ def online_read( df = execute_snowflake_statement(conn, query).fetch_pandas_all() for entity_key in entity_keys: - entity_key_bin = serialize_entity_key(entity_key, 2) + entity_key_bin = serialize_entity_key(entity_key, entity_key_serialization_version=config.entity_key_serialization_version) res = {} res_ts = None for index, row in df[df["entity_key"] == entity_key_bin].iterrows(): From d55d0a31336333d86e34777fd84a5f1dd34c8822 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Thu, 8 Sep 2022 09:49:49 -0700 Subject: [PATCH 2/3] Remove incorrect docs Signed-off-by: Felix Wang --- .../adding-support-for-a-new-online-store.md | 12 ------------ 1 file changed, 12 deletions(-) diff --git a/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md b/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md index 35ad98ed82a..ab88ebaa203 100644 --- a/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md +++ b/docs/how-to-guides/customizing-feast/adding-support-for-a-new-online-store.md @@ -214,18 +214,6 @@ def online_read( ``` {% endcode %} -### 1.3 Type Mapping - -Most online stores will have to perform some custom mapping of online store datatypes to feast value types. - -* The function to implement here are `source_datatype_to_feast_value_type` and `get_column_names_and_types` in your `DataSource` class. -* `source_datatype_to_feast_value_type` is used to convert your DataSource's datatypes to feast value types. -* `get_column_names_and_types` retrieves the column names and corresponding datasource types. - -Add any helper functions for type conversion to `sdk/python/feast/type_map.py`. - -* Be sure to implement correct type mapping so that Feast can process your feature columns without casting incorrectly that can potentially cause loss of information or incorrect data. - ## 2. Defining an OnlineStoreConfig class Additional configuration may be needed to allow the OnlineStore to talk to the backing store. For example, MySQL may need configuration information like the host at which the MySQL instance is running, credentials for connecting to the database, etc. From 5f8d5e2c54164d97e66651407d642314a4341d87 Mon Sep 17 00:00:00 2001 From: Felix Wang Date: Thu, 8 Sep 2022 10:48:51 -0700 Subject: [PATCH 3/3] Format Signed-off-by: Felix Wang --- .../cassandra_online_store.py | 6 ++++-- .../feast/infra/online_stores/snowflake.py | 18 ++++++++++++++---- 2 files changed, 18 insertions(+), 6 deletions(-) diff --git a/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py b/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py index 5aedbfa40d6..f89517c41eb 100644 --- a/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py +++ b/sdk/python/feast/infra/online_stores/contrib/cassandra_online_store/cassandra_online_store.py @@ -314,7 +314,8 @@ def online_write_batch( project = config.project for entity_key, values, timestamp, created_ts in data: entity_key_bin = serialize_entity_key( - entity_key, entity_key_serialization_version=config.entity_key_serialization_version, + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, ).hex() with tracing_span(name="remote_call"): self._write_rows( @@ -353,7 +354,8 @@ def online_read( for entity_key in entity_keys: entity_key_bin = serialize_entity_key( - entity_key, entity_key_serialization_version=config.entity_key_serialization_version, + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, ).hex() with tracing_span(name="remote_call"): diff --git a/sdk/python/feast/infra/online_stores/snowflake.py b/sdk/python/feast/infra/online_stores/snowflake.py index 7bd2092109f..9eadcb0d407 100644 --- a/sdk/python/feast/infra/online_stores/snowflake.py +++ b/sdk/python/feast/infra/online_stores/snowflake.py @@ -97,9 +97,13 @@ def online_write_batch( for j, (feature_name, val) in enumerate(values.items()): df.loc[j, "entity_feature_key"] = serialize_entity_key( - entity_key, entity_key_serialization_version=config.entity_key_serialization_version + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, ) + bytes(feature_name, encoding="utf-8") - df.loc[j, "entity_key"] = serialize_entity_key(entity_key, entity_key_serialization_version=config.entity_key_serialization_version) + df.loc[j, "entity_key"] = serialize_entity_key( + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, + ) df.loc[j, "feature_name"] = feature_name df.loc[j, "value"] = val.SerializeToString() df.loc[j, "event_ts"] = timestamp @@ -165,7 +169,10 @@ def online_read( ( "TO_BINARY(" + hexlify( - serialize_entity_key(combo[0], entity_key_serialization_version=config.entity_key_serialization_version) + serialize_entity_key( + combo[0], + entity_key_serialization_version=config.entity_key_serialization_version, + ) + bytes(combo[1], encoding="utf-8") ).__str__()[1:] + ")" @@ -187,7 +194,10 @@ def online_read( df = execute_snowflake_statement(conn, query).fetch_pandas_all() for entity_key in entity_keys: - entity_key_bin = serialize_entity_key(entity_key, entity_key_serialization_version=config.entity_key_serialization_version) + entity_key_bin = serialize_entity_key( + entity_key, + entity_key_serialization_version=config.entity_key_serialization_version, + ) res = {} res_ts = None for index, row in df[df["entity_key"] == entity_key_bin].iterrows():