Skip to content
Draft
Show file tree
Hide file tree
Changes from all 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
36 changes: 36 additions & 0 deletions sdk/python/feast/diff/registry_diff.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
from dataclasses import dataclass
from typing import Any, Dict, Iterable, List, Optional, Set, Tuple, TypeVar, cast

from google.protobuf.message import Message

from feast.base_feature_view import BaseFeatureView
from feast.data_source import DataSource
from feast.diff.property_diff import PropertyDiff, TransitionType
Expand Down Expand Up @@ -123,6 +125,38 @@ def tag_objects_for_keep_delete_update_add(
FIELDS_TO_IGNORE = {"project"}


def _clear_data_source_meta(message: Message) -> None:
"""Clears `meta` on every DataSource in `message`, including nested ones.

A data source's meta holds its created/last updated timestamps, which are
reset whenever the object is constructed, so comparing it would report a
change for every source embedded in an object that is being updated.
"""
if isinstance(message, DataSourceProto):
message.ClearField("meta")
for _, value in message.ListFields():
if isinstance(value, Message):
values = [value]
elif isinstance(value, (str, bytes)):
continue
elif hasattr(value, "values"): # map field
values = list(value.values())
elif hasattr(value, "__iter__"): # repeated field
values = list(value)
else:
continue
for v in values:
if isinstance(v, Message):
_clear_data_source_meta(v)


def _without_data_source_meta(spec: FeastObjectSpecProto) -> FeastObjectSpecProto:
stripped = type(spec)()
stripped.CopyFrom(spec)
_clear_data_source_meta(stripped)
return stripped


def diff_registry_objects(
current: FeastObject, new: FeastObject, object_type: FeastObjectType
) -> FeastObjectDiff:
Expand All @@ -143,6 +177,8 @@ def diff_registry_objects(
else:
current_spec = current_proto.spec
new_spec = new_proto.spec
current_spec = _without_data_source_meta(current_spec)
new_spec = _without_data_source_meta(new_spec)
if current != new:
for _field in current_spec.DESCRIPTOR.fields:
if _field.name in FIELDS_TO_IGNORE:
Expand Down
48 changes: 47 additions & 1 deletion sdk/python/tests/unit/diff/test_registry_diff.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
from datetime import datetime

import pandas as pd

from feast import Field, PushSource
from feast import Field, FileSource, PushSource
from feast.diff.registry_diff import (
diff_registry_objects,
tag_objects_for_keep_delete_update_add,
Expand Down Expand Up @@ -176,6 +178,50 @@ def test_diff_registry_objects_batch_to_push_source(simple_dataset_1):
)


def test_diff_registry_objects_ignores_data_source_meta():
# The registry copy and the freshly declared object each build their own
# source, so the sources' created/last updated timestamps always differ.
def source(path="data.parquet", timestamp_field="ts"):
return FileSource(name="src", path=path, timestamp_field=timestamp_field)

def feature_view(tags, source):
return FeatureView(
name="fv",
entities=[Entity(name="id", join_keys=["id"])],
source=PushSource(name="push", batch_source=source),
tags=tags,
)

registered = source()
registered.created_timestamp = datetime(2020, 1, 1)
registered.last_updated_timestamp = datetime(2020, 1, 2)

feast_object_diffs = diff_registry_objects(
feature_view({"when": "before"}, registered),
feature_view({"when": "after"}, source()),
"feature view",
)
assert [
d.property_name for d in feast_object_diffs.feast_object_property_diffs
] == ["tags"]

feast_object_diffs = diff_registry_objects(
feature_view({}, registered),
feature_view({}, source(path="other.parquet")),
"feature view",
)
assert [
d.property_name for d in feast_object_diffs.feast_object_property_diffs
] == ["batch_source", "stream_source"]

feast_object_diffs = diff_registry_objects(
registered, source(timestamp_field="other_ts"), "data source"
)
assert [
d.property_name for d in feast_object_diffs.feast_object_property_diffs
] == ["timestamp_field"]


def test_diff_registry_objects_permissions():
pre_changed = Permission(
name="reader",
Expand Down
Loading