| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent de67bdd commit edf25af
9 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -142,7 +142,7 @@ class VersionedOnlineReadNotSupported(FeastError): | |||
| 142 | 142 | def __init__(self, store_name: str, version: int): | |
| 143 | 143 | super().__init__( | |
| 144 | 144 | f"Versioned feature reads (@v{version}) are not yet supported by {store_name}. " | |
| 145 | - f"Currently only SQLite, PostgreSQL, MySQL, and FAISS support version-qualified feature references. " | ||
| 145 | + f"Currently only SQLite, PostgreSQL, MySQL, FAISS, Redis, and DynamoDB support version-qualified feature references. " | ||
| 146 | 146 | ) | |
| 147 | 147 | ||
| 148 | 148 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -24,7 +24,7 @@ | |||
| 24 | 24 | from pydantic import StrictBool, StrictStr | |
| 25 | 25 | ||
| 26 | 26 | from feast import Entity, FeatureView, utils | |
| 27 | - from feast.infra.online_stores.helpers import compute_entity_id | ||
| 27 | + from feast.infra.online_stores.helpers import compute_entity_id, compute_versioned_name | ||
| 28 | 28 | from feast.infra.online_stores.online_store import OnlineStore | |
| 29 | 29 | from feast.infra.supported_async_methods import SupportedAsyncMethods | |
| 30 | 30 | from feast.infra.utils.aws_utils import dynamo_write_items_async | |
@@ -1154,8 +1154,11 @@ def _initialize_dynamodb_resource( | |||
| 1154 | 1154 | def _get_table_name( | |
| 1155 | 1155 | online_config: DynamoDBOnlineStoreConfig, config: RepoConfig, table: FeatureView | |
| 1156 | 1156 | ) -> str: | |
| 1157 | + table_name = compute_versioned_name( | ||
| 1158 | + table, config.registry.enable_online_feature_view_versioning | ||
| 1159 | + ) | ||
| 1157 | 1160 | return online_config.table_name_template.format( | |
| 1158 | - project=config.project, table_name=table.name | ||
| 1161 | + project=config.project, table_name=table_name | ||
| 1159 | 1162 | ) | |
| 1160 | 1163 | ||
| 1161 | 1164 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -72,13 +72,18 @@ def _to_naive_utc(ts: datetime) -> datetime: | |||
| 72 | 72 | return ts.astimezone(tz=timezone.utc).replace(tzinfo=None) | |
| 73 | 73 | ||
| 74 | 74 | ||
| 75 | - def compute_table_id(project: str, table: Any, enable_versioning: bool = False) -> str: | ||
| 76 | - """Build the online-store table name, appending a version suffix when versioning is enabled.""" | ||
| 75 | + def compute_versioned_name(table: Any, enable_versioning: bool = False) -> str: | ||
| 76 | + """Return the table name with a ``_v{N}`` suffix when versioning is enabled.""" | ||
| 77 | 77 | name = table.name | |
| 78 | 78 | if enable_versioning: | |
| 79 | 79 | version = getattr(table.projection, "version_tag", None) | |
| 80 | 80 | if version is None: | |
| 81 | 81 | version = getattr(table, "current_version_number", None) | |
| 82 | 82 | if version is not None and version > 0: | |
| 83 | 83 | name = f"{table.name}_v{version}" | |
| 84 | - return f"{project}_{name}" | ||
| 84 | + return name | ||
| 85 | + | ||
| 86 | + | ||
| 87 | + def compute_table_id(project: str, table: Any, enable_versioning: bool = False) -> str: | ||
| 88 | + """Build the online-store table name, appending a version suffix when versioning is enabled.""" | ||
| 89 | + return f"{project}_{compute_versioned_name(table, enable_versioning)}" | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -280,6 +280,18 @@ def _check_versioned_read_support(self, grouped_refs): | |||
| 280 | 280 | supported_types.append(FaissOnlineStore) | |
| 281 | 281 | except ImportError: | |
| 282 | 282 | pass | |
| 283 | + try: | ||
| 284 | + from feast.infra.online_stores.redis import RedisOnlineStore | ||
| 285 | + | ||
| 286 | + supported_types.append(RedisOnlineStore) | ||
| 287 | + except Exception: | ||
| 288 | + pass | ||
| 289 | + try: | ||
| 290 | + from feast.infra.online_stores.dynamodb import DynamoDBOnlineStore | ||
| 291 | + | ||
| 292 | + supported_types.append(DynamoDBOnlineStore) | ||
| 293 | + except Exception: | ||
| 294 | + pass | ||
| 283 | 295 | ||
| 284 | 296 | if isinstance(self, tuple(supported_types)): | |
| 285 | 297 | return | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -32,7 +32,12 @@ | |||
| 32 | 32 | from pydantic import StrictStr | |
| 33 | 33 | ||
| 34 | 34 | from feast import Entity, FeatureView, RepoConfig, utils | |
| 35 | - from feast.infra.online_stores.helpers import _mmh3, _redis_key, _redis_key_prefix | ||
| 35 | + from feast.infra.online_stores.helpers import ( | ||
| 36 | + _mmh3, | ||
| 37 | + _redis_key, | ||
| 38 | + _redis_key_prefix, | ||
| 39 | + compute_versioned_name, | ||
| 40 | + ) | ||
| 36 | 41 | from feast.infra.online_stores.online_store import OnlineStore | |
| 37 | 42 | from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto | |
| 38 | 43 | from feast.protos.feast.types.Value_pb2 import Value as ValueProto | |
@@ -51,6 +56,13 @@ | |||
| 51 | 56 | logger = logging.getLogger(__name__) | |
| 52 | 57 | ||
| 53 | 58 | ||
| 59 | + def _versioned_fv_name(table: FeatureView, config: RepoConfig) -> str: | ||
| 60 | + """Return the feature view name with version suffix when versioning is enabled.""" | ||
| 61 | + return compute_versioned_name( | ||
| 62 | + table, config.registry.enable_online_feature_view_versioning | ||
| 63 | + ) | ||
| 64 | + | ||
| 65 | + | ||
| 54 | 66 | class RedisType(str, Enum): | |
| 55 | 67 | redis = "redis" | |
| 56 | 68 | redis_cluster = "redis_cluster" | |
@@ -123,8 +135,9 @@ def delete_table(self, config: RepoConfig, table: FeatureView): | |||
| 123 | 135 | deleted_count = 0 | |
| 124 | 136 | prefix = _redis_key_prefix(table.join_keys) | |
| 125 | 137 | ||
| 126 | - redis_hash_keys = [_mmh3(f"{table.name}:{f.name}") for f in table.features] | ||
| 127 | - redis_hash_keys.append(bytes(f"_ts:{table.name}", "utf8")) | ||
| 138 | + fv_name = _versioned_fv_name(table, config) | ||
| 139 | + redis_hash_keys = [_mmh3(f"{fv_name}:{f.name}") for f in table.features] | ||
| 140 | + redis_hash_keys.append(bytes(f"_ts:{fv_name}", "utf8")) | ||
| 128 | 141 | ||
| 129 | 142 | with client.pipeline(transaction=False) as pipe: | |
| 130 | 143 | for _k in client.scan_iter( | |
@@ -133,7 +146,7 @@ def delete_table(self, config: RepoConfig, table: FeatureView): | |||
| 133 | 146 | _tables = { | |
| 134 | 147 | _hk[4:] for _hk in client.hgetall(_k) if _hk.startswith(b"_ts:") | |
| 135 | 148 | } | |
| 136 | - if bytes(table.name, "utf8") not in _tables: | ||
| 149 | + if bytes(fv_name, "utf8") not in _tables: | ||
| 137 | 150 | continue | |
| 138 | 151 | if len(_tables) == 1: | |
| 139 | 152 | pipe.delete(_k) | |
@@ -142,7 +155,7 @@ def delete_table(self, config: RepoConfig, table: FeatureView): | |||
| 142 | 155 | deleted_count += 1 | |
| 143 | 156 | pipe.execute() | |
| 144 | 157 | ||
| 145 | - logger.debug(f"Deleted {deleted_count} rows for feature view {table.name}") | ||
| 158 | + logger.debug(f"Deleted {deleted_count} rows for feature view {fv_name}") | ||
| 146 | 159 | ||
| 147 | 160 | def update( | |
| 148 | 161 | self, | |
@@ -281,7 +294,7 @@ def online_write_batch( | |||
| 281 | 294 | client = self._get_client(online_store_config) | |
| 282 | 295 | project = config.project | |
| 283 | 296 | ||
| 284 | - feature_view = table.name | ||
| 297 | + feature_view = _versioned_fv_name(table, config) | ||
| 285 | 298 | ts_key = f"_ts:{feature_view}" | |
| 286 | 299 | keys = [] | |
| 287 | 300 | # redis pipelining optimization: send multiple commands to redis server without waiting for every reply | |
@@ -355,13 +368,15 @@ def _generate_hset_keys_for_features( | |||
| 355 | 368 | self, | |
| 356 | 369 | feature_view: FeatureView, | |
| 357 | 370 | requested_features: Optional[List[str]] = None, | |
| 371 | + fv_name_override: Optional[str] = None, | ||
| 358 | 372 | ) -> Tuple[List[str], List[str]]: | |
| 359 | 373 | if not requested_features: | |
| 360 | 374 | requested_features = [f.name for f in feature_view.features] | |
| 361 | 375 | ||
| 362 | - hset_keys = [_mmh3(f"{feature_view.name}:{k}") for k in requested_features] | ||
| 376 | + fv_name = fv_name_override or feature_view.name | ||
| 377 | + hset_keys = [_mmh3(f"{fv_name}:{k}") for k in requested_features] | ||
| 363 | 378 | ||
| 364 | - ts_key = f"_ts:{feature_view.name}" | ||
| 379 | + ts_key = f"_ts:{fv_name}" | ||
| 365 | 380 | hset_keys.append(ts_key) | |
| 366 | 381 | requested_features.append(ts_key) | |
| 367 | 382 | ||
@@ -390,9 +405,10 @@ def online_read( | |||
| 390 | 405 | ||
| 391 | 406 | client = self._get_client(online_store_config) | |
| 392 | 407 | feature_view = table | |
| 408 | + fv_name = _versioned_fv_name(table, config) | ||
| 393 | 409 | ||
| 394 | 410 | requested_features, hset_keys = self._generate_hset_keys_for_features( | |
| 395 | - feature_view, requested_features | ||
| 411 | + feature_view, requested_features, fv_name_override=fv_name | ||
| 396 | 412 | ) | |
| 397 | 413 | keys = self._generate_redis_keys_for_entities(config, entity_keys) | |
| 398 | 414 | ||
@@ -403,7 +419,7 @@ def online_read( | |||
| 403 | 419 | redis_values = pipe.execute() | |
| 404 | 420 | ||
| 405 | 421 | return self._convert_redis_values_to_protobuf( | |
| 406 | - redis_values, feature_view.name, requested_features | ||
| 422 | + redis_values, fv_name, requested_features | ||
| 407 | 423 | ) | |
| 408 | 424 | ||
| 409 | 425 | async def online_read_async( | |
@@ -418,9 +434,10 @@ async def online_read_async( | |||
| 418 | 434 | ||
| 419 | 435 | client = await self._get_client_async(online_store_config) | |
| 420 | 436 | feature_view = table | |
| 437 | + fv_name = _versioned_fv_name(table, config) | ||
| 421 | 438 | ||
| 422 | 439 | requested_features, hset_keys = self._generate_hset_keys_for_features( | |
| 423 | - feature_view, requested_features | ||
| 440 | + feature_view, requested_features, fv_name_override=fv_name | ||
| 424 | 441 | ) | |
| 425 | 442 | keys = self._generate_redis_keys_for_entities(config, entity_keys) | |
| 426 | 443 | ||
@@ -430,7 +447,7 @@ async def online_read_async( | |||
| 430 | 447 | redis_values = await pipe.execute() | |
| 431 | 448 | ||
| 432 | 449 | return self._convert_redis_values_to_protobuf( | |
| 433 | - redis_values, feature_view.name, requested_features | ||
| 450 | + redis_values, fv_name, requested_features | ||
| 434 | 451 | ) | |
| 435 | 452 | ||
| 436 | 453 | def _get_features_for_entity( | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,194 @@ | |||
| 1 | + """Integration tests for DynamoDB online store feature view versioning. | ||
| 2 | + | ||
| 3 | + Run with: pytest --integration sdk/python/tests/integration/online_store/test_dynamodb_versioning.py | ||
| 4 | + | ||
| 5 | + Uses moto to mock the DynamoDB service (no Docker required). | ||
| 6 | + """ | ||
| 7 | + | ||
| 8 | + import os | ||
| 9 | + from datetime import datetime, timedelta, timezone | ||
| 10 | + | ||
| 11 | + import pytest | ||
| 12 | + | ||
| 13 | + from feast import Entity, FeatureView | ||
| 14 | + from feast.field import Field | ||
| 15 | + from feast.infra.online_stores.dynamodb import DynamoDBOnlineStore | ||
| 16 | + from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto | ||
| 17 | + from feast.protos.feast.types.Value_pb2 import Value as ValueProto | ||
| 18 | + from feast.repo_config import RegistryConfig, RepoConfig | ||
| 19 | + from feast.types import Float32, Int64 | ||
| 20 | + from feast.value_type import ValueType | ||
| 21 | + | ||
| 22 | + | ||
| 23 | + def _make_feature_view(name="driver_stats", version="latest"): | ||
| 24 | + entity = Entity( | ||
| 25 | + name="driver_id", | ||
| 26 | + join_keys=["driver_id"], | ||
| 27 | + value_type=ValueType.INT64, | ||
| 28 | + ) | ||
| 29 | + return FeatureView( | ||
| 30 | + name=name, | ||
| 31 | + entities=[entity], | ||
| 32 | + ttl=timedelta(days=1), | ||
| 33 | + schema=[ | ||
| 34 | + Field(name="driver_id", dtype=Int64), | ||
| 35 | + Field(name="trips_today", dtype=Int64), | ||
| 36 | + Field(name="avg_rating", dtype=Float32), | ||
| 37 | + ], | ||
| 38 | + version=version, | ||
| 39 | + ) | ||
| 40 | + | ||
| 41 | + | ||
| 42 | + def _make_entity_key(driver_id: int) -> EntityKeyProto: | ||
| 43 | + entity_key = EntityKeyProto() | ||
| 44 | + entity_key.join_keys.append("driver_id") | ||
| 45 | + val = ValueProto() | ||
| 46 | + val.int64_val = driver_id | ||
| 47 | + entity_key.entity_values.append(val) | ||
| 48 | + return entity_key | ||
| 49 | + | ||
| 50 | + | ||
| 51 | + def _write_and_read(store, config, fv, driver_id=1001, trips=42): | ||
| 52 | + entity_key = _make_entity_key(driver_id) | ||
| 53 | + val = ValueProto() | ||
| 54 | + val.int64_val = trips | ||
| 55 | + now = datetime.now(tz=timezone.utc) | ||
| 56 | + store.online_write_batch( | ||
| 57 | + config, fv, [(entity_key, {"trips_today": val}, now, now)], None | ||
| 58 | + ) | ||
| 59 | + return store.online_read(config, fv, [entity_key], ["trips_today"]) | ||
| 60 | + | ||
| 61 | + | ||
| 62 | + def _make_config(enable_versioning=False): | ||
| 63 | + from feast.infra.online_stores.dynamodb import DynamoDBOnlineStoreConfig | ||
| 64 | + | ||
| 65 | + return RepoConfig( | ||
| 66 | + project="test_project", | ||
| 67 | + provider="local", | ||
| 68 | + online_store=DynamoDBOnlineStoreConfig( | ||
| 69 | + type="dynamodb", | ||
| 70 | + region="us-east-1", | ||
| 71 | + ), | ||
| 72 | + registry=RegistryConfig( | ||
| 73 | + path="/tmp/test_dynamodb_registry.pb", | ||
| 74 | + enable_online_feature_view_versioning=enable_versioning, | ||
| 75 | + ), | ||
| 76 | + entity_key_serialization_version=3, | ||
| 77 | + ) | ||
| 78 | + | ||
| 79 | + | ||
| 80 | + @pytest.mark.integration | ||
| 81 | + class TestDynamoDBVersioningIntegration: | ||
| 82 | + """Integration tests for DynamoDB versioning using moto mock.""" | ||
| 83 | + | ||
| 84 | + @pytest.fixture(autouse=True) | ||
| 85 | + def setup_dynamodb(self): | ||
| 86 | + try: | ||
| 87 | + from moto import mock_dynamodb | ||
| 88 | + except ImportError: | ||
| 89 | + pytest.skip("moto not installed") | ||
| 90 | + | ||
| 91 | + # Set dummy AWS credentials for moto | ||
| 92 | + os.environ["AWS_ACCESS_KEY_ID"] = "testing" | ||
| 93 | + os.environ["AWS_SECRET_ACCESS_KEY"] = "testing" # noqa: S105 # pragma: allowlist secret | ||
| 94 | + os.environ["AWS_SECURITY_TOKEN"] = "testing" # noqa: S105 # pragma: allowlist secret | ||
| 95 | + os.environ["AWS_SESSION_TOKEN"] = "testing" # noqa: S105 # pragma: allowlist secret | ||
| 96 | + os.environ["AWS_DEFAULT_REGION"] = "us-east-1" | ||
| 97 | + | ||
| 98 | + with mock_dynamodb(): | ||
| 99 | + yield | ||
| 100 | + | ||
| 101 | + def test_write_read_without_versioning(self): | ||
| 102 | + config = _make_config(enable_versioning=False) | ||
| 103 | + store = DynamoDBOnlineStore() | ||
| 104 | + fv = _make_feature_view() | ||
| 105 | + store.update(config, [], [fv], [], [], False) | ||
| 106 | + | ||
| 107 | + result = _write_and_read(store, config, fv) | ||
| 108 | + assert result[0][1] is not None | ||
| 109 | + assert result[0][1]["trips_today"].int64_val == 42 | ||
| 110 | + | ||
| 111 | + def test_write_read_with_versioning_v1(self): | ||
| 112 | + config = _make_config(enable_versioning=True) | ||
| 113 | + store = DynamoDBOnlineStore() | ||
| 114 | + fv = _make_feature_view() | ||
| 115 | + fv.current_version_number = 1 | ||
| 116 | + store.update(config, [], [fv], [], [], False) | ||
| 117 | + | ||
| 118 | + result = _write_and_read(store, config, fv) | ||
| 119 | + assert result[0][1] is not None | ||
| 120 | + assert result[0][1]["trips_today"].int64_val == 42 | ||
| 121 | + | ||
| 122 | + def test_version_isolation(self): | ||
| 123 | + """Data written to v1 is not visible from v2.""" | ||
| 124 | + config = _make_config(enable_versioning=True) | ||
| 125 | + store = DynamoDBOnlineStore() | ||
| 126 | + | ||
| 127 | + fv_v1 = _make_feature_view() | ||
| 128 | + fv_v1.current_version_number = 1 | ||
| 129 | + store.update(config, [], [fv_v1], [], [], False) | ||
| 130 | + _write_and_read(store, config, fv_v1, driver_id=1001, trips=10) | ||
| 131 | + | ||
| 132 | + fv_v2 = _make_feature_view() | ||
| 133 | + fv_v2.current_version_number = 2 | ||
| 134 | + store.update(config, [], [fv_v2], [], [], False) | ||
| 135 | + | ||
| 136 | + entity_key = _make_entity_key(1001) | ||
| 137 | + result = store.online_read(config, fv_v2, [entity_key], ["trips_today"]) | ||
| 138 | + assert result[0] == (None, None) | ||
| 139 | + | ||
| 140 | + result = store.online_read(config, fv_v1, [entity_key], ["trips_today"]) | ||
| 141 | + assert result[0][1] is not None | ||
| 142 | + assert result[0][1]["trips_today"].int64_val == 10 | ||
| 143 | + | ||
| 144 | + def test_projection_version_tag_routes_to_correct_table(self): | ||
| 145 | + """projection.version_tag routes reads to the correct versioned DynamoDB table.""" | ||
| 146 | + config = _make_config(enable_versioning=True) | ||
| 147 | + store = DynamoDBOnlineStore() | ||
| 148 | + | ||
| 149 | + fv_v1 = _make_feature_view() | ||
| 150 | + fv_v1.current_version_number = 1 | ||
| 151 | + store.update(config, [], [fv_v1], [], [], False) | ||
| 152 | + _write_and_read(store, config, fv_v1, driver_id=1001, trips=100) | ||
| 153 | + | ||
| 154 | + fv_v2 = _make_feature_view() | ||
| 155 | + fv_v2.current_version_number = 2 | ||
| 156 | + store.update(config, [], [fv_v2], [], [], False) | ||
| 157 | + _write_and_read(store, config, fv_v2, driver_id=1001, trips=200) | ||
| 158 | + | ||
| 159 | + fv_read = _make_feature_view() | ||
| 160 | + fv_read.projection.version_tag = 1 | ||
| 161 | + entity_key = _make_entity_key(1001) | ||
| 162 | + result = store.online_read(config, fv_read, [entity_key], ["trips_today"]) | ||
| 163 | + assert result[0][1]["trips_today"].int64_val == 100 | ||
| 164 | + | ||
| 165 | + fv_read2 = _make_feature_view() | ||
| 166 | + fv_read2.projection.version_tag = 2 | ||
| 167 | + result = store.online_read(config, fv_read2, [entity_key], ["trips_today"]) | ||
| 168 | + assert result[0][1]["trips_today"].int64_val == 200 | ||
| 169 | + | ||
| 170 | + def test_teardown_versioned_table(self): | ||
| 171 | + """teardown() drops the versioned DynamoDB table without error.""" | ||
| 172 | + config = _make_config(enable_versioning=True) | ||
| 173 | + store = DynamoDBOnlineStore() | ||
| 174 | + | ||
| 175 | + fv = _make_feature_view() | ||
| 176 | + fv.current_version_number = 1 | ||
| 177 | + store.update(config, [], [fv], [], [], False) | ||
| 178 | + _write_and_read(store, config, fv) | ||
| 179 | + | ||
| 180 | + # Should not raise | ||
| 181 | + store.teardown(config, [fv], []) | ||
| 182 | + | ||
| 183 | + def test_update_deletes_versioned_table(self): | ||
| 184 | + """update() with tables_to_delete correctly drops versioned DynamoDB tables.""" | ||
| 185 | + config = _make_config(enable_versioning=True) | ||
| 186 | + store = DynamoDBOnlineStore() | ||
| 187 | + | ||
| 188 | + fv = _make_feature_view() | ||
| 189 | + fv.current_version_number = 1 | ||
| 190 | + store.update(config, [], [fv], [], [], False) | ||
| 191 | + _write_and_read(store, config, fv, driver_id=1001, trips=50) | ||
| 192 | + | ||
| 193 | + # Delete the versioned table | ||
| 194 | + store.update(config, [fv], [], [], [], False) | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments