Skip to content
Merged
Show file tree
Hide file tree
Changes from 10 commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
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
311 changes: 180 additions & 131 deletions sdk/python/feast/feature_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -424,6 +424,109 @@ def _get_provider(self) -> Provider:
# TODO: Bake self.repo_path into self.config so that we dont only have one interface to paths

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] Race condition in state management

The state transition and rollback logic has potential race conditions when multiple materialization jobs run concurrently. The previous_states dictionary and feature view state updates are not atomic, which could lead to inconsistent states or lost rollbacks.

Suggested:

Suggested change
# TODO: Bake self.repo_path into self.config so that we dont only have one interface to paths
# Use locks or atomic operations for state management
import threading
state_lock = threading.Lock()
with state_lock:
for fv, job in zip(regular_fvs, jobs):
fv_status = job.status()
if fv_status == MaterializationJobStatus.ERROR:
failed_fvs.append(fv)
if first_error is None and job.error():
first_error = job.error()
else:
succeeded_fvs.append(fv)
if failed_fvs:
self._rollback_fv_states(failed_fvs, previous_states)

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.

The threading race condition described isn't present because previous_states is a local variable (line 2551) and all loops (for fv, job in zip(...)) are sequential Python for loops, not threaded. Each materialize() call creates its own dict.

However, reviewing this code path revealed a shared-object bug in our engine's _build_per_fv_jobs. Previously, all succeeded FVs received the same SparkApplicationMaterializationJob reference. During _wait_for_completion, polling sets _error on that object when the SparkApp transitions to FAILED. Later, when feature_store.py calls job.status() for each FV, they all hit the if self._error is not None: return ERROR early-return, even for the FVs that actually succeeded and wrote AVAILABLE_ONLINE to the registry.

return self.provider

def _rollback_fv_states(
self,
feature_views: list,
previous_states: dict,
) -> None:
"""Restore feature views to their pre-materialization states."""
for fv in feature_views:
prev = previous_states.get(fv.name)
if (
hasattr(fv, "state")
and prev is not None
and prev != FeatureViewState.STATE_UNSPECIFIED
):
fv.state = prev
self.registry.apply_feature_view(fv, self.project, commit=True)

def _transition_fv_to_materializing(
self,
feature_view,
already_transitioned: list,
previous_states: dict,
) -> None:
"""
Transition a feature view to MATERIALIZING state.

Rolls back all already-transitioned FVs if this one can't transition.
"""
previous_state = getattr(feature_view, "state", None)
if (
hasattr(feature_view, "state")
and feature_view.state != FeatureViewState.STATE_UNSPECIFIED
):
if not feature_view.state.can_transition_to(FeatureViewState.MATERIALIZING):
self._rollback_fv_states(already_transitioned, previous_states)
raise ValueError(
f"FeatureView {feature_view.name} cannot transition "
f"from {feature_view.state.name} 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,
provider,
tasks: list,
regular_fvs: list,
previous_states: dict,
fv_start_dates: dict,
) -> None:
"""
Submit all tasks to the engine in one call and process the results.

For each returned job: record watermark on success, roll back state on
error. If the engine itself raises, all states are rolled back.
"""
from feast.infra.common.materialization_job import (
MaterializationJobStatus,
)

batch_start = time.monotonic()
try:
jobs = provider.batch_engine.materialize(self.registry, tasks)
except Exception:
self._rollback_fv_states(regular_fvs, previous_states)
raise

first_error = None
succeeded_fvs = []
failed_fvs = []

for fv, job in zip(regular_fvs, jobs):
Comment thread
aniketpalu marked this conversation as resolved.
fv_status = job.status()

if fv_status == MaterializationJobStatus.ERROR:
failed_fvs.append(fv)
if first_error is None and job.error():
first_error = job.error()
else:
succeeded_fvs.append(fv)

if failed_fvs:
self._rollback_fv_states(failed_fvs, previous_states)

for fv in succeeded_fvs:
self.registry.apply_materialization(
fv,
self.project,
fv_start_dates[fv.name],
fv_start_dates["__end_date__"],
)

