Repository navigation
feat(spark): Add compute-on-read support for BatchFeatureView in get_… #6357
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
ntkathole
merged 7 commits into
feast-dev:master
from
SIDDHESH1564:feature/spark-bfv-compute-on-read
May 3, 2026
Merged
Changes from 1 commit
Commits
Show all changes
7 commits
Select commit
Hold shift + click to select a range
11d69be
feat(spark): add compute-on-read support for BatchFeatureView in get_…
SIDDHESH1564 c2a3d7a
Merge branch 'master' into feature/spark-bfv-compute-on-read
SIDDHESH1564 9390de9
Add None check for batch_source in _apply_bfv_transformations
SIDDHESH1564 a424bbf
refactor(spark): use feature_view_utils helpers for BFV transformatio…
SIDDHESH1564 c9e72de
fix(spark): apply time-range filter before BFV UDF and update temp vi…
SIDDHESH1564 502fad2
Merge branch 'master' into feature/spark-bfv-compute-on-read
SIDDHESH1564 e1e1764
refactor(spark): use resolve_feature_view_source_with_fallback for BF…
SIDDHESH1564 File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Next
Next commit
feat(spark): add compute-on-read support for BatchFeatureView in get_…
…historical_features Signed-off-by: Siddhesh Khairnar <khairnarsiddhesh4057@gmail.com>
- Loading branch information
commit 11d69be065e11c9e53c03201e6826535b6e1e8ef
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
188 changes: 188 additions & 0 deletions
188
...s/unit/infra/offline_stores/contrib/spark_offline_store/test_spark_bfv_compute_on_read.py
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,188 @@ | ||
| """ | ||
| Unit tests for BFV compute-on-read in SparkOfflineStore.get_historical_features(). | ||
|
|
||
| Verifies that BatchFeatureViews with a UDF have their transformation applied | ||
| during get_historical_features(), with the transformed DataFrame registered as | ||
| a temp view that replaces the raw table_subquery in the PIT join. | ||
| """ | ||
|
|
||
| from dataclasses import replace | ||
| from unittest.mock import MagicMock | ||
|
|
||
| import pytest | ||
|
|
||
| from feast.batch_feature_view import BatchFeatureView | ||
| from feast.feature_view import FeatureView | ||
| from feast.infra.offline_stores import offline_utils | ||
| from feast.infra.offline_stores.contrib.spark_offline_store.spark import ( | ||
| _apply_bfv_transformations, | ||
| ) | ||
| from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import ( | ||
| SparkSource, | ||
| ) | ||
| from feast.transformation.base import Transformation | ||
|
|
||
|
|
||
| @pytest.fixture() | ||
| def spark_session(): | ||
| mock = MagicMock() | ||
| mock.sql.return_value = MagicMock(name="source_df") | ||
| return mock | ||
|
|
||
|
|
||
| @pytest.fixture() | ||
| def spark_source(): | ||
| source = MagicMock(spec=SparkSource) | ||
| source.get_table_query_string.return_value = "`raw_events`" | ||
| return source | ||
|
|
||
|
|
||
| @pytest.fixture() | ||
| def base_query_context(): | ||
| return offline_utils.FeatureViewQueryContext( | ||
| name="my_bfv", | ||
| ttl=3600, | ||
| entities=["user_id"], | ||
| features=["avg_rating"], | ||
| field_mapping={}, | ||
| timestamp_field="event_timestamp", | ||
| created_timestamp_column=None, | ||
| table_subquery="`raw_events`", | ||
| entity_selections=["user_id AS user_id"], | ||
| min_event_timestamp="2023-01-01T00:00:00", | ||
| max_event_timestamp="2024-01-01T00:00:00", | ||
| date_partition_column=None, | ||
| ) | ||
|
|
||
|
|
||
| def _make_bfv(name: str, spark_source, has_udf: bool = True): | ||
| """Create a mock BatchFeatureView with optional UDF.""" | ||
| fv = MagicMock(spec=BatchFeatureView) | ||
| fv.name = name | ||
| fv.projection = MagicMock() | ||
| fv.projection.name_to_use.return_value = name | ||
| fv.batch_source = spark_source | ||
|
|
||
| if has_udf: | ||
| transformation = MagicMock(spec=Transformation) | ||
| transformed_df = MagicMock(name="transformed_df") | ||
| transformation.udf = MagicMock(return_value=transformed_df) | ||
| fv.feature_transformation = transformation | ||
| else: | ||
| fv.feature_transformation = None | ||
|
|
||
| return fv | ||
|
|
||
|
|
||
| def _make_plain_fv(name: str, spark_source): | ||
| """Create a mock plain FeatureView (not a BatchFeatureView).""" | ||
| fv = MagicMock(spec=FeatureView) | ||
| fv.name = name | ||
| fv.projection = MagicMock() | ||
| fv.projection.name_to_use.return_value = name | ||
| fv.batch_source = spark_source | ||
| return fv | ||
|
|
||
|
|
||
| class TestApplyBfvTransformations: | ||
| def test_bfv_with_udf_replaces_table_subquery( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """BFV with a UDF should have its table_subquery replaced with a temp view.""" | ||
| bfv = _make_bfv("my_bfv", spark_source) | ||
| contexts = [base_query_context] | ||
|
|
||
| result = _apply_bfv_transformations(spark_session, [bfv], contexts) | ||
|
|
||
| assert len(result) == 1 | ||
| assert result[0].table_subquery != "`raw_events`" | ||
| assert result[0].table_subquery.startswith("__feast_bfv_my_bfv_") | ||
|
|
||
| def test_bfv_udf_is_invoked_with_source_df( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """The UDF should be called with the DataFrame read from the raw source.""" | ||
| bfv = _make_bfv("my_bfv", spark_source) | ||
| contexts = [base_query_context] | ||
|
|
||
| _apply_bfv_transformations(spark_session, [bfv], contexts) | ||
|
|
||
| spark_session.sql.assert_called_once_with("SELECT * FROM `raw_events`") | ||
| source_df = spark_session.sql.return_value | ||
| bfv.feature_transformation.udf.assert_called_once_with(source_df) | ||
|
|
||
| def test_transformed_df_registered_as_temp_view( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """The transformed DataFrame should be registered as a temp view.""" | ||
| bfv = _make_bfv("my_bfv", spark_source) | ||
| transformed_df = bfv.feature_transformation.udf.return_value | ||
| contexts = [base_query_context] | ||
|
|
||
| result = _apply_bfv_transformations(spark_session, [bfv], contexts) | ||
|
|
||
| transformed_df.createOrReplaceTempView.assert_called_once() | ||
| view_name = transformed_df.createOrReplaceTempView.call_args[0][0] | ||
| assert view_name == result[0].table_subquery | ||
|
|
||
| def test_plain_feature_view_unchanged( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """Plain FeatureViews (not BFV) should pass through without modification.""" | ||
| fv = _make_plain_fv("my_bfv", spark_source) | ||
| contexts = [base_query_context] | ||
|
|
||
| result = _apply_bfv_transformations(spark_session, [fv], contexts) | ||
|
|
||
| assert result[0].table_subquery == "`raw_events`" | ||
| spark_session.sql.assert_not_called() | ||
|
|
||
| def test_bfv_without_udf_unchanged( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """BFV without a UDF should pass through without modification.""" | ||
| bfv = _make_bfv("my_bfv", spark_source, has_udf=False) | ||
| contexts = [base_query_context] | ||
|
|
||
| result = _apply_bfv_transformations(spark_session, [bfv], contexts) | ||
|
|
||
| assert result[0].table_subquery == "`raw_events`" | ||
| spark_session.sql.assert_not_called() | ||
|
|
||
| def test_mixed_views_only_transforms_bfvs( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """With mixed BFV and plain FVs, only BFVs with UDFs get transformed.""" | ||
| bfv = _make_bfv("my_bfv", spark_source) | ||
| plain_fv = _make_plain_fv("plain_fv", spark_source) | ||
|
|
||
| ctx_bfv = base_query_context | ||
| ctx_plain = replace( | ||
| base_query_context, | ||
| name="plain_fv", | ||
| features=["some_feature"], | ||
| ) | ||
|
|
||
| result = _apply_bfv_transformations( | ||
| spark_session, [bfv, plain_fv], [ctx_bfv, ctx_plain] | ||
| ) | ||
|
|
||
| assert result[0].table_subquery.startswith("__feast_bfv_my_bfv_") | ||
| assert result[1].table_subquery == "`raw_events`" | ||
|
|
||
| def test_other_context_fields_preserved( | ||
| self, spark_session, spark_source, base_query_context | ||
| ): | ||
| """All fields besides table_subquery should remain unchanged.""" | ||
| bfv = _make_bfv("my_bfv", spark_source) | ||
| contexts = [base_query_context] | ||
|
|
||
| result = _apply_bfv_transformations(spark_session, [bfv], contexts) | ||
|
|
||
| assert result[0].name == base_query_context.name | ||
| assert result[0].ttl == base_query_context.ttl | ||
| assert result[0].entities == base_query_context.entities | ||
| assert result[0].features == base_query_context.features | ||
| assert result[0].timestamp_field == base_query_context.timestamp_field | ||
| assert result[0].min_event_timestamp == base_query_context.min_event_timestamp | ||
| assert result[0].max_event_timestamp == base_query_context.max_event_timestamp |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.