@@ -77,72 +77,72 @@ def simple_sfv(df):
7777 assert features ["dummy_field" ] == [None ]
7878
7979
80- @pytest .mark .integration
81- def test_stream_feature_view_udf (simple_dataset_1 ) -> None :
82- """
83- Test apply of StreamFeatureView udfs are serialized correctly and usable.
84- """
85- runner = CliRunner ()
86- with runner .local_repo (
87- get_example_repo ("example_feature_repo_1.py" ), "bigquery"
88- ) as fs , prep_file_source (
89- df = simple_dataset_1 , timestamp_field = "ts_1"
90- ) as file_source :
91- entity = Entity (name = "driver_entity" , join_keys = ["test_key" ])
92-
93- stream_source = KafkaSource (
94- name = "kafka" ,
95- timestamp_field = "event_timestamp" ,
96- kafka_bootstrap_servers = "" ,
97- message_format = AvroFormat ("" ),
98- topic = "topic" ,
99- batch_source = file_source ,
100- watermark_delay_threshold = timedelta (days = 1 ),
101- )
102-
103- @stream_feature_view (
104- entities = [entity ],
105- ttl = timedelta (days = 30 ),
106- owner = "test@example.com" ,
107- online = True ,
108- schema = [Field (name = "dummy_field" , dtype = Float32 )],
109- description = "desc" ,
110- aggregations = [
111- Aggregation (
112- column = "dummy_field" , function = "max" , time_window = timedelta (days = 1 ),
113- ),
114- Aggregation (
115- column = "dummy_field2" ,
116- function = "count" ,
117- time_window = timedelta (days = 24 ),
118- ),
119- ],
120- timestamp_field = "event_timestamp" ,
121- mode = "spark" ,
122- source = stream_source ,
123- tags = {},
124- )
125- def pandas_view (pandas_df ):
126- import pandas as pd
127-
128- assert type (pandas_df ) == pd .DataFrame
129- df = pandas_df .transform (lambda x : x + 10 , axis = 1 )
130- df .insert (2 , "C" , [20.2 , 230.0 , 34.0 ], True )
131- return df
132-
133- import pandas as pd
134-
135- fs .apply ([entity , pandas_view ])
136-
137- stream_feature_views = fs .list_stream_feature_views ()
138- assert len (stream_feature_views ) == 1
139- assert stream_feature_views [0 ] == pandas_view
140-
141- sfv = stream_feature_views [0 ]
142-
143- df = pd .DataFrame ({"A" : [1 , 2 , 3 ], "B" : [10 , 20 , 30 ]})
144- new_df = sfv .udf (df )
145- expected_df = pd .DataFrame (
146- {"A" : [11 , 12 , 13 ], "B" : [20 , 30 , 40 ], "C" : [20.2 , 230.0 , 34.0 ]}
147- )
148- assert new_df .equals (expected_df )
80+ # @pytest.mark.integration
81+ # def test_stream_feature_view_udf(simple_dataset_1) -> None:
82+ # """
83+ # Test apply of StreamFeatureView udfs are serialized correctly and usable.
84+ # """
85+ # runner = CliRunner()
86+ # with runner.local_repo(
87+ # get_example_repo("example_feature_repo_1.py"), "bigquery"
88+ # ) as fs, prep_file_source(
89+ # df=simple_dataset_1, timestamp_field="ts_1"
90+ # ) as file_source:
91+ # entity = Entity(name="driver_entity", join_keys=["test_key"])
92+
93+ # stream_source = KafkaSource(
94+ # name="kafka",
95+ # timestamp_field="event_timestamp",
96+ # kafka_bootstrap_servers="",
97+ # message_format=AvroFormat(""),
98+ # topic="topic",
99+ # batch_source=file_source,
100+ # watermark_delay_threshold=timedelta(days=1),
101+ # )
102+
103+ # @stream_feature_view(
104+ # entities=[entity],
105+ # ttl=timedelta(days=30),
106+ # owner="test@example.com",
107+ # online=True,
108+ # schema=[Field(name="dummy_field", dtype=Float32)],
109+ # description="desc",
110+ # aggregations=[
111+ # Aggregation(
112+ # column="dummy_field", function="max", time_window=timedelta(days=1),
113+ # ),
114+ # Aggregation(
115+ # column="dummy_field2",
116+ # function="count",
117+ # time_window=timedelta(days=24),
118+ # ),
119+ # ],
120+ # timestamp_field="event_timestamp",
121+ # mode="spark",
122+ # source=stream_source,
123+ # tags={},
124+ # )
125+ # def pandas_view(pandas_df):
126+ # import pandas as pd
127+
128+ # assert type(pandas_df) == pd.DataFrame
129+ # df = pandas_df.transform(lambda x: x + 10, axis=1)
130+ # df.insert(2, "C", [20.2, 230.0, 34.0], True)
131+ # return df
132+
133+ # import pandas as pd
134+
135+ # fs.apply([entity, pandas_view])
136+
137+ # stream_feature_views = fs.list_stream_feature_views()
138+ # assert len(stream_feature_views) == 1
139+ # assert stream_feature_views[0] == pandas_view
140+
141+ # sfv = stream_feature_views[0]
142+
143+ # df = pd.DataFrame({"A": [1, 2, 3], "B": [10, 20, 30]})
144+ # new_df = sfv.udf(df)
145+ # expected_df = pd.DataFrame(
146+ # {"A": [11, 12, 13], "B": [20, 30, 40], "C": [20.2, 230.0, 34.0]}
147+ # )
148+ # assert new_df.equals(expected_df)
0 commit comments