FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

feat: Add feature view versioning support to Redis and DynamoDB onlin… · Marcus-Rosti/feast@edf25af · GitHub

forked from feast-dev/feast

Commit edf25af

Browse files
authored
feat: Add feature view versioning support to Redis and DynamoDB online stores (feast-dev#6257)
* feat: Add feature view versioning support to Redis and DynamoDB online stores Redis: Add _versioned_fv_name() helper that computes versioned feature view names (e.g. driver_stats_v2) used in hash field keys (_ts: and mmh3 feature hashes). This ensures version isolation within the same entity hash key. DynamoDB: Modify _get_table_name() to apply version suffix before template formatting, so each version gets its own DynamoDB table. Both stores are registered in _check_versioned_read_support() and the error message is updated accordingly. Closes feast-dev#6164, closes feast-dev#6163 Signed-off-by: yassinnouh21 <yassinnouh21@gmail.com> * fix: Add pragma allowlist for dummy moto credentials in DynamoDB test Signed-off-by: yassinnouh21 <yassinnouh21@gmail.com> * fix: Catch Exception instead of ImportError for Redis/DynamoDB in versioning gate Redis and DynamoDB modules raise FeastExtrasDependencyImportError (not ImportError) when their dependencies are missing. Use broad Exception catch so environments without redis/boto3 don't break. Signed-off-by: yassinnouh21 <yassinnouh21@gmail.com> * refactor: Extract shared compute_versioned_name helper to helpers.py Move version resolution logic into a shared compute_versioned_name() in helpers.py, reused by Redis, DynamoDB, and compute_table_id(). Addresses reviewer feedback to avoid duplicating version logic. Signed-off-by: yassinnouh21 <yassinnouh21@gmail.com> --------- Signed-off-by: yassinnouh21 <yassinnouh21@gmail.com>
1 parent de67bdd commit edf25af

9 files changed

Lines changed: 729 additions & 18 deletions

File tree

‎sdk/python/feast/errors.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -142,7 +142,7 @@ class VersionedOnlineReadNotSupported(FeastError):
142142
def __init__(self, store_name: str, version: int):
143143
super().__init__(
144144
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. "
146146
)
147147

148148

‎sdk/python/feast/infra/online_stores/dynamodb.py‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@
2424
from pydantic import StrictBool, StrictStr
2525

2626
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
2828
from feast.infra.online_stores.online_store import OnlineStore
2929
from feast.infra.supported_async_methods import SupportedAsyncMethods
3030
from feast.infra.utils.aws_utils import dynamo_write_items_async
@@ -1154,8 +1154,11 @@ def _initialize_dynamodb_resource(
11541154
def _get_table_name(
11551155
online_config: DynamoDBOnlineStoreConfig, config: RepoConfig, table: FeatureView
11561156
) -> str:
1157+
table_name = compute_versioned_name(
1158+
table, config.registry.enable_online_feature_view_versioning
1159+
)
11571160
return online_config.table_name_template.format(
1158-
project=config.project, table_name=table.name
1161+
project=config.project, table_name=table_name
11591162
)
11601163

11611164

‎sdk/python/feast/infra/online_stores/helpers.py‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -72,13 +72,18 @@ def _to_naive_utc(ts: datetime) -> datetime:
7272
return ts.astimezone(tz=timezone.utc).replace(tzinfo=None)
7373

7474

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."""
7777
name = table.name
7878
if enable_versioning:
7979
version = getattr(table.projection, "version_tag", None)
8080
if version is None:
8181
version = getattr(table, "current_version_number", None)
8282
if version is not None and version > 0:
8383
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)}"

‎sdk/python/feast/infra/online_stores/online_store.py‎

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -280,6 +280,18 @@ def _check_versioned_read_support(self, grouped_refs):
280280
supported_types.append(FaissOnlineStore)
281281
except ImportError:
282282
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
283295

284296
if isinstance(self, tuple(supported_types)):
285297
return

‎sdk/python/feast/infra/online_stores/redis.py‎

Lines changed: 29 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,12 @@
3232
from pydantic import StrictStr
3333

3434
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+
)
3641
from feast.infra.online_stores.online_store import OnlineStore
3742
from feast.protos.feast.types.EntityKey_pb2 import EntityKey as EntityKeyProto
3843
from feast.protos.feast.types.Value_pb2 import Value as ValueProto
@@ -51,6 +56,13 @@
5156
logger = logging.getLogger(__name__)
5257

5358

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+
5466
class RedisType(str, Enum):
5567
redis = "redis"
5668
redis_cluster = "redis_cluster"
@@ -123,8 +135,9 @@ def delete_table(self, config: RepoConfig, table: FeatureView):
123135
deleted_count = 0
124136
prefix = _redis_key_prefix(table.join_keys)
125137

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"))
128141

129142
with client.pipeline(transaction=False) as pipe:
130143
for _k in client.scan_iter(
@@ -133,7 +146,7 @@ def delete_table(self, config: RepoConfig, table: FeatureView):
133146
_tables = {
134147
_hk[4:] for _hk in client.hgetall(_k) if _hk.startswith(b"_ts:")
135148
}
136-
if bytes(table.name, "utf8") not in _tables:
149+
if bytes(fv_name, "utf8") not in _tables:
137150
continue
138151
if len(_tables) == 1:
139152
pipe.delete(_k)
@@ -142,7 +155,7 @@ def delete_table(self, config: RepoConfig, table: FeatureView):
142155
deleted_count += 1
143156
pipe.execute()
144157

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}")
146159

147160
def update(
148161
self,
@@ -281,7 +294,7 @@ def online_write_batch(
281294
client = self._get_client(online_store_config)
282295
project = config.project
283296

284-
feature_view = table.name
297+
feature_view = _versioned_fv_name(table, config)
285298
ts_key = f"_ts:{feature_view}"
286299
keys = []
287300
# redis pipelining optimization: send multiple commands to redis server without waiting for every reply
@@ -355,13 +368,15 @@ def _generate_hset_keys_for_features(
355368
self,
356369
feature_view: FeatureView,
357370
requested_features: Optional[List[str]] = None,
371+
fv_name_override: Optional[str] = None,
358372
) -> Tuple[List[str], List[str]]:
359373
if not requested_features:
360374
requested_features = [f.name for f in feature_view.features]
361375

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]
363378

364-
ts_key = f"_ts:{feature_view.name}"
379+
ts_key = f"_ts:{fv_name}"
365380
hset_keys.append(ts_key)
366381
requested_features.append(ts_key)
367382

@@ -390,9 +405,10 @@ def online_read(
390405

391406
client = self._get_client(online_store_config)
392407
feature_view = table
408+
fv_name = _versioned_fv_name(table, config)
393409

394410
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
396412
)
397413
keys = self._generate_redis_keys_for_entities(config, entity_keys)
398414

@@ -403,7 +419,7 @@ def online_read(
403419
redis_values = pipe.execute()
404420

405421
return self._convert_redis_values_to_protobuf(
406-
redis_values, feature_view.name, requested_features
422+
redis_values, fv_name, requested_features
407423
)
408424

409425
async def online_read_async(
@@ -418,9 +434,10 @@ async def online_read_async(
418434

419435
client = await self._get_client_async(online_store_config)
420436
feature_view = table
437+
fv_name = _versioned_fv_name(table, config)
421438

422439
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
424441
)
425442
keys = self._generate_redis_keys_for_entities(config, entity_keys)
426443

@@ -430,7 +447,7 @@ async def online_read_async(
430447
redis_values = await pipe.execute()
431448

432449
return self._convert_redis_values_to_protobuf(
433-
redis_values, feature_view.name, requested_features
450+
redis_values, fv_name, requested_features
434451
)
435452

436453
def _get_features_for_entity(
Lines changed: 194 additions & 0 deletions
Original file line numberDiff line numberDiff 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)

0 commit comments

Comments
 (0)

Back | FazBrowse Home | New Git URL