Skip to content
Merged
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
Next Next commit
refactor: Normalize created timestamp on read, not in the join predicate
The cutoff predicate cast created_timestamp to UTC while the event-timestamp
comparison beside it did not. read_fv normalizes the event timestamp when it
reads the source and left the created timestamp alone, so the predicate was
compensating for a missing normalization at the one comparison site.

Normalize both on read instead. The predicate then needs no cast and matches
its neighbour. deduplicate() orders by created_timestamp_column regardless of
the cutoff flag, so the normalization is unconditional rather than gated on it;
casting a column that is already tz-aware compiles away, so this leaves the
emitted query unchanged for tz-aware sources and retires the "mutate only if
tz-naive" TODO.

Signed-off-by: David <david-adeniji@hotmail.co.uk>
  • Loading branch information
addenergyx committed Jul 31, 2026
commit adbdebeb28c8b98d3d9e42af5723af2ccf7f8f42
23 changes: 13 additions & 10 deletions sdk/python/feast/infra/offline_stores/ibis.py
Original file line number Diff line number Diff line change
Expand Up @@ -188,13 +188,18 @@ def read_fv(
fv_table = fv_table.rename({new_name: old_name})

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.
utc_columns = [timestamp_field]
if created_timestamp_column:
utc_columns.append(created_timestamp_column)

# TODO mutate only if tz-naive
fv_table = fv_table.mutate(
**{
timestamp_field: fv_table[timestamp_field].cast(
dt.Timestamp(timezone="UTC")
)
column: fv_table[column].cast(dt.Timestamp(timezone="UTC"))
for column in utc_columns
}
)

Expand All @@ -220,8 +225,8 @@ def read_fv(

return (
fv_table,
feature_view.batch_source.timestamp_field,
feature_view.batch_source.created_timestamp_column,
timestamp_field,
created_timestamp_column,
feature_view.projection.join_key_map
or {e.name: e.name for e in feature_view.entity_columns},
feature_refs,
Expand Down Expand Up @@ -439,10 +444,8 @@ def point_in_time_join(

if filter_by_created_timestamp and created_timestamp_field:

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.

This predicate casts the created timestamp to UTC but the existing event-timestamp comparison on the same join doesn't cast. Compare with the existing predicate at line 433-434 which does no cast.

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 @ntkathole, made a change so read_fv now normalizes the created timestamp on read, as it already did for the event timestamp, so the predicate no longer casts. Made it unconditional since deduplicate() reads that column regardless of the flag.

predicates.append(
feature_table[created_timestamp_field].cast(
dt.Timestamp(timezone="UTC")
)
<= entity_table[event_timestamp_col]
feature_table[created_timestamp_field]
<= entity_table[event_timestamp_col],
)

if ttl:
Expand Down
Loading