From 8c0079755a583d0ec0b0c5b7afbc145d4cf3cb39 Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Tue, 7 Jun 2022 14:08:01 -0700 Subject: [PATCH 1/6] fix: Bugfixes for how registry is loaded Signed-off-by: Achal Shah --- sdk/python/feast/feature_store.py | 11 +- sdk/python/feast/infra/registry_stores/sql.py | 5 +- sdk/python/feast/registry.py | 138 +++++++++--------- sdk/python/feast/repo_operations.py | 3 +- 4 files changed, 81 insertions(+), 76 deletions(-) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 7a5a8299ebc..917e38a70c3 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -182,10 +182,13 @@ def refresh_registry(self): downloaded synchronously, which may increase latencies if the triggering method is get_online_features(). """ registry_config = self.config.get_registry_config() - registry = Registry(registry_config, repo_path=self.repo_path) - registry.refresh() - - self._registry = registry + if registry_config.registry_type == "sql": + self._registry = SqlRegistry(registry_config, None) + else: + r = Registry(registry_config, repo_path=self.repo_path) + r._initialize_registry() + self._registry = r + self._registry.refresh() @log_exceptions_and_usage def list_entities(self, allow_cache: bool = False) -> List[Entity]: diff --git a/sdk/python/feast/infra/registry_stores/sql.py b/sdk/python/feast/infra/registry_stores/sql.py index d34bb2fa8b6..503aaf86880 100644 --- a/sdk/python/feast/infra/registry_stores/sql.py +++ b/sdk/python/feast/infra/registry_stores/sql.py @@ -472,7 +472,7 @@ def update_infra(self, infra: Infra, project: str, commit: bool = True): pass def get_infra(self, project: str, allow_cache: bool = False) -> Infra: - pass + return Infra() def apply_user_metadata( self, @@ -550,7 +550,8 @@ def proto(self) -> RegistryProto: (self.list_validation_references, r.validation_references), ]: objs: List[Any] = lister(project) # type: ignore - registry_proto_field.extend([obj.to_proto() for obj in objs]) + if objs: + registry_proto_field.extend([obj.to_proto() for obj in objs]) return r diff --git a/sdk/python/feast/registry.py b/sdk/python/feast/registry.py index e993533c8b4..afb19299314 100644 --- a/sdk/python/feast/registry.py +++ b/sdk/python/feast/registry.py @@ -663,6 +663,75 @@ def commit(self): def refresh(self): """Refreshes the state of the registry cache by fetching the registry state from the remote registry store.""" + @staticmethod + def _message_to_sorted_dict(message: Message) -> Dict[str, Any]: + return json.loads(MessageToJson(message, sort_keys=True)) + + def to_dict(self, project: str) -> Dict[str, List[Any]]: + """Returns a dictionary representation of the registry contents for the specified project. + + For each list in the dictionary, the elements are sorted by name, so this + method can be used to compare two registries. + + Args: + project: Feast project to convert to a dict + """ + registry_dict: Dict[str, Any] = defaultdict(list) + registry_dict["project"] = project + for data_source in sorted( + self.list_data_sources(project=project), key=lambda ds: ds.name + ): + registry_dict["dataSources"].append( + self._message_to_sorted_dict(data_source.to_proto()) + ) + for entity in sorted( + self.list_entities(project=project), key=lambda entity: entity.name + ): + registry_dict["entities"].append( + self._message_to_sorted_dict(entity.to_proto()) + ) + for feature_view in sorted( + self.list_feature_views(project=project), + key=lambda feature_view: feature_view.name, + ): + registry_dict["featureViews"].append( + self._message_to_sorted_dict(feature_view.to_proto()) + ) + for feature_service in sorted( + self.list_feature_services(project=project), + key=lambda feature_service: feature_service.name, + ): + registry_dict["featureServices"].append( + self._message_to_sorted_dict(feature_service.to_proto()) + ) + for on_demand_feature_view in sorted( + self.list_on_demand_feature_views(project=project), + key=lambda on_demand_feature_view: on_demand_feature_view.name, + ): + odfv_dict = self._message_to_sorted_dict(on_demand_feature_view.to_proto()) + odfv_dict["spec"]["userDefinedFunction"]["body"] = dill.source.getsource( + on_demand_feature_view.udf + ) + registry_dict["onDemandFeatureViews"].append(odfv_dict) + for request_feature_view in sorted( + self.list_request_feature_views(project=project), + key=lambda request_feature_view: request_feature_view.name, + ): + registry_dict["requestFeatureViews"].append( + self._message_to_sorted_dict(request_feature_view.to_proto()) + ) + for saved_dataset in sorted( + self.list_saved_datasets(project=project), key=lambda item: item.name + ): + registry_dict["savedDatasets"].append( + self._message_to_sorted_dict(saved_dataset.to_proto()) + ) + for infra_object in sorted(self.get_infra(project=project).infra_objects): + registry_dict["infra"].append( + self._message_to_sorted_dict(infra_object.to_proto()) + ) + return registry_dict + class Registry(BaseRegistry): """ @@ -1587,75 +1656,6 @@ def teardown(self): def proto(self) -> RegistryProto: return self.cached_registry_proto or RegistryProto() - def to_dict(self, project: str) -> Dict[str, List[Any]]: - """Returns a dictionary representation of the registry contents for the specified project. - - For each list in the dictionary, the elements are sorted by name, so this - method can be used to compare two registries. - - Args: - project: Feast project to convert to a dict - """ - registry_dict: Dict[str, Any] = defaultdict(list) - registry_dict["project"] = project - for data_source in sorted( - self.list_data_sources(project=project), key=lambda ds: ds.name - ): - registry_dict["dataSources"].append( - self._message_to_sorted_dict(data_source.to_proto()) - ) - for entity in sorted( - self.list_entities(project=project), key=lambda entity: entity.name - ): - registry_dict["entities"].append( - self._message_to_sorted_dict(entity.to_proto()) - ) - for feature_view in sorted( - self.list_feature_views(project=project), - key=lambda feature_view: feature_view.name, - ): - registry_dict["featureViews"].append( - self._message_to_sorted_dict(feature_view.to_proto()) - ) - for feature_service in sorted( - self.list_feature_services(project=project), - key=lambda feature_service: feature_service.name, - ): - registry_dict["featureServices"].append( - self._message_to_sorted_dict(feature_service.to_proto()) - ) - for on_demand_feature_view in sorted( - self.list_on_demand_feature_views(project=project), - key=lambda on_demand_feature_view: on_demand_feature_view.name, - ): - odfv_dict = self._message_to_sorted_dict(on_demand_feature_view.to_proto()) - odfv_dict["spec"]["userDefinedFunction"]["body"] = dill.source.getsource( - on_demand_feature_view.udf - ) - registry_dict["onDemandFeatureViews"].append(odfv_dict) - for request_feature_view in sorted( - self.list_request_feature_views(project=project), - key=lambda request_feature_view: request_feature_view.name, - ): - registry_dict["requestFeatureViews"].append( - self._message_to_sorted_dict(request_feature_view.to_proto()) - ) - for saved_dataset in sorted( - self.list_saved_datasets(project=project), key=lambda item: item.name - ): - registry_dict["savedDatasets"].append( - self._message_to_sorted_dict(saved_dataset.to_proto()) - ) - for infra_object in sorted(self.get_infra(project=project).infra_objects): - registry_dict["infra"].append( - self._message_to_sorted_dict(infra_object.to_proto()) - ) - return registry_dict - - @staticmethod - def _message_to_sorted_dict(message: Message) -> Dict[str, Any]: - return json.loads(MessageToJson(message, sort_keys=True)) - def _prepare_registry_for_changes(self): """Prepares the Registry for changes by refreshing the cache if necessary.""" try: diff --git a/sdk/python/feast/repo_operations.py b/sdk/python/feast/repo_operations.py index 37daa6500eb..c94bbf67622 100644 --- a/sdk/python/feast/repo_operations.py +++ b/sdk/python/feast/repo_operations.py @@ -20,6 +20,7 @@ from feast.feature_service import FeatureService from feast.feature_store import FeatureStore from feast.feature_view import DUMMY_ENTITY, FeatureView +from feast.infra.registry_stores.sql import SqlRegistry from feast.names import adjectives, animals from feast.on_demand_feature_view import OnDemandFeatureView from feast.registry import FEAST_OBJECT_TYPES, FeastObjectType, Registry @@ -319,7 +320,7 @@ def registry_dump(repo_config: RepoConfig, repo_path: Path) -> str: """For debugging only: output contents of the metadata registry""" registry_config = repo_config.get_registry_config() project = repo_config.project - registry = Registry(registry_config=registry_config, repo_path=repo_path) + registry = SqlRegistry(registry_config=registry_config, repo_path=repo_path) registry_dict = registry.to_dict(project=project) return json.dumps(registry_dict, indent=2, sort_keys=True) From 52c36993c1709bad85cfeb60eac4a827cc588911 Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Tue, 7 Jun 2022 14:45:22 -0700 Subject: [PATCH 2/6] more fixes Signed-off-by: Achal Shah --- sdk/python/feast/registry.py | 9 +++++++++ sdk/python/feast/repo_operations.py | 3 +-- 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/sdk/python/feast/registry.py b/sdk/python/feast/registry.py index afb19299314..4d2a75e77f3 100644 --- a/sdk/python/feast/registry.py +++ b/sdk/python/feast/registry.py @@ -758,6 +758,15 @@ def get_user_metadata( cached_registry_proto_created: Optional[datetime] = None cached_registry_proto_ttl: timedelta + def __new__( + cls, registry_config: Optional[RegistryConfig], repo_path: Optional[Path] + ): + if registry_config and registry_config.registry_type == "sql": + from feast.infra.registry_stores.sql import SqlRegistry + + # all big numbers should be ClassB objects: + return SqlRegistry(registry_config, repo_path) + def __init__( self, registry_config: Optional[RegistryConfig], repo_path: Optional[Path] ): diff --git a/sdk/python/feast/repo_operations.py b/sdk/python/feast/repo_operations.py index c94bbf67622..37daa6500eb 100644 --- a/sdk/python/feast/repo_operations.py +++ b/sdk/python/feast/repo_operations.py @@ -20,7 +20,6 @@ from feast.feature_service import FeatureService from feast.feature_store import FeatureStore from feast.feature_view import DUMMY_ENTITY, FeatureView -from feast.infra.registry_stores.sql import SqlRegistry from feast.names import adjectives, animals from feast.on_demand_feature_view import OnDemandFeatureView from feast.registry import FEAST_OBJECT_TYPES, FeastObjectType, Registry @@ -320,7 +319,7 @@ def registry_dump(repo_config: RepoConfig, repo_path: Path) -> str: """For debugging only: output contents of the metadata registry""" registry_config = repo_config.get_registry_config() project = repo_config.project - registry = SqlRegistry(registry_config=registry_config, repo_path=repo_path) + registry = Registry(registry_config=registry_config, repo_path=repo_path) registry_dict = registry.to_dict(project=project) return json.dumps(registry_dict, indent=2, sort_keys=True) From dc41325b0fddef2a702d08f180848d7526f867af Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Tue, 7 Jun 2022 16:49:38 -0700 Subject: [PATCH 3/6] more fixes Signed-off-by: Achal Shah --- sdk/python/feast/registry.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdk/python/feast/registry.py b/sdk/python/feast/registry.py index 4d2a75e77f3..d9977fadcc2 100644 --- a/sdk/python/feast/registry.py +++ b/sdk/python/feast/registry.py @@ -764,8 +764,9 @@ def __new__( if registry_config and registry_config.registry_type == "sql": from feast.infra.registry_stores.sql import SqlRegistry - # all big numbers should be ClassB objects: return SqlRegistry(registry_config, repo_path) + else: + return super(Registry, cls).__new__(cls) def __init__( self, registry_config: Optional[RegistryConfig], repo_path: Optional[Path] From 6aa3becda3762c4c28fd8fd9284d45596a094a69 Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Tue, 7 Jun 2022 22:48:55 -0700 Subject: [PATCH 4/6] increase timeout Signed-off-by: Achal Shah --- sdk/python/tests/integration/registration/test_sql_registry.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/python/tests/integration/registration/test_sql_registry.py b/sdk/python/tests/integration/registration/test_sql_registry.py index c96d83ce0a1..1fe9ff5cecf 100644 --- a/sdk/python/tests/integration/registration/test_sql_registry.py +++ b/sdk/python/tests/integration/registration/test_sql_registry.py @@ -85,7 +85,7 @@ def mysql_registry(): log_string_to_wait_for = "/usr/sbin/mysqld: ready for connections. Version: '8.0.29' socket: '/var/run/mysqld/mysqld.sock' port: 3306" waited = wait_for_logs( - container=container, predicate=log_string_to_wait_for, timeout=30, interval=10, + container=container, predicate=log_string_to_wait_for, timeout=60, interval=10, ) logger.info("Waited for %s seconds until mysql container was up", waited) container_port = container.get_exposed_port(3306) From 8444039bd2297b4534e0b54dfb693e0ef3464c86 Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Wed, 8 Jun 2022 12:37:22 -0700 Subject: [PATCH 5/6] more fixes Signed-off-by: Achal Shah --- sdk/python/feast/registry.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/sdk/python/feast/registry.py b/sdk/python/feast/registry.py index d9977fadcc2..c8b00befc6b 100644 --- a/sdk/python/feast/registry.py +++ b/sdk/python/feast/registry.py @@ -761,6 +761,8 @@ def get_user_metadata( def __new__( cls, registry_config: Optional[RegistryConfig], repo_path: Optional[Path] ): + # We override __new__ so that we can inspect registry_config and create a SqlRegistry without callers + # needing to make any changes. if registry_config and registry_config.registry_type == "sql": from feast.infra.registry_stores.sql import SqlRegistry From 7e963b9f39b836c83fadfb39fab8425f7873c562 Mon Sep 17 00:00:00 2001 From: Achal Shah Date: Thu, 9 Jun 2022 09:17:14 -0700 Subject: [PATCH 6/6] more fixes Signed-off-by: Achal Shah --- sdk/python/feast/feature_store.py | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 917e38a70c3..7a5a8299ebc 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -182,13 +182,10 @@ def refresh_registry(self): downloaded synchronously, which may increase latencies if the triggering method is get_online_features(). """ registry_config = self.config.get_registry_config() - if registry_config.registry_type == "sql": - self._registry = SqlRegistry(registry_config, None) - else: - r = Registry(registry_config, repo_path=self.repo_path) - r._initialize_registry() - self._registry = r - self._registry.refresh() + registry = Registry(registry_config, repo_path=self.repo_path) + registry.refresh() + + self._registry = registry @log_exceptions_and_usage def list_entities(self, allow_cache: bool = False) -> List[Entity]: