@@ -27,27 +27,6 @@ def run_demo():
27
27
print ("\n --- Online features retrieved (instead) through a feature service---" )
28
28
fetch_online_features (store , use_feature_service = True )
29
29
30
- print ("\n --- Simulate a stream event ingestion of the hourly stats df ---" )
31
- event_df = pd .DataFrame .from_dict (
32
- {
33
- "driver_id" : [1001 ],
34
- "event_timestamp" : [
35
- datetime (2021 , 5 , 13 , 10 , 59 , 42 ),
36
- ],
37
- "created" : [
38
- datetime (2021 , 5 , 13 , 10 , 59 , 42 ),
39
- ],
40
- "conv_rate" : [1.0 ],
41
- "acc_rate" : [1.0 ],
42
- "avg_daily_trips" : [1000 ],
43
- }
44
- )
45
- print (event_df )
46
- store .push ("driver_stats_push_source" , event_df , to = PushMode .ONLINE )
47
-
48
- print ("\n --- Online features again with updated values from a stream push---" )
49
- fetch_online_features (store , use_feature_service = True )
50
-
51
30
print ("\n --- Run feast teardown ---" )
52
31
subprocess .run (["feast" , "teardown" ])
53
32
@@ -92,17 +71,19 @@ def fetch_online_features(store, use_feature_service: bool):
92
71
# {join_key: entity_value}
93
72
{
94
73
"driver_id" : 1001 ,
74
+ "customer_id" : 201 ,
95
75
"val_to_add" : 1000 ,
96
76
"val_to_add_2" : 2000 ,
97
77
},
98
78
{
99
79
"driver_id" : 1002 ,
80
+ "customer_id" : 202 ,
100
81
"val_to_add" : 1001 ,
101
82
"val_to_add_2" : 2002 ,
102
83
},
103
84
]
104
85
if use_feature_service :
105
- features_to_fetch = store .get_feature_service ("driver_activity_v1 " )
86
+ features_to_fetch = store .get_feature_service ("driver_activity " )
106
87
else :
107
88
features_to_fetch = [
108
89
"driver_hourly_stats:acc_rate" ,
0 commit comments