Skip to content
Merged
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
7 changes: 6 additions & 1 deletion sdk/python/feast/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -883,7 +883,12 @@ def _augment_response_with_on_demand_transforms(
"""
from feast.online_response import OnlineResponse

requested_odfv_map = {odfv.name: odfv for odfv in requested_on_demand_feature_views}
# Feature refs name an ODFV by its alias when it was added to a
# FeatureService with ``with_name``, so key the lookup the same way.
requested_odfv_map = {
(odfv.projection.name_alias if odfv.projection else None) or odfv.name: odfv
for odfv in requested_on_demand_feature_views
}
requested_odfv_feature_names = requested_odfv_map.keys()

odfv_feature_refs = defaultdict(list)
Expand Down
58 changes: 58 additions & 0 deletions sdk/python/tests/unit/test_on_demand_pandas_transformation.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@

from feast import (
Entity,
FeatureService,
FeatureStore,
FeatureView,
FileSource,
Expand Down Expand Up @@ -597,3 +598,60 @@ def aliased_guard(inputs: pd.DataFrame) -> pd.DataFrame:
# the aliased-but-declared source must survive; the unrelated FV stays hidden
assert response["saw_declared"] == [True]
assert response["saw_undeclared"] == [False]


@pytest.mark.parametrize("mode", ["pandas", "python"])
def test_feature_service_serves_aliased_odfv_features(mode):
"""An ODFV added to a FeatureService under ``with_name`` must still be
computed. Its feature refs use the alias, so the transform lookup has to be
keyed by the alias too, not only by the ODFV's registered name."""
with tempfile.TemporaryDirectory() as data_dir:
store, driver, src, fv1, fv2, driver_df = _two_fv_store(data_dir)

if mode == "pandas":

@on_demand_feature_view(
sources=[fv1],
schema=[Field(name="conv_rate_x10", dtype=Float64)],
mode="pandas",
)
def scaled(inputs: pd.DataFrame) -> pd.DataFrame:
out = pd.DataFrame()
out["conv_rate_x10"] = inputs["conv_rate"] * 10
return out

else:

@on_demand_feature_view(
sources=[fv1],
schema=[Field(name="conv_rate_x10", dtype=Float64)],
mode="python",
)
def scaled(inputs: dict) -> dict:
return {"conv_rate_x10": [v * 10 for v in inputs["conv_rate"]]}

service = FeatureService(
name="aliased_service",
features=[fv1, scaled.with_name("scaled_alias")],
)
store.apply([driver, src, fv1, scaled, service])
store.write_to_online_store(feature_view_name="fv1", df=driver_df)

response = store.get_online_features(
entity_rows=[{"driver_id": 1001}],
features=store.get_feature_service("aliased_service"),
).to_dict()
assert response["conv_rate_x10"] == [
pytest.approx(response["conv_rate"][0] * 10)
]

if mode != "pandas":
return
response = store.get_online_features(
entity_rows=[{"driver_id": 1001}],
features=store.get_feature_service("aliased_service"),
full_feature_names=True,
).to_dict()
assert response["scaled_alias__conv_rate_x10"] == [
pytest.approx(response["fv1__conv_rate"][0] * 10)
]