Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
Shorten comments in the created timestamp cutoff
Signed-off-by: David <david-adeniji@hotmail.co.uk>
  • Loading branch information
addenergyx committed Jul 31, 2026
commit 3d2ea3e0d43d18f0606659d1132402be640f8412
11 changes: 5 additions & 6 deletions sdk/python/feast/infra/offline_stores/dask.py
Original file line number Diff line number Diff line change
Expand Up @@ -1207,17 +1207,16 @@ def _apply_created_timestamp_cutoff(
entity_df_event_timestamp_col: str,
preserved_columns: Set[str],
) -> dd.DataFrame:
# Blanking rather than dropping preserves entity-dataframe cardinality when every
# candidate is too new. A blanked row looks like an unmatched left join, which
# _drop_duplicates already resolves (nulls sort first, keep="last"). The isna() term
# keeps a matched row whose source timestamp is null, as _filter_ttl does.
# Versions created after the entity timestamp. The isna() term leaves rows with a
# null source timestamp untouched, matching _filter_ttl.
too_new = ~df_to_join[timestamp_field].isna() & ~(
df_to_join[created_timestamp_column]
<= df_to_join[entity_df_event_timestamp_col]
)

# One assign, not a per-column loop: chained assignments make optimization
# super-linear in column count. Lazy so it fuses into the next persist.
# Blank instead of drop, so the entity row survives when every candidate is too new.
# _drop_duplicates then prefers a real match (nulls sort first, keep="last").
# Single assign, left lazy: a per-column loop is super-linear to optimize.
return df_to_join.assign(
**{
column: df_to_join[column].mask(too_new)
Expand Down
4 changes: 2 additions & 2 deletions sdk/python/feast/infra/offline_stores/ibis.py
Original file line number Diff line number Diff line change
Expand Up @@ -190,8 +190,8 @@ def read_fv(
timestamp_field = feature_view.batch_source.timestamp_field
created_timestamp_column = feature_view.batch_source.created_timestamp_column

# deduplicate() orders by created_timestamp_column regardless of the cutoff, so
# normalize both unconditionally. An already-tz-aware cast compiles away.
# deduplicate() orders by created_timestamp_column whether or not the cutoff is
# on, so this cannot be gated on it. Casting an already-UTC column is a no-op.
utc_columns = [timestamp_field]
if created_timestamp_column:
utc_columns.append(created_timestamp_column)
Expand Down
Loading