Skip to content
Merged
Show file tree
Hide file tree
Changes from 6 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
244 changes: 199 additions & 45 deletions sdk/python/feast/feature_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
from fastapi import (
Depends,
FastAPI,
Query,
Request,
Response,
WebSocket,
Expand All @@ -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
Expand Down Expand Up @@ -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):
Comment on lines 97 to 122

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Suggestion] Missing version parameter documentation and validation

The version parameter was added to both MaterializeRequest and MaterializeIncrementalRequest but there's no documentation about what values are valid or how it affects the materialization behavior.

Suggested:

Suggested change
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):
version: Optional[str] = Field(
None,
description="Optional version to materialize (e.g., 'v2'). Requires feature_views with exactly one entry."
)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done

Expand Down Expand Up @@ -333,6 +337,121 @@ 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(
Comment on lines +397 to +409

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] Silent exception handling could hide real errors

The _reset_stuck_materializing_to_generated and _check_already_materializing functions silently ignore all exceptions when accessing feature views. This could hide legitimate errors like registry corruption or network issues, making debugging difficult.

Suggested:

Suggested change
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(
except (FeatureViewNotFoundException, KeyError):
# Expected when FV doesn't exist
pass
except Exception as e:
logger.warning(f"Unexpected error checking state for {fv_name}: {e}")
pass

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Accepted

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 _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:
"""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,
Expand Down Expand Up @@ -798,36 +917,49 @@ 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 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))
if async_mode:
if force:
_reset_stuck_materializing_to_generated(store, fv_names)
else:
conflict = _check_already_materializing(store, fv_names)
if conflict:
return conflict

# State transitions (MATERIALIZING / AVAILABLE_ONLINE) are owned
# by store.materialize(); server only accepts and runs in background.
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}",

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Critical] Background task exception handling insufficient

The async materialization runs in a background thread but doesn't handle critical failures properly. If materialization fails, only the FV state is reset to GENERATED, but there's no way for the client to know the operation failed since the API already returned 202. This could lead to silent failures in production.

Suggested:

Suggested change
# by store.materialize(); server only accepts and runs in background.
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}",
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)
# TODO: Consider implementing a status endpoint or webhook callback
# for clients to check materialization status

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks — agreed that fire-and-forget 202 means clients don’t get the failure on the HTTP response. That’s intentional for this PR: async acceptance + observe completion via FeatureView.state / materialization_intervals (and run_async=False when the client needs the sync HTTP path to block/error).

A dedicated status endpoint / job handle / webhook is out of scope here and planned as a follow-up (along the RemoteComputeEngine / job-tracking direction discussed earlier). I’ve added a TODO in the failure path comment pointing to that follow-up.

exc_info=True,
)
_update_fv_state(store, fv_names, FeatureViewState.GENERATED)

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,
Expand All @@ -839,27 +971,49 @@ async def materialize(request: MaterializeRequest) -> None:
)

Comment on lines 998 to 1016

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Nitpick] Inconsistent parameter passing in sync mode

In the synchronous code path for materialize, the version parameter is passed but the feature_views parameter uses the original request.feature_views instead of the resolved fv_names list.

Suggested:

Suggested change
await run_in_threadpool(
store.materialize,
start_date,
end_date,
fv_names, # Use resolved names for consistency
disable_event_timestamp=request.disable_event_timestamp,
full_feature_names=request.full_feature_names,
version=request.version, # Don't forget version parameter
)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sync path now passes resolved fv_names and version=request.version (aligned with async).

@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 force:
_reset_stuck_materializing_to_generated(store, fv_names)
else:
conflict = _check_already_materializing(store, fv_names)
if conflict:
return conflict

# State transitions owned by store.materialize_incremental().
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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

non-blocking but I think it's better if we used dedicated executor, instead of default shared pool, for materialization so that it won't block other operations if multiple executors in progress

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed, added separate executor


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,
)
Expand Down
Loading
Loading