From f61be7d7663bea0571bbddfe71af670d504593b7 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Tue, 28 Jul 2026 16:39:44 +0530 Subject: [PATCH 1/8] feat(server): Add async materialization via ?async=true query param When ?async=true is passed to /materialize and /materialize-incremental, the endpoint fires off materialization in a background thread and returns 202 Accepted immediately. Concurrent requests for an already-MATERIALIZING FV are rejected with 409. Without ?async=true, behavior is unchanged. Client-side additions: - remote=True param on store.materialize() / materialize_incremental() to delegate to the feature server (URL/TLS from online_store config) - wait=False support with store.poll_materialization() for status polling - FeatureView state set to MATERIALIZING before 202, reset on failure Server-side additions: - ?async=true on existing /materialize and /materialize-incremental - ?force=true to override stuck MATERIALIZING state - Four module-level helpers for testability: _authorize_materialize_views, _check_already_materializing, _update_fv_state, _parse_materialize_timestamps Registry fixes: - SQL and Snowflake registries now set FV state to AVAILABLE_ONLINE in apply_materialization() (parity with file-based registry) Addresses feast-dev/feast#4526 Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 229 +++++++++++++++---- sdk/python/feast/feature_store.py | 90 +++++++- sdk/python/feast/feature_view.py | 1 + sdk/python/feast/infra/registry/snowflake.py | 20 +- sdk/python/feast/infra/registry/sql.py | 19 +- 5 files changed, 295 insertions(+), 64 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 1f374236790..7464ca876ab 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -13,6 +13,7 @@ # limitations under the License. import asyncio +import functools import os import sys import threading @@ -30,6 +31,7 @@ from fastapi import ( Depends, FastAPI, + Query, Request, Response, WebSocket, @@ -53,6 +55,7 @@ ) from feast.feast_object import FeastObject from feast.feature_server_utils import convert_response_to_dict +from feast.feature_view import FeatureViewState from feast.feature_view_utils import get_feature_view_from_feature_store from feast.filter_models import ComparisonFilter, CompoundFilter from feast.permissions.action import WRITE, AuthzedAction @@ -93,12 +96,14 @@ class MaterializeRequest(BaseModel): feature_views: Optional[List[str]] = None disable_event_timestamp: bool = False full_feature_names: bool = False + version: Optional[str] = None class MaterializeIncrementalRequest(BaseModel): end_ts: str feature_views: Optional[List[str]] = None full_feature_names: bool = False + version: Optional[str] = None class GetOnlineFeaturesRequest(BaseModel): @@ -332,6 +337,96 @@ async def load_static_artifacts(app: FastAPI, store): logger.warning(f"Failed to load static artifacts: {e}") +def _authorize_materialize_views( + store: "feast.FeatureStore", + feature_view_names: Optional[List[str]], +) -> List[str]: + """Resolve + authorize feature views for materialization. + + Returns the resolved list of FV names (all eligible FVs when + feature_view_names is None). + """ + feature_views_to_materialize = store._get_feature_views_to_materialize( + feature_view_names + ) + for fv in feature_views_to_materialize: + assert_permissions( + resource=fv, + actions=[AuthzedAction.WRITE_ONLINE], + ) + return [fv.name for fv in feature_views_to_materialize] + + +def _check_already_materializing( + store: "feast.FeatureStore", + fv_names: List[str], +) -> Optional[JSONResponse]: + """Return a 409 JSONResponse if any requested FV is already MATERIALIZING.""" + conflicting: List[str] = [] + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: + conflicting.append(fv_name) + except Exception: + pass + if conflicting: + return JSONResponse( + status_code=409, + content={ + "error": ( + f"Cannot start async materialization — the following feature " + f"views are already in MATERIALIZING state: {conflicting}. " + f"Use ?force=true to override." + ), + "feature_views": conflicting, + }, + ) + return None + + +def _update_fv_state( + store: "feast.FeatureStore", + fv_names: List[str], + state: FeatureViewState, +) -> None: + """Set FV state in the registry for each named feature view.""" + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + fv.state = state + store.registry.apply_feature_view(fv, store.project) + except Exception: + logger.warning(f"Failed to set state={state} for {fv_name}") + + +def _parse_materialize_timestamps( + request: "MaterializeRequest", +) -> tuple: + """Parse and validate start/end timestamps from a MaterializeRequest.""" + if request.disable_event_timestamp: + now = datetime.now() + return datetime(1970, 1, 1), now + + if not request.start_ts or not request.end_ts: + raise ValueError( + "start_ts and end_ts are required when disable_event_timestamp is False" + ) + try: + start_date = utils.make_tzaware(parser.parse(request.start_ts)) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + except (ValueError, TypeError) as e: + raise ValueError(f"Invalid timestamp format: {e}") from e + + if start_date >= end_date: + raise ValueError(f"start_ts ({start_date}) must be before end_ts ({end_date})") + return start_date, end_date + + def get_app( store: "feast.FeatureStore", registry_ttl_sec: int = DEFAULT_FEATURE_SERVER_REGISTRY_TTL, @@ -797,36 +892,47 @@ async def chat_ui(): return Response(content=content, media_type="text/html") @app.post("/materialize", dependencies=[Depends(inject_user_details)]) - async def materialize(request: MaterializeRequest) -> None: + async def materialize( + request: MaterializeRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None - ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + fv_names = _authorize_materialize_views(store, request.feature_views) + start_date, end_date = _parse_materialize_timestamps(request) + + if async_mode: + if not force: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize(): + try: + store.materialize( + start_date, + end_date, + fv_names, + disable_event_timestamp=request.disable_event_timestamp, + full_feature_names=request.full_feature_names, + version=request.version, + ) + except Exception as e: + logger.error( + f"Async materialization failed for {fv_names}: {e}", + exc_info=True, + ) + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) - if request.disable_event_timestamp: - now = datetime.now() - start_date = datetime(1970, 1, 1) - end_date = now - else: - if not request.start_ts or not request.end_ts: - raise ValueError( - "start_ts and end_ts are required when disable_event_timestamp is False" - ) - start_date = utils.make_tzaware(parser.parse(request.start_ts)) - end_date = utils.make_tzaware(parser.parse(request.end_ts)) + loop = asyncio.get_running_loop() + loop.run_in_executor(None, _run_materialize) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, + ) await run_in_threadpool( store.materialize, @@ -838,27 +944,48 @@ async def materialize(request: MaterializeRequest) -> None: ) @app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)]) - async def materialize_incremental(request: MaterializeIncrementalRequest) -> None: + async def materialize_incremental( + request: MaterializeIncrementalRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize-incremental"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None + fv_names = _authorize_materialize_views(store, request.feature_views) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + + if async_mode: + if not force: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize_incremental(): + try: + store.materialize_incremental( + end_date, + fv_names, + full_feature_names=request.full_feature_names, + ) + except Exception as e: + logger.error( + f"Async materialize-incremental failed for {fv_names}: {e}", + exc_info=True, + ) + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) + + loop = asyncio.get_running_loop() + loop.run_in_executor(None, _run_materialize_incremental) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + await run_in_threadpool( store.materialize_incremental, - utils.make_tzaware(parser.parse(request.end_ts)), + end_date, request.feature_views, full_feature_names=request.full_feature_names, ) @@ -1004,6 +1131,7 @@ def __init__( store=store, registry_ttl_sec=options["registry_ttl_sec"], ) + self._store = store self._options = options self._metrics_enabled = metrics_enabled super().__init__() @@ -1015,15 +1143,19 @@ def load_config(self): self.cfg.set("worker_class", "uvicorn_worker.UvicornWorker") if self._metrics_enabled: - self.cfg.set("post_worker_init", _gunicorn_post_worker_init) + self.cfg.set( + "post_worker_init", + functools.partial(_gunicorn_post_worker_init, self._store), + ) self.cfg.set("child_exit", _gunicorn_child_exit) def load(self): return self._app - def _gunicorn_post_worker_init(worker): - """Start per-worker resource monitoring after Gunicorn forks.""" + def _gunicorn_post_worker_init(store: "feast.FeatureStore", worker): + """Start per-worker resource and freshness monitoring after Gunicorn forks.""" feast_metrics.init_worker_monitoring() + feast_metrics.init_worker_freshness_monitoring(store) def _gunicorn_child_exit(server, worker): """Clean up Prometheus metric files for a dead worker.""" @@ -1061,6 +1193,7 @@ def start_server( store, metrics_config=flags, start_resource_monitoring=not uses_gunicorn, + start_freshness_monitoring=not uses_gunicorn, ) logger.debug("start_server called") diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index a35b02f3319..e5ad27b7101 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -510,7 +510,7 @@ def _transition_fv_to_materializing( Rolls back all already-transitioned FVs if this one can't transition. """ - previous_state = getattr(feature_view, "state", None) + previous_states[feature_view.name] = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED @@ -523,7 +523,6 @@ def _transition_fv_to_materializing( ) feature_view.state = FeatureViewState.MATERIALIZING self.registry.apply_feature_view(feature_view, self.project, commit=True) - previous_states[feature_view.name] = previous_state def _submit_and_process_materialization_jobs( self, @@ -1884,6 +1883,21 @@ def _emit_openlineage_apply(self, objects: List[Any]): def teardown(self): """Tears down all local and cloud resources for the feature store.""" + from feast.constants import PROTECTED_PROJECT_TAG + + # Prevent teardown of protected projects + try: + current = self.registry.get_project(name=self.project, allow_cache=False) + if current and current.tags.get(PROTECTED_PROJECT_TAG) == "true": + raise ValueError( + f'Teardown is not allowed on protected project "{self.project}". ' + "Protected projects are managed externally and cannot be torn down via Feast." + ) + except ValueError: + raise + except Exception: + pass + tables: List[BaseFeatureView] = [] tables.extend(self.list_feature_views()) tables.extend(self.list_label_views()) @@ -1891,7 +1905,10 @@ def teardown(self): entities = self.list_entities() self._get_provider().teardown_infra(self.project, tables, entities) # type: ignore[arg-type] - self.registry.teardown() + + for project in self.list_projects(): + self.registry.delete_project(project.name) + self._teardown_openlineage() def _teardown_openlineage(self): @@ -2409,6 +2426,31 @@ def _materialize_odfv( ) self.write_to_online_store(feature_view.name, df=transformed_df) + def _get_remote_materialize_url(self) -> str: + """Get the feature server URL from online_store.path for remote materialization.""" + online_cfg = self.config.online_store + url = getattr(online_cfg, "path", None) + if not url: + raise ValueError( + "online_store.path must be set to use remote materialization. " + "Configure online_store with type: remote and a valid path." + ) + return url.rstrip("/") + + def _get_remote_http_session(self): + """Get an HTTP session with auth configured for the feature server.""" + import requests + + auth_config = getattr(self.config, "auth_config", None) + if auth_config and getattr(auth_config, "type", "no_auth") != "no_auth": + from feast.permissions.client.http_auth_requests_wrapper import ( + get_http_auth_requests_session, + ) + + return get_http_auth_requests_session(auth_config) + + return requests.Session() + def materialize_incremental( self, end_date: datetime, @@ -2542,8 +2584,6 @@ def tqdm_builder(length): ) else: for feature_view, start_date in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") @@ -2575,7 +2615,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None @@ -2734,8 +2773,6 @@ def tqdm_builder(length): ) else: for feature_view, fv_start in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") @@ -2768,7 +2805,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None @@ -4732,6 +4768,9 @@ def list_projects( """ Retrieves the list of projects from the registry. + Protected projects (feast.dev/protected-project=true) are automatically + excluded from the results. + Args: allow_cache: Whether to allow returning projects from a cached registry. tags: Filter by tags. @@ -4739,7 +4778,10 @@ def list_projects( Returns: A list of projects. """ - return self.registry.list_projects(allow_cache=allow_cache, tags=tags) + from feast.constants import PROTECTED_PROJECT_TAG + + projects = self.registry.list_projects(allow_cache=allow_cache, tags=tags) + return [p for p in projects if p.tags.get(PROTECTED_PROJECT_TAG) != "true"] def get_project(self, name: Optional[str]) -> Project: """ @@ -4766,11 +4808,29 @@ def delete_project(self, name: str, commit: bool = True) -> None: Raises: ProjectNotFoundException: The project could not be found. + ValueError: If the project is protected. """ + from feast.constants import PROTECTED_PROJECT_TAG + + try: + project = self.registry.get_project(name=name, allow_cache=False) + if project and project.tags.get(PROTECTED_PROJECT_TAG) == "true": + raise ValueError( + f'Cannot delete protected project "{name}". ' + "Protected projects are managed externally." + ) + except ValueError: + raise + except Exception: + pass return self.registry.delete_project(name, commit=commit) def list_saved_datasets( - self, allow_cache: bool = False, tags: Optional[dict[str, str]] = None + self, + allow_cache: bool = False, + tags: Optional[dict[str, str]] = None, + namespace: Optional[str] = None, + collection: Optional[str] = None, ) -> List[SavedDataset]: """ Retrieves the list of saved datasets from the registry. @@ -4778,12 +4838,18 @@ def list_saved_datasets( Args: allow_cache: Whether to allow returning saved datasets from a cached registry. tags: Filter by tags. + namespace: Filter by logical namespace grouping. + collection: Filter by collection sub-grouping within namespace. Returns: A list of saved datasets. """ return self.registry.list_saved_datasets( - self.project, allow_cache=allow_cache, tags=tags + self.project, + allow_cache=allow_cache, + tags=tags, + namespace=namespace, + collection=collection, ) async def initialize(self) -> None: diff --git a/sdk/python/feast/feature_view.py b/sdk/python/feast/feature_view.py index a5d3c8d9537..0b7e28f555b 100644 --- a/sdk/python/feast/feature_view.py +++ b/sdk/python/feast/feature_view.py @@ -416,6 +416,7 @@ def __eq__(self, other): or normalize_version_string(self.version) != normalize_version_string(other.version) or self.org != other.org + or self.state != other.state ): return False diff --git a/sdk/python/feast/infra/registry/snowflake.py b/sdk/python/feast/infra/registry/snowflake.py index 5590e1b7574..fc75db20e9c 100644 --- a/sdk/python/feast/infra/registry/snowflake.py +++ b/sdk/python/feast/infra/registry/snowflake.py @@ -264,6 +264,7 @@ def apply_entity(self, entity: Entity, project: str, commit: bool = True): def apply_feature_service( self, feature_service: FeatureService, project: str, commit: bool = True ): + feature_service.prepare_for_apply(self, project, allow_cache=True) return self._apply_object( "FEATURE_SERVICES", project, @@ -945,13 +946,19 @@ def list_saved_datasets( project: str, allow_cache: bool = False, tags: Optional[dict[str, str]] = None, + namespace: Optional[str] = None, + collection: Optional[str] = None, ) -> List[SavedDataset]: if allow_cache: registry_proto = self._refresh_cached_registry_if_necessary() return proto_registry_utils.list_saved_datasets( - registry_proto, project, tags + registry_proto, + project, + tags, + namespace=namespace, + collection=collection, ) - return self._list_objects( + results = self._list_objects( "SAVED_DATASETS", project, SavedDatasetProto, @@ -959,6 +966,11 @@ def list_saved_datasets( "SAVED_DATASET_PROTO", tags=tags, ) + if namespace is not None: + results = [sd for sd in results if sd.namespace == namespace] + if collection is not None: + results = [sd for sd in results if sd.collection == collection] + return results def list_stream_feature_views( self, @@ -1135,6 +1147,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( fv_table_str, project, diff --git a/sdk/python/feast/infra/registry/sql.py b/sdk/python/feast/infra/registry/sql.py index edd89347be2..00e5f9c5253 100644 --- a/sdk/python/feast/infra/registry/sql.py +++ b/sdk/python/feast/infra/registry/sql.py @@ -1031,6 +1031,7 @@ def apply_feature_view( def apply_feature_service( self, feature_service: FeatureService, project: str, commit: bool = True ): + feature_service.prepare_for_apply(self, project, allow_cache=True) return self._apply_object( feature_services, project, @@ -1076,9 +1077,14 @@ def _list_feature_views( ) def _list_saved_datasets( - self, project: str, tags: Optional[dict[str, str]] = None, **kwargs + self, + project: str, + tags: Optional[dict[str, str]] = None, + namespace: Optional[str] = None, + collection: Optional[str] = None, + **kwargs, ) -> List[SavedDataset]: - return self._list_objects( + results = self._list_objects( saved_datasets, project, SavedDatasetProto, @@ -1087,6 +1093,11 @@ def _list_saved_datasets( tags=tags, **kwargs, ) + if namespace is not None: + results = [sd for sd in results if sd.namespace == namespace] + if collection is not None: + results = [sd for sd in results if sd.collection == collection] + return results def _list_on_demand_feature_views( self, project: str, tags: Optional[dict[str, str]], **kwargs @@ -1217,6 +1228,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( table, project, "feature_view_name", fv, "feature_view_proto" ) From 6c3aaa78e3098581ab6964a362591b64f9d2de97 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Tue, 28 Jul 2026 16:45:12 +0530 Subject: [PATCH 2/8] feat(server): Add async materialization via ?async=true query param When ?async=true is passed to /materialize and /materialize-incremental, the endpoint fires off materialization in a background thread and returns 202 Accepted immediately. Concurrent requests for an already-MATERIALIZING FV are rejected with 409. Without ?async=true, behavior is unchanged. Client-side additions: - remote=True param on store.materialize() / materialize_incremental() to delegate to the feature server (URL/TLS from online_store config) - wait=False support with store.poll_materialization() for status polling - FeatureView state set to MATERIALIZING before 202, reset on failure Server-side additions: - ?async=true on existing /materialize and /materialize-incremental - ?force=true to override stuck MATERIALIZING state - Four module-level helpers for testability: _authorize_materialize_views, _check_already_materializing, _update_fv_state, _parse_materialize_timestamps Registry fixes: - SQL and Snowflake registries now set FV state to AVAILABLE_ONLINE in apply_materialization() (parity with file-based registry) Addresses feast-dev/feast#4526 Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 216 +++++++++++++++---- sdk/python/feast/feature_store.py | 34 ++- sdk/python/feast/feature_view.py | 1 + sdk/python/feast/infra/registry/snowflake.py | 4 + sdk/python/feast/infra/registry/sql.py | 4 + 5 files changed, 206 insertions(+), 53 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 7412931ce8a..7464ca876ab 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -31,6 +31,7 @@ from fastapi import ( Depends, FastAPI, + Query, Request, Response, WebSocket, @@ -54,6 +55,7 @@ ) from feast.feast_object import FeastObject from feast.feature_server_utils import convert_response_to_dict +from feast.feature_view import FeatureViewState from feast.feature_view_utils import get_feature_view_from_feature_store from feast.filter_models import ComparisonFilter, CompoundFilter from feast.permissions.action import WRITE, AuthzedAction @@ -94,12 +96,14 @@ class MaterializeRequest(BaseModel): feature_views: Optional[List[str]] = None disable_event_timestamp: bool = False full_feature_names: bool = False + version: Optional[str] = None class MaterializeIncrementalRequest(BaseModel): end_ts: str feature_views: Optional[List[str]] = None full_feature_names: bool = False + version: Optional[str] = None class GetOnlineFeaturesRequest(BaseModel): @@ -333,6 +337,96 @@ async def load_static_artifacts(app: FastAPI, store): logger.warning(f"Failed to load static artifacts: {e}") +def _authorize_materialize_views( + store: "feast.FeatureStore", + feature_view_names: Optional[List[str]], +) -> List[str]: + """Resolve + authorize feature views for materialization. + + Returns the resolved list of FV names (all eligible FVs when + feature_view_names is None). + """ + feature_views_to_materialize = store._get_feature_views_to_materialize( + feature_view_names + ) + for fv in feature_views_to_materialize: + assert_permissions( + resource=fv, + actions=[AuthzedAction.WRITE_ONLINE], + ) + return [fv.name for fv in feature_views_to_materialize] + + +def _check_already_materializing( + store: "feast.FeatureStore", + fv_names: List[str], +) -> Optional[JSONResponse]: + """Return a 409 JSONResponse if any requested FV is already MATERIALIZING.""" + conflicting: List[str] = [] + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: + conflicting.append(fv_name) + except Exception: + pass + if conflicting: + return JSONResponse( + status_code=409, + content={ + "error": ( + f"Cannot start async materialization — the following feature " + f"views are already in MATERIALIZING state: {conflicting}. " + f"Use ?force=true to override." + ), + "feature_views": conflicting, + }, + ) + return None + + +def _update_fv_state( + store: "feast.FeatureStore", + fv_names: List[str], + state: FeatureViewState, +) -> None: + """Set FV state in the registry for each named feature view.""" + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + fv.state = state + store.registry.apply_feature_view(fv, store.project) + except Exception: + logger.warning(f"Failed to set state={state} for {fv_name}") + + +def _parse_materialize_timestamps( + request: "MaterializeRequest", +) -> tuple: + """Parse and validate start/end timestamps from a MaterializeRequest.""" + if request.disable_event_timestamp: + now = datetime.now() + return datetime(1970, 1, 1), now + + if not request.start_ts or not request.end_ts: + raise ValueError( + "start_ts and end_ts are required when disable_event_timestamp is False" + ) + try: + start_date = utils.make_tzaware(parser.parse(request.start_ts)) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + except (ValueError, TypeError) as e: + raise ValueError(f"Invalid timestamp format: {e}") from e + + if start_date >= end_date: + raise ValueError(f"start_ts ({start_date}) must be before end_ts ({end_date})") + return start_date, end_date + + def get_app( store: "feast.FeatureStore", registry_ttl_sec: int = DEFAULT_FEATURE_SERVER_REGISTRY_TTL, @@ -798,36 +892,47 @@ async def chat_ui(): return Response(content=content, media_type="text/html") @app.post("/materialize", dependencies=[Depends(inject_user_details)]) - async def materialize(request: MaterializeRequest) -> None: + async def materialize( + request: MaterializeRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None - ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + fv_names = _authorize_materialize_views(store, request.feature_views) + start_date, end_date = _parse_materialize_timestamps(request) + + if async_mode: + if not force: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize(): + try: + store.materialize( + start_date, + end_date, + fv_names, + disable_event_timestamp=request.disable_event_timestamp, + full_feature_names=request.full_feature_names, + version=request.version, + ) + except Exception as e: + logger.error( + f"Async materialization failed for {fv_names}: {e}", + exc_info=True, + ) + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) - if request.disable_event_timestamp: - now = datetime.now() - start_date = datetime(1970, 1, 1) - end_date = now - else: - if not request.start_ts or not request.end_ts: - raise ValueError( - "start_ts and end_ts are required when disable_event_timestamp is False" - ) - start_date = utils.make_tzaware(parser.parse(request.start_ts)) - end_date = utils.make_tzaware(parser.parse(request.end_ts)) + loop = asyncio.get_running_loop() + loop.run_in_executor(None, _run_materialize) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, + ) await run_in_threadpool( store.materialize, @@ -839,27 +944,48 @@ async def materialize(request: MaterializeRequest) -> None: ) @app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)]) - async def materialize_incremental(request: MaterializeIncrementalRequest) -> None: + async def materialize_incremental( + request: MaterializeIncrementalRequest, + async_mode: bool = Query(False, alias="async"), + force: bool = Query(False), + ): with feast_metrics.track_request_latency("/materialize-incremental"): - if request.feature_views: - for feature_view in request.feature_views: - resource = await _get_feast_object(feature_view, True) - assert_permissions( - resource=resource, - actions=[AuthzedAction.WRITE_ONLINE], - ) - else: - feature_views_to_materialize = store._get_feature_views_to_materialize( - None + fv_names = _authorize_materialize_views(store, request.feature_views) + end_date = utils.make_tzaware(parser.parse(request.end_ts)) + + if async_mode: + if not force: + conflict = _check_already_materializing(store, fv_names) + if conflict: + return conflict + + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + + def _run_materialize_incremental(): + try: + store.materialize_incremental( + end_date, + fv_names, + full_feature_names=request.full_feature_names, + ) + except Exception as e: + logger.error( + f"Async materialize-incremental failed for {fv_names}: {e}", + exc_info=True, + ) + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) + + loop = asyncio.get_running_loop() + loop.run_in_executor(None, _run_materialize_incremental) + + return JSONResponse( + status_code=202, + content={"status": "accepted", "feature_views": fv_names}, ) - for fv in feature_views_to_materialize: - assert_permissions( - resource=fv, - actions=[AuthzedAction.WRITE_ONLINE], - ) + await run_in_threadpool( store.materialize_incremental, - utils.make_tzaware(parser.parse(request.end_ts)), + end_date, request.feature_views, full_feature_names=request.full_feature_names, ) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 34a77310bac..e5ad27b7101 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -510,7 +510,7 @@ def _transition_fv_to_materializing( Rolls back all already-transitioned FVs if this one can't transition. """ - previous_state = getattr(feature_view, "state", None) + previous_states[feature_view.name] = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED @@ -523,7 +523,6 @@ def _transition_fv_to_materializing( ) feature_view.state = FeatureViewState.MATERIALIZING self.registry.apply_feature_view(feature_view, self.project, commit=True) - previous_states[feature_view.name] = previous_state def _submit_and_process_materialization_jobs( self, @@ -2427,6 +2426,31 @@ def _materialize_odfv( ) self.write_to_online_store(feature_view.name, df=transformed_df) + def _get_remote_materialize_url(self) -> str: + """Get the feature server URL from online_store.path for remote materialization.""" + online_cfg = self.config.online_store + url = getattr(online_cfg, "path", None) + if not url: + raise ValueError( + "online_store.path must be set to use remote materialization. " + "Configure online_store with type: remote and a valid path." + ) + return url.rstrip("/") + + def _get_remote_http_session(self): + """Get an HTTP session with auth configured for the feature server.""" + import requests + + auth_config = getattr(self.config, "auth_config", None) + if auth_config and getattr(auth_config, "type", "no_auth") != "no_auth": + from feast.permissions.client.http_auth_requests_wrapper import ( + get_http_auth_requests_session, + ) + + return get_http_auth_requests_session(auth_config) + + return requests.Session() + def materialize_incremental( self, end_date: datetime, @@ -2560,8 +2584,6 @@ def tqdm_builder(length): ) else: for feature_view, start_date in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") @@ -2593,7 +2615,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None @@ -2752,8 +2773,6 @@ def tqdm_builder(length): ) else: for feature_view, fv_start in regular_fvs_with_dates: - # Transition state to MATERIALIZING before starting. - # Only enforce when the state machine is active (not STATE_UNSPECIFIED). previous_state = getattr(feature_view, "state", None) if ( hasattr(feature_view, "state") @@ -2786,7 +2805,6 @@ def tqdm_builder(length): ) except Exception: fv_success = False - # Roll back state to previous value on failure. if ( hasattr(feature_view, "state") and previous_state is not None diff --git a/sdk/python/feast/feature_view.py b/sdk/python/feast/feature_view.py index a5d3c8d9537..0b7e28f555b 100644 --- a/sdk/python/feast/feature_view.py +++ b/sdk/python/feast/feature_view.py @@ -416,6 +416,7 @@ def __eq__(self, other): or normalize_version_string(self.version) != normalize_version_string(other.version) or self.org != other.org + or self.state != other.state ): return False diff --git a/sdk/python/feast/infra/registry/snowflake.py b/sdk/python/feast/infra/registry/snowflake.py index 0438e997ab1..fc75db20e9c 100644 --- a/sdk/python/feast/infra/registry/snowflake.py +++ b/sdk/python/feast/infra/registry/snowflake.py @@ -1147,6 +1147,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( fv_table_str, project, diff --git a/sdk/python/feast/infra/registry/sql.py b/sdk/python/feast/infra/registry/sql.py index 0a037e519c2..00e5f9c5253 100644 --- a/sdk/python/feast/infra/registry/sql.py +++ b/sdk/python/feast/infra/registry/sql.py @@ -1228,6 +1228,10 @@ def apply_materialization( FeatureViewNotFoundException, ) fv.materialization_intervals.append((start_date, end_date)) + if hasattr(fv, "state"): + from feast.feature_view import FeatureViewState + + fv.state = FeatureViewState.AVAILABLE_ONLINE self._apply_object( table, project, "feature_view_name", fv, "feature_view_proto" ) From 1384c604650fa7a2e2ea6838dde98a81cdcebc16 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Tue, 28 Jul 2026 17:43:15 +0530 Subject: [PATCH 3/8] feat(sdk): Auto-delegate materialize to feature server when online_store is remote When online_store.type == remote, materialize() and materialize_incremental() POST to the feature server with ?async=true (fire-and-forget) instead of running the engine locally. Optional force=True maps to ?force=true for stuck MATERIALIZING recovery. URL/TLS/auth come from online_store config. Shared _delegate_remote_materialize() builds query params and posts; RemoteComputeEngine remains a follow-up for provider-layer remoting. Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_store.py | 90 +++++++++++++++++++++++++++++++ 1 file changed, 90 insertions(+) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index e5ad27b7101..8ca497c71c1 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -2451,12 +2451,71 @@ def _get_remote_http_session(self): return requests.Session() + def _is_remote_topology(self) -> bool: + """True when this client talks to a remote feature server for online ops.""" + return getattr(self.config.online_store, "type", None) == "remote" + + def _post_to_feature_server( + self, + endpoint: str, + payload: Dict[str, Any], + query_params: Optional[Dict[str, str]] = None, + ) -> Dict[str, Any]: + """POST JSON to the feature server; raise on 409 / 4xx / 5xx.""" + url = f"{self._get_remote_materialize_url()}{endpoint}" + session = self._get_remote_http_session() + cert = getattr(self.config.online_store, "cert", "") or "" + verify: Any = cert if cert else True + try: + response = session.post( + url, + json=payload, + params=query_params or {}, + verify=verify, + ) + if response.status_code == 409: + try: + detail = response.json() + message = detail.get("error", response.text) + except Exception: + message = response.text + raise RuntimeError( + f"Remote materialization conflict (409): {message}" + ) from None + if response.status_code >= 400: + raise RuntimeError( + f"Remote materialization failed " + f"({response.status_code}): {response.text}" + ) + if not response.content: + return {} + try: + return response.json() + except Exception: + return {"status": "accepted", "raw": response.text} + finally: + session.close() + + def _delegate_remote_materialize( + self, + endpoint: str, + payload: Dict[str, Any], + force: bool = False, + ) -> None: + """Fire-and-forget POST to feature server with ?async=true.""" + query_params = {"async": "true"} + if force: + query_params["force"] = "true" + result = self._post_to_feature_server(endpoint, payload, query_params) + _logger.info("Remote materialization accepted (%s): %s", endpoint, result) + def materialize_incremental( self, end_date: datetime, feature_views: Optional[List[str]] = None, full_feature_names: bool = False, version: Optional[str] = None, + force: bool = False, ) -> None: """ Materialize incremental new data from the offline store into the online store. @@ -2475,6 +2534,8 @@ def materialize_incremental( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. + force (bool): When using remote topology, pass force=true to override stuck + MATERIALIZING state on the feature server. Ignored for local topology. Raises: Exception: A feature view being materialized does not have a TTL set. @@ -2490,6 +2551,19 @@ def materialize_incremental( ... """ + if self._is_remote_topology(): + payload: Dict[str, Any] = { + "end_ts": end_date.isoformat(), + "feature_views": feature_views, + "full_feature_names": full_feature_names, + } + if version is not None: + payload["version"] = version + self._delegate_remote_materialize( + "/materialize-incremental", payload, force=force + ) + return + parsed_version = self._validate_materialize_version(version, feature_views) feature_views_to_materialize = self._get_feature_views_to_materialize( feature_views, version=parsed_version @@ -2676,6 +2750,7 @@ def materialize( disable_event_timestamp: bool = False, full_feature_names: bool = False, version: Optional[str] = None, + force: bool = False, ) -> None: """ Materialize data from the offline store into the online store. @@ -2694,6 +2769,8 @@ def materialize( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. + force (bool): When using remote topology, pass force=true to override stuck + MATERIALIZING state on the feature server. Ignored for local topology. Examples: Materialize all features into the online store over the interval @@ -2708,6 +2785,19 @@ def materialize( ... """ + if self._is_remote_topology(): + payload: Dict[str, Any] = { + "start_ts": start_date.isoformat(), + "end_ts": end_date.isoformat(), + "feature_views": feature_views, + "disable_event_timestamp": disable_event_timestamp, + "full_feature_names": full_feature_names, + } + if version is not None: + payload["version"] = version + self._delegate_remote_materialize("/materialize", payload, force=force) + return + if utils.make_tzaware(start_date) > utils.make_tzaware(end_date): raise ValueError( f"The given start_date {start_date} is greater than the given end_date {end_date}." From b5cb173cb6f789abfb35c071f539bc80696db976 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Tue, 28 Jul 2026 19:16:11 +0530 Subject: [PATCH 4/8] =?UTF-8?q?fix(server):=20Let=20store=20own=20MATERIAL?= =?UTF-8?q?IZING=20transitions=20on=20async=20path=20Remove=20server-side?= =?UTF-8?q?=20MATERIALIZING=20pre-set=20before=20202=20which=20conflicted?= =?UTF-8?q?=20with=20store.materialize()=20state=20machine=20(MATERIALIZIN?= =?UTF-8?q?G=20=E2=86=92=20MATERIALIZING=20rejected).=20Async=20now=20acce?= =?UTF-8?q?pts=20and=20runs=20materialize=20in=20the=20background;=20store?= =?UTF-8?q?=20owns=20transitions.=20=3Fforce=3Dtrue=20resets=20stuck=20MAT?= =?UTF-8?q?ERIALIZING=20FVs=20to=20GENERATED=20so=20normal=20materialize?= =?UTF-8?q?=20can=20proceed.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 40 +++++++++++++++++++++++++----- 1 file changed, 34 insertions(+), 6 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 7464ca876ab..dcd2f9d1a46 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -404,6 +404,31 @@ def _update_fv_state( logger.warning(f"Failed to set state={state} for {fv_name}") +def _reset_stuck_materializing_to_generated( + store: "feast.FeatureStore", + fv_names: List[str], +) -> None: + """Reset FVs currently in MATERIALIZING to GENERATED (force override). + + Leaves other states untouched so store.materialize() can transition normally. + """ + stuck: List[str] = [] + for fv_name in fv_names: + try: + fv = store.registry.get_feature_view( + fv_name, store.project, allow_cache=False + ) + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: + stuck.append(fv_name) + except Exception: + pass + if stuck: + _update_fv_state(store, stuck, FeatureViewState.GENERATED) + logger.info( + "Force reset MATERIALIZING → GENERATED for feature views: %s", stuck + ) + + def _parse_materialize_timestamps( request: "MaterializeRequest", ) -> tuple: @@ -902,13 +927,15 @@ async def materialize( start_date, end_date = _parse_materialize_timestamps(request) if async_mode: - if not force: + if force: + _reset_stuck_materializing_to_generated(store, fv_names) + else: conflict = _check_already_materializing(store, fv_names) if conflict: return conflict - _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) - + # State transitions (MATERIALIZING / AVAILABLE_ONLINE) are owned + # by store.materialize(); server only accepts and runs in background. def _run_materialize(): try: store.materialize( @@ -954,13 +981,14 @@ async def materialize_incremental( end_date = utils.make_tzaware(parser.parse(request.end_ts)) if async_mode: - if not force: + if force: + _reset_stuck_materializing_to_generated(store, fv_names) + else: conflict = _check_already_materializing(store, fv_names) if conflict: return conflict - _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) - + # State transitions owned by store.materialize_incremental(). def _run_materialize_incremental(): try: store.materialize_incremental( From d8e7a9974bdffbd20cbf23920db2cc6980cead63 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Fri, 31 Jul 2026 16:20:30 +0530 Subject: [PATCH 5/8] =?UTF-8?q?fix:=20Address=20#6649=20review=20=E2=80=94?= =?UTF-8?q?=20run=5Fasync,=20race=20reserve,=20version=20threading=20-=20R?= =?UTF-8?q?evert=20FeatureView.state=20from=20=5F=5Feq=5F=5F=20-=20Add=20r?= =?UTF-8?q?un=5Fasync=20for=20remote=20sync=20vs=20async=20HTTP=20-=20Rese?= =?UTF-8?q?rve=20MATERIALIZING=20before=20202;=20idempotent=20store=20tran?= =?UTF-8?q?sitions=20-=20Thread=20version=20through=20authorize=20and=20al?= =?UTF-8?q?l=20materialize=20server=20paths=20-=20Narrow=20silent=20except?= =?UTF-8?q?s;=20document=20version=20on=20request=20models?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 67 +++++++++++++++++++++-------- sdk/python/feast/feature_store.py | 69 ++++++++++++++++++++++++------ sdk/python/feast/feature_view.py | 1 - 3 files changed, 104 insertions(+), 33 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index dcd2f9d1a46..32212ef892b 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -42,7 +42,7 @@ from fastapi.logger import logger from fastapi.responses import JSONResponse from fastapi.staticfiles import StaticFiles -from pydantic import BaseModel, field_validator +from pydantic import BaseModel, Field, field_validator import feast from feast import metrics as feast_metrics @@ -96,14 +96,26 @@ class MaterializeRequest(BaseModel): feature_views: Optional[List[str]] = None disable_event_timestamp: bool = False full_feature_names: bool = False - version: Optional[str] = None + version: Optional[str] = Field( + None, + description=( + "Optional version to materialize (e.g. 'v2'). Requires feature_views " + "with exactly one entry and registry.enable_online_feature_view_versioning." + ), + ) class MaterializeIncrementalRequest(BaseModel): end_ts: str feature_views: Optional[List[str]] = None full_feature_names: bool = False - version: Optional[str] = None + version: Optional[str] = Field( + None, + description=( + "Optional version to materialize (e.g. 'v2'). Requires feature_views " + "with exactly one entry and registry.enable_online_feature_view_versioning." + ), + ) class GetOnlineFeaturesRequest(BaseModel): @@ -340,14 +352,16 @@ async def load_static_artifacts(app: FastAPI, store): def _authorize_materialize_views( store: "feast.FeatureStore", feature_view_names: Optional[List[str]], + version: Optional[str] = None, ) -> List[str]: """Resolve + authorize feature views for materialization. Returns the resolved list of FV names (all eligible FVs when feature_view_names is None). """ + parsed_version = store._validate_materialize_version(version, feature_view_names) feature_views_to_materialize = store._get_feature_views_to_materialize( - feature_view_names + feature_view_names, version=parsed_version ) for fv in feature_views_to_materialize: assert_permissions( @@ -370,8 +384,12 @@ def _check_already_materializing( ) if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: conflicting.append(fv_name) - except Exception: + except (FeatureViewNotFoundException, KeyError): pass + except Exception as e: + logger.warning( + f"Unexpected error checking MATERIALIZING state for {fv_name}: {e}" + ) if conflicting: return JSONResponse( status_code=409, @@ -400,18 +418,17 @@ def _update_fv_state( ) fv.state = state store.registry.apply_feature_view(fv, store.project) - except Exception: - logger.warning(f"Failed to set state={state} for {fv_name}") + except (FeatureViewNotFoundException, KeyError): + logger.warning(f"Feature view {fv_name} not found; skip state={state}") + except Exception as e: + logger.warning(f"Failed to set state={state} for {fv_name}: {e}") def _reset_stuck_materializing_to_generated( store: "feast.FeatureStore", fv_names: List[str], ) -> None: - """Reset FVs currently in MATERIALIZING to GENERATED (force override). - - Leaves other states untouched so store.materialize() can transition normally. - """ + """Reset FVs currently in MATERIALIZING to GENERATED (force override).""" stuck: List[str] = [] for fv_name in fv_names: try: @@ -420,8 +437,10 @@ def _reset_stuck_materializing_to_generated( ) if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: stuck.append(fv_name) - except Exception: + except (FeatureViewNotFoundException, KeyError): pass + except Exception as e: + logger.warning(f"Unexpected error while force-resetting {fv_name}: {e}") if stuck: _update_fv_state(store, stuck, FeatureViewState.GENERATED) logger.info( @@ -923,7 +942,9 @@ async def materialize( force: bool = Query(False), ): with feast_metrics.track_request_latency("/materialize"): - fv_names = _authorize_materialize_views(store, request.feature_views) + fv_names = _authorize_materialize_views( + store, request.feature_views, version=request.version + ) start_date, end_date = _parse_materialize_timestamps(request) if async_mode: @@ -934,8 +955,10 @@ async def materialize( if conflict: return conflict - # State transitions (MATERIALIZING / AVAILABLE_ONLINE) are owned - # by store.materialize(); server only accepts and runs in background. + # Reserve MATERIALIZING before 202 so concurrent requests hit 409. + # store.materialize() treats already-MATERIALIZING as a no-op. + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + def _run_materialize(): try: store.materialize( @@ -965,9 +988,10 @@ def _run_materialize(): store.materialize, start_date, end_date, - request.feature_views, + fv_names, disable_event_timestamp=request.disable_event_timestamp, full_feature_names=request.full_feature_names, + version=request.version, ) @app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)]) @@ -977,7 +1001,9 @@ async def materialize_incremental( force: bool = Query(False), ): with feast_metrics.track_request_latency("/materialize-incremental"): - fv_names = _authorize_materialize_views(store, request.feature_views) + fv_names = _authorize_materialize_views( + store, request.feature_views, version=request.version + ) end_date = utils.make_tzaware(parser.parse(request.end_ts)) if async_mode: @@ -988,13 +1014,15 @@ async def materialize_incremental( if conflict: return conflict - # State transitions owned by store.materialize_incremental(). + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) + def _run_materialize_incremental(): try: store.materialize_incremental( end_date, fv_names, full_feature_names=request.full_feature_names, + version=request.version, ) except Exception as e: logger.error( @@ -1014,8 +1042,9 @@ def _run_materialize_incremental(): await run_in_threadpool( store.materialize_incremental, end_date, - request.feature_views, + fv_names, full_feature_names=request.full_feature_names, + version=request.version, ) @app.exception_handler(Exception) diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 6d6480d8f34..f3bc1cd8fa9 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -509,8 +509,15 @@ def _transition_fv_to_materializing( Transition a feature view to MATERIALIZING state. Rolls back all already-transitioned FVs if this one can't transition. + Already MATERIALIZING is a no-op (async server may have reserved the state + before returning 202); rollback target is GENERATED in that case. """ - previous_states[feature_view.name] = getattr(feature_view, "state", None) + current = getattr(feature_view, "state", None) + if current == FeatureViewState.MATERIALIZING: + previous_states[feature_view.name] = FeatureViewState.GENERATED + return + + previous_states[feature_view.name] = current if ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED @@ -2501,13 +2508,32 @@ def _delegate_remote_materialize( endpoint: str, payload: Dict[str, Any], force: bool = False, + run_async: bool = True, ) -> None: - """Fire-and-forget POST to feature server with ?async=true.""" - query_params = {"async": "true"} + """POST materialize to the feature server. + + When run_async=True (default), sends ?async=true and returns after 202. + When run_async=False, omits async and blocks until the server finishes + synchronous materialization (HTTP response). + force=True is only valid with run_async=True (server force applies to async). + """ + if force and not run_async: + raise ValueError( + "force=True requires run_async=True. " + "force only overrides stuck MATERIALIZING on the async path." + ) + + query_params: Dict[str, str] = {} + if run_async: + query_params["async"] = "true" if force: query_params["force"] = "true" - result = self._post_to_feature_server(endpoint, payload, query_params) - _logger.info("Remote materialization accepted (%s): %s", endpoint, result) + + result = self._post_to_feature_server(endpoint, payload, query_params or None) + if run_async: + _logger.info("Remote materialization accepted (%s): %s", endpoint, result) + else: + _logger.info("Remote materialization completed (%s): %s", endpoint, result) def materialize_incremental( self, @@ -2516,6 +2542,7 @@ def materialize_incremental( full_feature_names: bool = False, version: Optional[str] = None, force: bool = False, + run_async: bool = True, ) -> None: """ Materialize incremental new data from the offline store into the online store. @@ -2534,8 +2561,11 @@ def materialize_incremental( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. - force (bool): When using remote topology, pass force=true to override stuck - MATERIALIZING state on the feature server. Ignored for local topology. + force (bool): When using remote topology with run_async=True, pass force=true to + override stuck MATERIALIZING state on the feature server. Ignored for local topology. + run_async (bool): When using remote topology, if True (default) POST with ?async=true + and return after 202. If False, POST without async and block until the server + finishes sync materialization. Ignored for local topology. Raises: Exception: A feature view being materialized does not have a TTL set. @@ -2560,7 +2590,10 @@ def materialize_incremental( if version is not None: payload["version"] = version self._delegate_remote_materialize( - "/materialize-incremental", payload, force=force + "/materialize-incremental", + payload, + force=force, + run_async=run_async, ) return @@ -2659,7 +2692,9 @@ def tqdm_builder(length): else: for feature_view, start_date in regular_fvs_with_dates: previous_state = getattr(feature_view, "state", None) - if ( + if previous_state == FeatureViewState.MATERIALIZING: + previous_state = FeatureViewState.GENERATED + elif ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED ): @@ -2751,6 +2786,7 @@ def materialize( full_feature_names: bool = False, version: Optional[str] = None, force: bool = False, + run_async: bool = True, ) -> None: """ Materialize data from the offline store into the online store. @@ -2769,8 +2805,11 @@ def materialize( feature view name. version (str): Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry and enable_online_feature_view_versioning to be enabled. - force (bool): When using remote topology, pass force=true to override stuck - MATERIALIZING state on the feature server. Ignored for local topology. + force (bool): When using remote topology with run_async=True, pass force=true to + override stuck MATERIALIZING state on the feature server. Ignored for local topology. + run_async (bool): When using remote topology, if True (default) POST with ?async=true + and return after 202. If False, POST without async and block until the server + finishes sync materialization. Ignored for local topology. Examples: Materialize all features into the online store over the interval @@ -2795,7 +2834,9 @@ def materialize( } if version is not None: payload["version"] = version - self._delegate_remote_materialize("/materialize", payload, force=force) + self._delegate_remote_materialize( + "/materialize", payload, force=force, run_async=run_async + ) return if utils.make_tzaware(start_date) > utils.make_tzaware(end_date): @@ -2864,7 +2905,9 @@ def tqdm_builder(length): else: for feature_view, fv_start in regular_fvs_with_dates: previous_state = getattr(feature_view, "state", None) - if ( + if previous_state == FeatureViewState.MATERIALIZING: + previous_state = FeatureViewState.GENERATED + elif ( hasattr(feature_view, "state") and feature_view.state != FeatureViewState.STATE_UNSPECIFIED ): diff --git a/sdk/python/feast/feature_view.py b/sdk/python/feast/feature_view.py index 626ff5c91cf..6d1a77b1114 100644 --- a/sdk/python/feast/feature_view.py +++ b/sdk/python/feast/feature_view.py @@ -414,7 +414,6 @@ def __eq__(self, other): or normalize_version_string(self.version) != normalize_version_string(other.version) or self.org != other.org - or self.state != other.state ): return False From c602fe300aa2fa60d7e270dee8de46714a8149ff Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Wed, 5 Aug 2026 23:18:41 +0530 Subject: [PATCH 6/8] fix(server): Only reset MATERIALIZING FVs on async materialize failure Avoid wiping AVAILABLE_ONLINE when a post-success exception fires (e.g. SparkApp CR already cleaned up). Reuse the stuck-state helper so true failures still return to GENERATED. Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 32212ef892b..9defb87a27c 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -974,7 +974,7 @@ def _run_materialize(): f"Async materialization failed for {fv_names}: {e}", exc_info=True, ) - _update_fv_state(store, fv_names, FeatureViewState.GENERATED) + _reset_stuck_materializing_to_generated(store, fv_names) loop = asyncio.get_running_loop() loop.run_in_executor(None, _run_materialize) @@ -1029,7 +1029,7 @@ def _run_materialize_incremental(): f"Async materialize-incremental failed for {fv_names}: {e}", exc_info=True, ) - _update_fv_state(store, fv_names, FeatureViewState.GENERATED) + _reset_stuck_materializing_to_generated(store, fv_names) loop = asyncio.get_running_loop() loop.run_in_executor(None, _run_materialize_incremental) From 7909fc4a71371049a3a68b8aa36f533b42d9422d Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Wed, 5 Aug 2026 23:24:19 +0530 Subject: [PATCH 7/8] fix: Default run_async to False and cover version forwarding Preserve sync semantics for remote users; opt into fire-and-forget with run_async=True. Add unit tests that version is threaded through authorize and both sync/async materialize server paths. Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_store.py | 24 +-- .../test_materialize_version_forwarding.py | 169 ++++++++++++++++++ 2 files changed, 181 insertions(+), 12 deletions(-) create mode 100644 sdk/python/tests/unit/test_materialize_version_forwarding.py diff --git a/sdk/python/feast/feature_store.py b/sdk/python/feast/feature_store.py index 3f6d8f89c6d..546f6f7680d 100644 --- a/sdk/python/feast/feature_store.py +++ b/sdk/python/feast/feature_store.py @@ -2516,13 +2516,13 @@ def _delegate_remote_materialize( endpoint: str, payload: Dict[str, Any], force: bool = False, - run_async: bool = True, + run_async: bool = False, ) -> None: """POST materialize to the feature server. - When run_async=True (default), sends ?async=true and returns after 202. - When run_async=False, omits async and blocks until the server finishes - synchronous materialization (HTTP response). + When run_async=False (default), omits async and blocks until the server + finishes synchronous materialization (HTTP response). + When run_async=True, sends ?async=true and returns after 202. force=True is only valid with run_async=True (server force applies to async). """ if force and not run_async: @@ -2550,7 +2550,7 @@ def materialize_incremental( full_feature_names: bool = False, version: Optional[str] = None, force: bool = False, - run_async: bool = True, + run_async: bool = False, ) -> None: """ Materialize incremental new data from the offline store into the online store. @@ -2571,9 +2571,9 @@ def materialize_incremental( with exactly one entry and enable_online_feature_view_versioning to be enabled. force (bool): When using remote topology with run_async=True, pass force=true to override stuck MATERIALIZING state on the feature server. Ignored for local topology. - run_async (bool): When using remote topology, if True (default) POST with ?async=true - and return after 202. If False, POST without async and block until the server - finishes sync materialization. Ignored for local topology. + run_async (bool): When using remote topology, if False (default) POST without async and + block until the server finishes sync materialization. If True, POST with ?async=true + and return after 202. Ignored for local topology. Raises: Exception: A feature view being materialized does not have a TTL set. @@ -2794,7 +2794,7 @@ def materialize( full_feature_names: bool = False, version: Optional[str] = None, force: bool = False, - run_async: bool = True, + run_async: bool = False, ) -> None: """ Materialize data from the offline store into the online store. @@ -2815,9 +2815,9 @@ def materialize( with exactly one entry and enable_online_feature_view_versioning to be enabled. force (bool): When using remote topology with run_async=True, pass force=true to override stuck MATERIALIZING state on the feature server. Ignored for local topology. - run_async (bool): When using remote topology, if True (default) POST with ?async=true - and return after 202. If False, POST without async and block until the server - finishes sync materialization. Ignored for local topology. + run_async (bool): When using remote topology, if False (default) POST without async and + block until the server finishes sync materialization. If True, POST with ?async=true + and return after 202. Ignored for local topology. Examples: Materialize all features into the online store over the interval diff --git a/sdk/python/tests/unit/test_materialize_version_forwarding.py b/sdk/python/tests/unit/test_materialize_version_forwarding.py new file mode 100644 index 00000000000..3f675e7269d --- /dev/null +++ b/sdk/python/tests/unit/test_materialize_version_forwarding.py @@ -0,0 +1,169 @@ +# Copyright 2026 The Feast Authors +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +"""Server tests: version is forwarded on materialize sync and async paths.""" + +from concurrent.futures import Future +from inspect import signature +from unittest.mock import AsyncMock, MagicMock, patch + +from fastapi.testclient import TestClient + +from feast.feature_server import _authorize_materialize_views, get_app +from feast.feature_store import FeatureStore +from feast.feature_view import FeatureViewState + + +def _mock_store_for_materialize(): + fs = MagicMock() + fs.project = "test_project" + fs.initialize = AsyncMock() + fs.close = AsyncMock() + fs._validate_materialize_version.return_value = "v2" + fv = MagicMock() + fv.name = "driver_hourly_stats" + fv.state = FeatureViewState.GENERATED + fs._get_feature_views_to_materialize.return_value = [fv] + fs.registry.get_feature_view.return_value = fv + fs.materialize = MagicMock() + fs.materialize_incremental = MagicMock() + return fs + + +def _run_executor_inline(loop, executor, func, *args): + """Make asyncio run_in_executor execute synchronously for tests.""" + fut: Future = Future() + try: + fut.set_result(func(*args) if args else func()) + except Exception as exc: + fut.set_exception(exc) + return fut + + +def test_run_async_defaults_to_false(): + for method_name in ( + "materialize", + "materialize_incremental", + "_delegate_remote_materialize", + ): + param = signature(getattr(FeatureStore, method_name)).parameters["run_async"] + assert param.default is False, ( + f"{method_name}.run_async default should be False" + ) + + +@patch("feast.feature_server.assert_permissions") +def test_authorize_materialize_views_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + names = _authorize_materialize_views(fs, ["driver_hourly_stats"], version="v2") + assert names == ["driver_hourly_stats"] + fs._validate_materialize_version.assert_called_once_with( + "v2", ["driver_hourly_stats"] + ) + fs._get_feature_views_to_materialize.assert_called_once_with( + ["driver_hourly_stats"], version="v2" + ) + + +@patch("feast.feature_server.assert_permissions") +def test_materialize_sync_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 200 + fs._validate_materialize_version.assert_called_with("v2", ["driver_hourly_stats"]) + fs.materialize.assert_called_once() + _, kwargs = fs.materialize.call_args + assert kwargs.get("version") == "v2" + assert fs.materialize.call_args.args[2] == ["driver_hourly_stats"] + + +@patch("feast.feature_server.assert_permissions") +def test_materialize_incremental_sync_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize-incremental", + json={ + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 200 + fs._validate_materialize_version.assert_called_with("v2", ["driver_hourly_stats"]) + fs.materialize_incremental.assert_called_once() + _, kwargs = fs.materialize_incremental.call_args + assert kwargs.get("version") == "v2" + assert fs.materialize_incremental.call_args.args[1] == ["driver_hourly_stats"] + + +@patch( + "asyncio.BaseEventLoop.run_in_executor", + new=_run_executor_inline, +) +@patch("feast.feature_server.assert_permissions") +def test_materialize_async_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize?async=true", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 202 + fs.materialize.assert_called_once() + _, kwargs = fs.materialize.call_args + assert kwargs.get("version") == "v2" + + +@patch( + "asyncio.BaseEventLoop.run_in_executor", + new=_run_executor_inline, +) +@patch("feast.feature_server.assert_permissions") +def test_materialize_incremental_async_forwards_version(_mock_perms): + fs = _mock_store_for_materialize() + client = TestClient(get_app(fs)) + + response = client.post( + "/materialize-incremental?async=true", + json={ + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + "version": "v2", + }, + ) + + assert response.status_code == 202 + fs.materialize_incremental.assert_called_once() + _, kwargs = fs.materialize_incremental.call_args + assert kwargs.get("version") == "v2" From a2810c86694fe2b1d9c0f08706e8ebd6c4a96eb7 Mon Sep 17 00:00:00 2001 From: Aniket Paluskar Date: Thu, 6 Aug 2026 14:09:29 +0530 Subject: [PATCH 8/8] fix(server): Use dedicated executor for async materialize Isolate long materialize waits from the default/shared thread pool used by online serving. Pool size via FEAST_MATERIALIZE_MAX_WORKERS (default 2); shut down on app lifespan exit without blocking. Signed-off-by: Aniket Paluskar --- sdk/python/feast/feature_server.py | 24 ++++++++++++-- .../test_materialize_version_forwarding.py | 32 ++++++++++++++++++- 2 files changed, 53 insertions(+), 3 deletions(-) diff --git a/sdk/python/feast/feature_server.py b/sdk/python/feast/feature_server.py index 9defb87a27c..e42e28d6db9 100644 --- a/sdk/python/feast/feature_server.py +++ b/sdk/python/feast/feature_server.py @@ -20,6 +20,7 @@ import time import traceback from collections import defaultdict +from concurrent.futures import ThreadPoolExecutor from contextlib import asynccontextmanager from datetime import datetime from importlib import resources as importlib_resources @@ -550,6 +551,22 @@ def get_app( else: logger.debug("Offline write batching is DISABLED") + # Dedicated pool for async materialize so long Spark/offline waits do not + # starve the default executor used by online serving and run_in_threadpool. + _mat_workers_raw = os.environ.get("FEAST_MATERIALIZE_MAX_WORKERS", "2") + try: + materialize_max_workers = max(1, int(_mat_workers_raw)) + except ValueError: + logger.warning( + "Invalid FEAST_MATERIALIZE_MAX_WORKERS=%r; using default 2", + _mat_workers_raw, + ) + materialize_max_workers = 2 + materialize_executor = ThreadPoolExecutor( + max_workers=materialize_max_workers, + thread_name_prefix="feast-materialize", + ) + def stop_refresh(): nonlocal shutting_down shutting_down = True @@ -583,6 +600,9 @@ async def lifespan(app: FastAPI): stop_refresh() if offline_batcher is not None: offline_batcher.shutdown() + # wait=False: do not block process exit on in-flight materialize + # (same fire-and-forget contract as returning 202 mid-job). + materialize_executor.shutdown(wait=False) await store.close() app = FastAPI(lifespan=lifespan) @@ -977,7 +997,7 @@ def _run_materialize(): _reset_stuck_materializing_to_generated(store, fv_names) loop = asyncio.get_running_loop() - loop.run_in_executor(None, _run_materialize) + loop.run_in_executor(materialize_executor, _run_materialize) return JSONResponse( status_code=202, @@ -1032,7 +1052,7 @@ def _run_materialize_incremental(): _reset_stuck_materializing_to_generated(store, fv_names) loop = asyncio.get_running_loop() - loop.run_in_executor(None, _run_materialize_incremental) + loop.run_in_executor(materialize_executor, _run_materialize_incremental) return JSONResponse( status_code=202, diff --git a/sdk/python/tests/unit/test_materialize_version_forwarding.py b/sdk/python/tests/unit/test_materialize_version_forwarding.py index 3f675e7269d..cd7f6b40d61 100644 --- a/sdk/python/tests/unit/test_materialize_version_forwarding.py +++ b/sdk/python/tests/unit/test_materialize_version_forwarding.py @@ -13,7 +13,7 @@ # limitations under the License. """Server tests: version is forwarded on materialize sync and async paths.""" -from concurrent.futures import Future +from concurrent.futures import Future, ThreadPoolExecutor from inspect import signature from unittest.mock import AsyncMock, MagicMock, patch @@ -120,6 +120,36 @@ def test_materialize_incremental_sync_forwards_version(_mock_perms): assert fs.materialize_incremental.call_args.args[1] == ["driver_hourly_stats"] +@patch("feast.feature_server.assert_permissions") +def test_materialize_async_uses_dedicated_executor(_mock_perms): + """Async materialize must not use the shared default pool (executor=None).""" + fs = _mock_store_for_materialize() + seen: dict = {} + + def capture_executor(loop, executor, func, *args): + seen["executor"] = executor + return _run_executor_inline(loop, executor, func, *args) + + with patch( + "asyncio.BaseEventLoop.run_in_executor", + new=capture_executor, + ): + client = TestClient(get_app(fs)) + response = client.post( + "/materialize?async=true", + json={ + "start_ts": "2021-01-01T00:00:00", + "end_ts": "2021-01-02T00:00:00", + "feature_views": ["driver_hourly_stats"], + }, + ) + + assert response.status_code == 202 + assert seen.get("executor") is not None + assert isinstance(seen["executor"], ThreadPoolExecutor) + assert "feast-materialize" in seen["executor"]._thread_name_prefix + + @patch( "asyncio.BaseEventLoop.run_in_executor", new=_run_executor_inline,