|
31 | 31 | from fastapi import ( |
32 | 32 | Depends, |
33 | 33 | FastAPI, |
| 34 | + Query, |
34 | 35 | Request, |
35 | 36 | Response, |
36 | 37 | WebSocket, |
|
41 | 42 | from fastapi.logger import logger |
42 | 43 | from fastapi.responses import JSONResponse |
43 | 44 | from fastapi.staticfiles import StaticFiles |
44 | | -from pydantic import BaseModel, field_validator |
| 45 | +from pydantic import BaseModel, Field, field_validator |
45 | 46 |
|
46 | 47 | import feast |
47 | 48 | from feast import metrics as feast_metrics |
|
54 | 55 | ) |
55 | 56 | from feast.feast_object import FeastObject |
56 | 57 | from feast.feature_server_utils import convert_response_to_dict |
| 58 | +from feast.feature_view import FeatureViewState |
57 | 59 | from feast.feature_view_utils import get_feature_view_from_feature_store |
58 | 60 | from feast.filter_models import ComparisonFilter, CompoundFilter |
59 | 61 | from feast.permissions.action import WRITE, AuthzedAction |
@@ -94,12 +96,26 @@ class MaterializeRequest(BaseModel): |
94 | 96 | feature_views: Optional[List[str]] = None |
95 | 97 | disable_event_timestamp: bool = False |
96 | 98 | full_feature_names: bool = False |
| 99 | + version: Optional[str] = Field( |
| 100 | + None, |
| 101 | + description=( |
| 102 | + "Optional version to materialize (e.g. 'v2'). Requires feature_views " |
| 103 | + "with exactly one entry and registry.enable_online_feature_view_versioning." |
| 104 | + ), |
| 105 | + ) |
97 | 106 |
|
98 | 107 |
|
99 | 108 | class MaterializeIncrementalRequest(BaseModel): |
100 | 109 | end_ts: str |
101 | 110 | feature_views: Optional[List[str]] = None |
102 | 111 | full_feature_names: bool = False |
| 112 | + version: Optional[str] = Field( |
| 113 | + None, |
| 114 | + description=( |
| 115 | + "Optional version to materialize (e.g. 'v2'). Requires feature_views " |
| 116 | + "with exactly one entry and registry.enable_online_feature_view_versioning." |
| 117 | + ), |
| 118 | + ) |
103 | 119 |
|
104 | 120 |
|
105 | 121 | class GetOnlineFeaturesRequest(BaseModel): |
@@ -333,6 +349,128 @@ async def load_static_artifacts(app: FastAPI, store): |
333 | 349 | logger.warning(f"Failed to load static artifacts: {e}") |
334 | 350 |
|
335 | 351 |
|
| 352 | +def _authorize_materialize_views( |
| 353 | + store: "feast.FeatureStore", |
| 354 | + feature_view_names: Optional[List[str]], |
| 355 | + version: Optional[str] = None, |
| 356 | +) -> List[str]: |
| 357 | + """Resolve + authorize feature views for materialization. |
| 358 | +
|
| 359 | + Returns the resolved list of FV names (all eligible FVs when |
| 360 | + feature_view_names is None). |
| 361 | + """ |
| 362 | + parsed_version = store._validate_materialize_version(version, feature_view_names) |
| 363 | + feature_views_to_materialize = store._get_feature_views_to_materialize( |
| 364 | + feature_view_names, version=parsed_version |
| 365 | + ) |
| 366 | + for fv in feature_views_to_materialize: |
| 367 | + assert_permissions( |
| 368 | + resource=fv, |
| 369 | + actions=[AuthzedAction.WRITE_ONLINE], |
| 370 | + ) |
| 371 | + return [fv.name for fv in feature_views_to_materialize] |
| 372 | + |
| 373 | + |
| 374 | +def _check_already_materializing( |
| 375 | + store: "feast.FeatureStore", |
| 376 | + fv_names: List[str], |
| 377 | +) -> Optional[JSONResponse]: |
| 378 | + """Return a 409 JSONResponse if any requested FV is already MATERIALIZING.""" |
| 379 | + conflicting: List[str] = [] |
| 380 | + for fv_name in fv_names: |
| 381 | + try: |
| 382 | + fv = store.registry.get_feature_view( |
| 383 | + fv_name, store.project, allow_cache=False |
| 384 | + ) |
| 385 | + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: |
| 386 | + conflicting.append(fv_name) |
| 387 | + except (FeatureViewNotFoundException, KeyError): |
| 388 | + pass |
| 389 | + except Exception as e: |
| 390 | + logger.warning( |
| 391 | + f"Unexpected error checking MATERIALIZING state for {fv_name}: {e}" |
| 392 | + ) |
| 393 | + if conflicting: |
| 394 | + return JSONResponse( |
| 395 | + status_code=409, |
| 396 | + content={ |
| 397 | + "error": ( |
| 398 | + f"Cannot start async materialization — the following feature " |
| 399 | + f"views are already in MATERIALIZING state: {conflicting}. " |
| 400 | + f"Use ?force=true to override." |
| 401 | + ), |
| 402 | + "feature_views": conflicting, |
| 403 | + }, |
| 404 | + ) |
| 405 | + return None |
| 406 | + |
| 407 | + |
| 408 | +def _update_fv_state( |
| 409 | + store: "feast.FeatureStore", |
| 410 | + fv_names: List[str], |
| 411 | + state: FeatureViewState, |
| 412 | +) -> None: |
| 413 | + """Set FV state in the registry for each named feature view.""" |
| 414 | + for fv_name in fv_names: |
| 415 | + try: |
| 416 | + fv = store.registry.get_feature_view( |
| 417 | + fv_name, store.project, allow_cache=False |
| 418 | + ) |
| 419 | + fv.state = state |
| 420 | + store.registry.apply_feature_view(fv, store.project) |
| 421 | + except (FeatureViewNotFoundException, KeyError): |
| 422 | + logger.warning(f"Feature view {fv_name} not found; skip state={state}") |
| 423 | + except Exception as e: |
| 424 | + logger.warning(f"Failed to set state={state} for {fv_name}: {e}") |
| 425 | + |
| 426 | + |
| 427 | +def _reset_stuck_materializing_to_generated( |
| 428 | + store: "feast.FeatureStore", |
| 429 | + fv_names: List[str], |
| 430 | +) -> None: |
| 431 | + """Reset FVs currently in MATERIALIZING to GENERATED (force override).""" |
| 432 | + stuck: List[str] = [] |
| 433 | + for fv_name in fv_names: |
| 434 | + try: |
| 435 | + fv = store.registry.get_feature_view( |
| 436 | + fv_name, store.project, allow_cache=False |
| 437 | + ) |
| 438 | + if getattr(fv, "state", None) == FeatureViewState.MATERIALIZING: |
| 439 | + stuck.append(fv_name) |
| 440 | + except (FeatureViewNotFoundException, KeyError): |
| 441 | + pass |
| 442 | + except Exception as e: |
| 443 | + logger.warning(f"Unexpected error while force-resetting {fv_name}: {e}") |
| 444 | + if stuck: |
| 445 | + _update_fv_state(store, stuck, FeatureViewState.GENERATED) |
| 446 | + logger.info( |
| 447 | + "Force reset MATERIALIZING → GENERATED for feature views: %s", stuck |
| 448 | + ) |
| 449 | + |
| 450 | + |
| 451 | +def _parse_materialize_timestamps( |
| 452 | + request: "MaterializeRequest", |
| 453 | +) -> tuple: |
| 454 | + """Parse and validate start/end timestamps from a MaterializeRequest.""" |
| 455 | + if request.disable_event_timestamp: |
| 456 | + now = datetime.now() |
| 457 | + return datetime(1970, 1, 1), now |
| 458 | + |
| 459 | + if not request.start_ts or not request.end_ts: |
| 460 | + raise ValueError( |
| 461 | + "start_ts and end_ts are required when disable_event_timestamp is False" |
| 462 | + ) |
| 463 | + try: |
| 464 | + start_date = utils.make_tzaware(parser.parse(request.start_ts)) |
| 465 | + end_date = utils.make_tzaware(parser.parse(request.end_ts)) |
| 466 | + except (ValueError, TypeError) as e: |
| 467 | + raise ValueError(f"Invalid timestamp format: {e}") from e |
| 468 | + |
| 469 | + if start_date >= end_date: |
| 470 | + raise ValueError(f"start_ts ({start_date}) must be before end_ts ({end_date})") |
| 471 | + return start_date, end_date |
| 472 | + |
| 473 | + |
336 | 474 | def get_app( |
337 | 475 | store: "feast.FeatureStore", |
338 | 476 | registry_ttl_sec: int = DEFAULT_FEATURE_SERVER_REGISTRY_TTL, |
@@ -798,70 +936,115 @@ async def chat_ui(): |
798 | 936 | return Response(content=content, media_type="text/html") |
799 | 937 |
|
800 | 938 | @app.post("/materialize", dependencies=[Depends(inject_user_details)]) |
801 | | - async def materialize(request: MaterializeRequest) -> None: |
| 939 | + async def materialize( |
| 940 | + request: MaterializeRequest, |
| 941 | + async_mode: bool = Query(False, alias="async"), |
| 942 | + force: bool = Query(False), |
| 943 | + ): |
802 | 944 | with feast_metrics.track_request_latency("/materialize"): |
803 | | - if request.feature_views: |
804 | | - for feature_view in request.feature_views: |
805 | | - resource = await _get_feast_object(feature_view, True) |
806 | | - assert_permissions( |
807 | | - resource=resource, |
808 | | - actions=[AuthzedAction.WRITE_ONLINE], |
809 | | - ) |
810 | | - else: |
811 | | - feature_views_to_materialize = store._get_feature_views_to_materialize( |
812 | | - None |
813 | | - ) |
814 | | - for fv in feature_views_to_materialize: |
815 | | - assert_permissions( |
816 | | - resource=fv, |
817 | | - actions=[AuthzedAction.WRITE_ONLINE], |
818 | | - ) |
| 945 | + fv_names = _authorize_materialize_views( |
| 946 | + store, request.feature_views, version=request.version |
| 947 | + ) |
| 948 | + start_date, end_date = _parse_materialize_timestamps(request) |
819 | 949 |
|
820 | | - if request.disable_event_timestamp: |
821 | | - now = datetime.now() |
822 | | - start_date = datetime(1970, 1, 1) |
823 | | - end_date = now |
824 | | - else: |
825 | | - if not request.start_ts or not request.end_ts: |
826 | | - raise ValueError( |
827 | | - "start_ts and end_ts are required when disable_event_timestamp is False" |
828 | | - ) |
829 | | - start_date = utils.make_tzaware(parser.parse(request.start_ts)) |
830 | | - end_date = utils.make_tzaware(parser.parse(request.end_ts)) |
| 950 | + if async_mode: |
| 951 | + if force: |
| 952 | + _reset_stuck_materializing_to_generated(store, fv_names) |
| 953 | + else: |
| 954 | + conflict = _check_already_materializing(store, fv_names) |
| 955 | + if conflict: |
| 956 | + return conflict |
| 957 | + |
| 958 | + # Reserve MATERIALIZING before 202 so concurrent requests hit 409. |
| 959 | + # store.materialize() treats already-MATERIALIZING as a no-op. |
| 960 | + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) |
| 961 | + |
| 962 | + def _run_materialize(): |
| 963 | + try: |
| 964 | + store.materialize( |
| 965 | + start_date, |
| 966 | + end_date, |
| 967 | + fv_names, |
| 968 | + disable_event_timestamp=request.disable_event_timestamp, |
| 969 | + full_feature_names=request.full_feature_names, |
| 970 | + version=request.version, |
| 971 | + ) |
| 972 | + except Exception as e: |
| 973 | + logger.error( |
| 974 | + f"Async materialization failed for {fv_names}: {e}", |
| 975 | + exc_info=True, |
| 976 | + ) |
| 977 | + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) |
| 978 | + |
| 979 | + loop = asyncio.get_running_loop() |
| 980 | + loop.run_in_executor(None, _run_materialize) |
| 981 | + |
| 982 | + return JSONResponse( |
| 983 | + status_code=202, |
| 984 | + content={"status": "accepted", "feature_views": fv_names}, |
| 985 | + ) |
831 | 986 |
|
832 | 987 | await run_in_threadpool( |
833 | 988 | store.materialize, |
834 | 989 | start_date, |
835 | 990 | end_date, |
836 | | - request.feature_views, |
| 991 | + fv_names, |
837 | 992 | disable_event_timestamp=request.disable_event_timestamp, |
838 | 993 | full_feature_names=request.full_feature_names, |
| 994 | + version=request.version, |
839 | 995 | ) |
840 | 996 |
|
841 | 997 | @app.post("/materialize-incremental", dependencies=[Depends(inject_user_details)]) |
842 | | - async def materialize_incremental(request: MaterializeIncrementalRequest) -> None: |
| 998 | + async def materialize_incremental( |
| 999 | + request: MaterializeIncrementalRequest, |
| 1000 | + async_mode: bool = Query(False, alias="async"), |
| 1001 | + force: bool = Query(False), |
| 1002 | + ): |
843 | 1003 | with feast_metrics.track_request_latency("/materialize-incremental"): |
844 | | - if request.feature_views: |
845 | | - for feature_view in request.feature_views: |
846 | | - resource = await _get_feast_object(feature_view, True) |
847 | | - assert_permissions( |
848 | | - resource=resource, |
849 | | - actions=[AuthzedAction.WRITE_ONLINE], |
850 | | - ) |
851 | | - else: |
852 | | - feature_views_to_materialize = store._get_feature_views_to_materialize( |
853 | | - None |
| 1004 | + fv_names = _authorize_materialize_views( |
| 1005 | + store, request.feature_views, version=request.version |
| 1006 | + ) |
| 1007 | + end_date = utils.make_tzaware(parser.parse(request.end_ts)) |
| 1008 | + |
| 1009 | + if async_mode: |
| 1010 | + if force: |
| 1011 | + _reset_stuck_materializing_to_generated(store, fv_names) |
| 1012 | + else: |
| 1013 | + conflict = _check_already_materializing(store, fv_names) |
| 1014 | + if conflict: |
| 1015 | + return conflict |
| 1016 | + |
| 1017 | + _update_fv_state(store, fv_names, FeatureViewState.MATERIALIZING) |
| 1018 | + |
| 1019 | + def _run_materialize_incremental(): |
| 1020 | + try: |
| 1021 | + store.materialize_incremental( |
| 1022 | + end_date, |
| 1023 | + fv_names, |
| 1024 | + full_feature_names=request.full_feature_names, |
| 1025 | + version=request.version, |
| 1026 | + ) |
| 1027 | + except Exception as e: |
| 1028 | + logger.error( |
| 1029 | + f"Async materialize-incremental failed for {fv_names}: {e}", |
| 1030 | + exc_info=True, |
| 1031 | + ) |
| 1032 | + _update_fv_state(store, fv_names, FeatureViewState.GENERATED) |
| 1033 | + |
| 1034 | + loop = asyncio.get_running_loop() |
| 1035 | + loop.run_in_executor(None, _run_materialize_incremental) |
| 1036 | + |
| 1037 | + return JSONResponse( |
| 1038 | + status_code=202, |
| 1039 | + content={"status": "accepted", "feature_views": fv_names}, |
854 | 1040 | ) |
855 | | - for fv in feature_views_to_materialize: |
856 | | - assert_permissions( |
857 | | - resource=fv, |
858 | | - actions=[AuthzedAction.WRITE_ONLINE], |
859 | | - ) |
| 1041 | + |
860 | 1042 | await run_in_threadpool( |
861 | 1043 | store.materialize_incremental, |
862 | | - utils.make_tzaware(parser.parse(request.end_ts)), |
863 | | - request.feature_views, |
| 1044 | + end_date, |
| 1045 | + fv_names, |
864 | 1046 | full_feature_names=request.full_feature_names, |
| 1047 | + version=request.version, |
865 | 1048 | ) |
866 | 1049 |
|
867 | 1050 | @app.exception_handler(Exception) |
|
0 commit comments