Skip to content

Commit 0ea48e8

Browse files
committed
Merge PR feast-dev#6649 (feat/remote-materialize-v2) into merged-spark-e2e
2 parents 79b33ce + d8e7a99 commit 0ea48e8

4 files changed

Lines changed: 400 additions & 58 deletions

File tree

‎sdk/python/feast/feature_server.py‎

Lines changed: 231 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@
3131
from fastapi import (
3232
Depends,
3333
FastAPI,
34+
Query,
3435
Request,
3536
Response,
3637
WebSocket,
@@ -41,7 +42,7 @@
4142
from fastapi.logger import logger
4243
from fastapi.responses import JSONResponse
4344
from fastapi.staticfiles import StaticFiles
44-
from pydantic import BaseModel, field_validator
45+
from pydantic import BaseModel, Field, field_validator
4546

4647
import feast
4748
from feast import metrics as feast_metrics
@@ -54,6 +55,7 @@
5455
)
5556
from feast.feast_object import FeastObject
5657
from feast.feature_server_utils import convert_response_to_dict
58+
from feast.feature_view import FeatureViewState
5759
from feast.feature_view_utils import get_feature_view_from_feature_store
5860
from feast.filter_models import ComparisonFilter, CompoundFilter
5961
from feast.permissions.action import WRITE, AuthzedAction
@@ -94,12 +96,26 @@ class MaterializeRequest(BaseModel):
9496
feature_views: Optional[List[str]] = None
9597
disable_event_timestamp: bool = False
9698
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+
)
97106

98107

99108
class MaterializeIncrementalRequest(BaseModel):
100109
end_ts: str
101110
feature_views: Optional[List[str]] = None
102111
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+
)
103119

104120

105121
class GetOnlineFeaturesRequest(BaseModel):
@@ -333,6 +349,128 @@ async def load_static_artifacts(app: FastAPI, store):
333349
logger.warning(f"Failed to load static artifacts: {e}")
334350

335351

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+
336474
def get_app(
337475
store: "feast.FeatureStore",
338476
registry_ttl_sec: int = DEFAULT_FEATURE_SERVER_REGISTRY_TTL,
@@ -798,70 +936,115 @@ async def chat_ui():
798936
return Response(content=content, media_type="text/html")
799937

800938
@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+
):
802944
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)
819949

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+
)
831986

832987
await run_in_threadpool(
833988
store.materialize,
834989
start_date,
835990
end_date,
836-
request.feature_views,
991+
fv_names,
837992
disable_event_timestamp=request.disable_event_timestamp,
838993
full_feature_names=request.full_feature_names,
994+
version=request.version,
839995
)
840996

841997
@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+
):
8431003
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},
8541040
)
855-
for fv in feature_views_to_materialize:
856-
assert_permissions(
857-
resource=fv,
858-
actions=[AuthzedAction.WRITE_ONLINE],
859-
)
1041+
8601042
await run_in_threadpool(
8611043
store.materialize_incremental,
862-
utils.make_tzaware(parser.parse(request.end_ts)),
863-
request.feature_views,
1044+
end_date,
1045+
fv_names,
8641046
full_feature_names=request.full_feature_names,
1047+
version=request.version,
8651048
)
8661049

8671050
@app.exception_handler(Exception)

0 commit comments

Comments
 (0)