_tracker = _get_track_materialization()
if _tracker is not None:
elapsed = time.monotonic() - batch_start
for fv in succeeded_fvs:
_tracker(fv.name, True, elapsed)
for fv in failed_fvs:
_tracker(fv.name, False, elapsed)

if first_error:
raise first_error

@property
def openlineage_emitter(self) -> Optional[Any]:
"""Gets the OpenLineage emitter of this feature store."""
Expand Down Expand Up @@ -2249,7 +2352,21 @@ def materialize_incremental(

_mat_start = time.monotonic()
try:
# TODO paging large loads
from feast.infra.common.materialization_job import (
MaterializationTask,
)

provider = self._get_provider()
end_date_tz = utils.make_tzaware(end_date) or _utc_now()

def tqdm_builder(length):
return tqdm(total=length, ncols=100)

tasks: list = []
regular_fvs: list = []
previous_states: dict = {}
fv_start_dates: dict = {"__end_date__": end_date_tz}
Comment thread
aniketpalu marked this conversation as resolved.
Outdated

for feature_view in feature_views_to_materialize:
if isinstance(feature_view, OnDemandFeatureView):
if feature_view.write_to_online_store:
Expand Down Expand Up @@ -2278,96 +2395,55 @@ def materialize_incremental(
)
continue

start_date = feature_view.most_recent_end_time
if start_date is None:
fv_start_date = feature_view.most_recent_end_time
if fv_start_date is None:
if feature_view.ttl is None:
raise Exception(
f"No start time found for feature view {feature_view.name}. materialize_incremental() requires"
f" either a ttl to be set or for materialize() to have been run at least once."
)
elif feature_view.ttl.total_seconds() > 0:
start_date = _utc_now() - feature_view.ttl
fv_start_date = _utc_now() - feature_view.ttl
else:
# TODO(felixwang9817): Find the earliest timestamp for this specific feature
# view from the offline store, and set the start date to that timestamp.
print(
f"Since the ttl is 0 for feature view {Style.BRIGHT + Fore.GREEN}{feature_view.name}{Style.RESET_ALL}, "
"the start date will be set to 1 year before the current time."
)
start_date = _utc_now() - timedelta(weeks=52)
provider = self._get_provider()
fv_start_date = _utc_now() - timedelta(weeks=52)

fv_start_date = utils.make_tzaware(fv_start_date)
fv_start_dates[feature_view.name] = fv_start_date

print(
f"{Style.BRIGHT + Fore.GREEN}{feature_view.name}{Style.RESET_ALL}"
f" from {Style.BRIGHT + Fore.GREEN}{utils.make_tzaware(start_date.replace(microsecond=0))}{Style.RESET_ALL}"
f" from {Style.BRIGHT + Fore.GREEN}{utils.make_tzaware(fv_start_date.replace(microsecond=0))}{Style.RESET_ALL}"
f" to {Style.BRIGHT + Fore.GREEN}{utils.make_tzaware(end_date.replace(microsecond=0))}{Style.RESET_ALL}:"
)

def tqdm_builder(length):
return tqdm(total=length, ncols=100)

start_date = utils.make_tzaware(start_date)
end_date = utils.make_tzaware(end_date) or _utc_now()

# 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")
and feature_view.state != FeatureViewState.STATE_UNSPECIFIED
):
if not feature_view.state.can_transition_to(
FeatureViewState.MATERIALIZING
):
raise ValueError(
f"FeatureView {feature_view.name} cannot transition "
f"from {feature_view.state.name} to MATERIALIZING."
)
feature_view.state = FeatureViewState.MATERIALIZING
self.registry.apply_feature_view(
feature_view, self.project, commit=True
)

