Skip to content

Commit 09505d4

Browse files
committed
fix: Fixed file registry cache sync
Signed-off-by: ntkathole <nikhilkathole2683@gmail.com>
1 parent fc3ea20 commit 09505d4

2 files changed

Lines changed: 54 additions & 2 deletions

File tree

‎sdk/python/feast/infra/registry/file.py‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import json
2+
import os
23
import uuid
34
from pathlib import Path
45
from typing import Optional
@@ -44,6 +45,8 @@ def _write_registry(self, registry_proto: RegistryProto):
4445
file_dir.mkdir(exist_ok=True)
4546
with open(self._filepath, mode="wb", buffering=0) as f:
4647
f.write(registry_proto.SerializeToString())
48+
f.flush()
49+
os.fsync(f.fileno())
4750

4851
def set_project_metadata(self, project: str, key: str, value: str):
4952
"""Set a custom project metadata key-value pair in the registry proto (file backend)."""

‎sdk/python/feast/infra/registry/registry.py‎

Lines changed: 51 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -221,6 +221,13 @@ def __init__(
221221
else False
222222
)
223223

224+
self.cache_mode = (
225+
registry_config.cache_mode if registry_config is not None else "sync"
226+
)
227+
228+
self._file_mtime = None
229+
self._file_path = None
230+
224231
if registry_config:
225232
registry_store_type = registry_config.registry_store_type
226233
registry_path = registry_config.path
@@ -238,6 +245,13 @@ def __init__(
238245
)
239246
)
240247

248+
from feast.infra.registry.file import FileRegistryStore
249+
250+
if isinstance(self._registry_store, FileRegistryStore):
251+
self._file_path = self._registry_store._filepath
252+
if self._file_path.exists():
253+
self._file_mtime = self._file_path.stat().st_mtime
254+
241255
try:
242256
registry_proto = self._registry_store.get_registry_proto()
243257
self.cached_registry_proto = registry_proto
@@ -883,6 +897,11 @@ def commit(self):
883897
"""Commits the state of the registry cache to the remote registry store."""
884898
if self.cached_registry_proto:
885899
self._registry_store.update_registry_proto(self.cached_registry_proto)
900+
if self._file_path is not None and self._file_path.exists():
901+
try:
902+
self._file_mtime = self._file_path.stat().st_mtime
903+
except (OSError, FileNotFoundError):
904+
pass
886905

887906
def refresh(self, project: Optional[str] = None):
888907
"""Refreshes the state of the registry cache by fetching the registry state from the remote registry store."""
@@ -939,6 +958,23 @@ def _get_registry_proto(
939958
Returns: Returns a RegistryProto object which represents the state of the registry
940959
"""
941960
with self._refresh_lock:
961+
# For file-based registries in sync mode, check file modification time
962+
# to detect changes immediately, not just based on TTL
963+
file_modified = False
964+
if (
965+
allow_cache
966+
and self.cache_mode == "sync"
967+
and self._file_path is not None
968+
and self._file_path.exists()
969+
):
970+
try:
971+
current_mtime = self._file_path.stat().st_mtime
972+
if self._file_mtime is None or current_mtime > self._file_mtime:
973+
file_modified = True
974+
self._file_mtime = current_mtime
975+
except (OSError, FileNotFoundError):
976+
file_modified = True
977+
942978
expired = (self.cached_registry_proto_created is None) or (
943979
self.cached_registry_proto_ttl.total_seconds()
944980
> 0 # 0 ttl means infinity
@@ -951,12 +987,25 @@ def _get_registry_proto(
951987
)
952988
)
953989

954-
if allow_cache and not expired:
990+
# Refresh if expired or file was modified (for sync mode with file registry)
991+
if allow_cache and not expired and not file_modified:
955992
return self.cached_registry_proto
956-
logger.info("Registry cache expired, so refreshing")
993+
994+
if file_modified:
995+
logger.info("Registry file modified, so refreshing")
996+
else:
997+
logger.info("Registry cache expired, so refreshing")
998+
957999
registry_proto = self._registry_store.get_registry_proto()
9581000
self.cached_registry_proto = registry_proto
9591001
self.cached_registry_proto_created = _utc_now()
1002+
1003+
if self._file_path is not None and self._file_path.exists():
1004+
try:
1005+
self._file_mtime = self._file_path.stat().st_mtime
1006+
except (OSError, FileNotFoundError):
1007+
pass
1008+
9601009
return registry_proto
9611010

9621011
def _check_conflicting_feature_view_names(self, feature_view: BaseFeatureView):

0 commit comments

Comments
 (0)