fv_start = time.monotonic()
fv_success = True
try:
provider.materialize_single_feature_view(
config=self.config,
feature_view=feature_view,
start_date=start_date,
end_date=end_date,
registry=self.registry,
self._transition_fv_to_materializing(
feature_view, regular_fvs, previous_states
)
regular_fvs.append(feature_view)
tasks.append(
MaterializationTask(
project=self.project,
feature_view=feature_view,
start_time=fv_start_date,
end_time=end_date_tz,
tqdm_builder=tqdm_builder,
)
except Exception:
fv_success = False
# Roll back state to previous value on failure.
if (
hasattr(feature_view, "state")
and previous_state is not None
and previous_state != FeatureViewState.STATE_UNSPECIFIED
):
feature_view.state = previous_state
self.registry.apply_feature_view(
feature_view, self.project, commit=True
)
raise
finally:
_tracker = _get_track_materialization()
if _tracker is not None:
_tracker(
feature_view.name,
fv_success,
time.monotonic() - fv_start,
)
)

if not isinstance(feature_view, OnDemandFeatureView):
self.registry.apply_materialization(
feature_view,
self.project,
start_date,
end_date,
)
if tasks:
self._submit_and_process_materialization_jobs(
provider,
tasks,
regular_fvs,
previous_states,
fv_start_dates,
)

materialized_fv_names = [
fv.name
Expand Down Expand Up @@ -2459,7 +2535,22 @@ def materialize(

_mat_start = time.monotonic()
try:
# TODO paging large loads
from feast.infra.common.materialization_job import (
MaterializationTask,
)

provider = self._get_provider()
start_date = utils.make_tzaware(start_date)
end_date = utils.make_tzaware(end_date)

def tqdm_builder(length):
return tqdm(total=length, ncols=100)

tasks: list = []
regular_fvs: list = []
previous_states: dict = {}
fv_start_dates: dict = {"__end_date__": end_date}

for feature_view in feature_views_to_materialize:
if isinstance(feature_view, OnDemandFeatureView):
if feature_view.write_to_online_store:
Expand All @@ -2473,76 +2564,34 @@ def materialize(
full_feature_names=full_feature_names,
)
continue
provider = self._get_provider()

print(
f"{Style.BRIGHT + Fore.GREEN}{feature_view.name}{Style.RESET_ALL}:"
)

def tqdm_builder(length):
return tqdm(total=length, ncols=100)

start_date = utils.make_tzaware(start_date)
end_date = utils.make_tzaware(end_date)

# 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")
and feature_view.state != FeatureViewState.STATE_UNSPECIFIED
):
if not feature_view.state.can_transition_to(
FeatureViewState.MATERIALIZING
):
raise ValueError(
f"FeatureView {feature_view.name} cannot transition "
f"from {feature_view.state.name} to MATERIALIZING."
)
feature_view.state = FeatureViewState.MATERIALIZING
self.registry.apply_feature_view(
feature_view, self.project, commit=True
)

fv_start = time.monotonic()
fv_success = True
try:
provider.materialize_single_feature_view(
config=self.config,
feature_view=feature_view,
start_date=start_date,
end_date=end_date,
registry=self.registry,
self._transition_fv_to_materializing(
feature_view, regular_fvs, previous_states
)
regular_fvs.append(feature_view)
fv_start_dates[feature_view.name] = start_date
tasks.append(
MaterializationTask(
project=self.project,
feature_view=feature_view,
start_time=start_date,
end_time=end_date,
tqdm_builder=tqdm_builder,
disable_event_timestamp=disable_event_timestamp,
)
except Exception:
fv_success = False
# Roll back state to previous value on failure.
if (
hasattr(feature_view, "state")
and previous_state is not None
and previous_state != FeatureViewState.STATE_UNSPECIFIED
):
feature_view.state = previous_state
self.registry.apply_feature_view(
feature_view, self.project, commit=True
)
raise
finally:
_tracker = _get_track_materialization()
if _tracker is not None:
_tracker(
feature_view.name,
fv_success,
time.monotonic() - fv_start,
)
)

self.registry.apply_materialization(
feature_view,
self.project,
start_date,
end_date,
if tasks:
self._submit_and_process_materialization_jobs(
provider,
tasks,
regular_fvs,
previous_states,
fv_start_dates,
)

materialized_fv_names = [
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
feast/.git
feast/__pycache__
feast/**/__pycache__
feast/.mypy_cache
feast/tests
feast/.pixi
**/*.pyc
**/.pytest_cache
Loading
Loading