Skip to content
Open
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
142 changes: 137 additions & 5 deletions docs/getting-started/quickstart.md
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,7 @@ entity_key_serialization_version: 3
# This is an example feature definition file

from datetime import timedelta
from typing import Any

import pandas as pd

Expand Down Expand Up @@ -194,7 +195,8 @@ input_request = RequestSource(


# Define an on demand feature view which can generate new features based on
# existing feature views and RequestSource features
# existing feature views and RequestSource features. The transformation runs at
# read time, every time the features are requested. Pandas mode is the default.
@on_demand_feature_view(
sources=[driver_stats_fv, input_request],
schema=[
Expand All @@ -209,6 +211,36 @@ def transformed_conv_rate(inputs: pd.DataFrame) -> pd.DataFrame:
return df


# The same transformation written in native Python mode (mode="python"). The UDF
# receives a dict mapping each input feature name to a list of values (one per
# row) and returns a dict with the same shape.
#
# Only the features the UDF needs are selected from the source feature view.
# This is required here: driver_stats_fv also has Map / Struct / Json fields, and
# Python mode feature inference cannot generate sample values for those types.
@on_demand_feature_view(
sources=[driver_stats_fv[["conv_rate"]], input_request],
schema=[
Field(name="conv_rate_plus_val1_python", dtype=Float64),
Field(name="conv_rate_plus_val2_python", dtype=Float64),
],
mode="python",
)
def transformed_conv_rate_python(inputs: dict[str, Any]) -> dict[str, Any]:
return {
"conv_rate_plus_val1_python": [
conv_rate + val_to_add
for conv_rate, val_to_add in zip(inputs["conv_rate"], inputs["val_to_add"])
],
"conv_rate_plus_val2_python": [
conv_rate + val_to_add_2
for conv_rate, val_to_add_2 in zip(
inputs["conv_rate"], inputs["val_to_add_2"]
)
],
}


# This groups features into a model version
driver_activity_v1 = FeatureService(
name="driver_activity_v1",
Expand Down Expand Up @@ -265,6 +297,32 @@ driver_activity_v3 = FeatureService(
name="driver_activity_v3",
features=[driver_stats_fresh_fv, transformed_conv_rate_fresh],
)


# Setting write_to_online_store=True runs the transformation at write time instead:
# the derived features are computed once when data is materialized or written to
# the online store, and are then served like any other pre-computed feature.
# Because the results are persisted, the view must declare its entities and can
# only depend on other feature views (not on request-time data).
@on_demand_feature_view(
entities=[driver],
sources=[driver_stats_fv[["conv_rate", "acc_rate", "avg_daily_trips"]]],
schema=[
Field(name="conv_rate_x_acc_rate", dtype=Float64),
Field(name="expected_daily_conversions", dtype=Float64),
],
mode="pandas",
write_to_online_store=True,
)
def transformed_conv_rate_on_write(inputs: pd.DataFrame) -> pd.DataFrame:
df = pd.DataFrame()
df["conv_rate_x_acc_rate"] = (inputs["conv_rate"] * inputs["acc_rate"]).astype(
"float64"
)
df["expected_daily_conversions"] = (
inputs["avg_daily_trips"] * inputs["conv_rate"]
).astype("float64")
return df
```
{% endtab %}
{% endtabs %}
Expand Down Expand Up @@ -330,13 +388,16 @@ Created entity driver
Created feature view driver_hourly_stats
Created feature view driver_hourly_stats_fresh
Created on demand feature view transformed_conv_rate
Created on demand feature view transformed_conv_rate_python
Created on demand feature view transformed_conv_rate_fresh
Created on demand feature view transformed_conv_rate_on_write
Created feature service driver_activity_v3
Created feature service driver_activity_v1
Created feature service driver_activity_v2

Created sqlite table my_project_driver_hourly_stats_fresh
Created sqlite table my_project_driver_hourly_stats
Created sqlite table my_project_transformed_conv_rate_on_write
```
{% endtab %}
{% endtabs %}
Expand Down Expand Up @@ -503,6 +564,8 @@ print(training_df.head())

We now serialize the latest values of features from the materialization window to prepare for serving. `materialize-incremental` starts from each feature view's most recent materialization end date. On the first run, it starts from the current time minus the feature view's `ttl`. In this example, the initial lookback is one day (`ttl` was set on the `FeatureView` instances in `feature_definitions.py`).

On demand feature views with `write_to_online_store=True` (here `transformed_conv_rate_on_write`) are materialized too: their source features are read from the offline store, the transformation is applied once, and the results are written to the online store.

{% tabs %}
{% tab title="Bash (with timestamp)" %}
```bash
Expand All @@ -522,13 +585,13 @@ feast materialize --disable-event-timestamp
{% tabs %}
{% tab title="Output" %}
```bash
Materializing 2 feature views to 2024-04-19 10:59:58-04:00 into the sqlite online store.
Materializing 3 feature views to 2024-04-19 10:59:58-04:00 into the sqlite online store.

driver_hourly_stats from 2024-04-18 15:00:46-04:00 to 2024-04-19 10:59:58-04:00:
100%|████████████████████████████████████████████████████████████████| 5/5 [00:00<00:00, 370.32it/s]
driver_hourly_stats_fresh from 2024-04-18 15:00:46-04:00 to 2024-04-19 10:59:58-04:00:
100%|███████████████████████████████████████████████████████████████| 5/5 [00:00<00:00, 1046.64it/s]
Materializing 2 feature views to 2024-04-19 10:59:58-04:00 into the sqlite online store.
transformed_conv_rate_on_write:
```
{% endtab %}
{% endtabs %}
Expand Down Expand Up @@ -628,7 +691,76 @@ pprint(feature_vector)
{% endtab %}
{% endtabs %}

## Step 9: Browse your features with the Web UI (experimental)
### Step 9: Transforming features on read vs. on write

The template defines on demand feature views in both flavors. `transformed_conv_rate` (Pandas mode) and
`transformed_conv_rate_python` (native Python mode) run their transformation at read time, so they can combine
stored features with request-time data such as `val_to_add` and `val_to_add_2`. Both produce the same values:

{% tabs %}
{% tab title="Python" %}
```python
feature_vector = store.get_online_features(
features=[
"transformed_conv_rate:conv_rate_plus_val1",
"transformed_conv_rate_python:conv_rate_plus_val1_python",
],
entity_rows=[
{"driver_id": 1004, "val_to_add": 1000, "val_to_add_2": 2000},
{"driver_id": 1005, "val_to_add": 1001, "val_to_add_2": 2002},
],
).to_dict()
pprint(feature_vector)
```
{% endtab %}
{% endtabs %}

{% tabs %}
{% tab title="Output" %}
```bash
{
'conv_rate_plus_val1': [1000.0459796078503, 1001.811998963356],
'conv_rate_plus_val1_python': [1000.0459796078503, 1001.811998963356],
'driver_id': [1004, 1005]
}
```
{% endtab %}
{% endtabs %}

`transformed_conv_rate_on_write` instead has `write_to_online_store=True`: its transformation already ran during
`materialize-incremental` (and runs again whenever new rows are written with `store.write_to_online_store(...)`), so
reading it is a plain lookup with no transformation at request time:

{% tabs %}
{% tab title="Python" %}
```python
feature_vector = store.get_online_features(
features=[
"transformed_conv_rate_on_write:conv_rate_x_acc_rate",
"transformed_conv_rate_on_write:expected_daily_conversions",
],
entity_rows=[{"driver_id": 1004}, {"driver_id": 1005}],
).to_dict()
pprint(feature_vector)
```
{% endtab %}
{% endtabs %}

{% tabs %}
{% tab title="Output" %}
```bash
{
'conv_rate_x_acc_rate': [0.005327459424734116, 0.4349033236503601],
'driver_id': [1004, 1005],
'expected_daily_conversions': [32.231705103069544, 391.3835003376007]
}
```
{% endtab %}
{% endtabs %}

See [On demand feature views](../reference/beta-on-demand-feature-view.md) for more details on the transformation modes and on write-time transformations.

## Step 10: Browse your features with the Web UI (experimental)

View all registered features, data sources, entities, and feature services with the Web UI.

Expand Down Expand Up @@ -660,7 +792,7 @@ INFO: Uvicorn running on http://0.0.0.0:8888 (Press CTRL+C to quit)

![](../reference/ui.png)

## Step 10: Re-examine `test_workflow.py`
## Step 11: Re-examine `test_workflow.py`
Take a look at `test_workflow.py` again. It showcases many sample flows on how to interact with Feast. You'll see these
show up in the upcoming concepts + architecture + tutorial pages as well.

Expand Down
14 changes: 14 additions & 0 deletions sdk/python/feast/infra/online_stores/sqlite.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
from feast.infra.online_stores.online_store import OnlineStore
from feast.infra.online_stores.vector_store import VectorStoreConfig
from feast.labeling.label_view import LabelView
from feast.on_demand_feature_view import OnDemandFeatureView
from feast.protos.feast.core.InfraObject_pb2 import InfraObject as InfraObjectProto
from feast.protos.feast.core.Registry_pb2 import Registry as RegistryProto
from feast.protos.feast.core.SqliteTable_pb2 import SqliteTable as SqliteTableProto
Expand Down Expand Up @@ -495,6 +496,19 @@ def plan(
)
)

# On demand feature views with write_to_online_store=True persist their
# transformed features, so they need a table just like regular feature views.
for odfv_proto in desired_registry_proto.on_demand_feature_views:
if odfv_proto.spec.write_to_online_store:
odfv = OnDemandFeatureView.from_proto(odfv_proto)
infra_objects.append(
SqliteTable(
path=self._get_db_path(config),
name=_table_id(project, odfv, versioning),
include_value_num=include_value_num,
)
)

return infra_objects

def teardown(
Expand Down
5 changes: 5 additions & 0 deletions sdk/python/feast/templates/local/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@ uses this repo. A quick view of what's in this repository's `feature_repo/` dire

* `data/` contains raw demo parquet data
* `feature_repo/feature_definitions.py` contains demo feature definitions
- `driver_hourly_stats` / `driver_hourly_stats_fresh`: regular feature views backed by a parquet file and a push source
- `transformed_conv_rate` / `transformed_conv_rate_python`: on demand feature views that transform features at
read time, written in Pandas mode and in native Python mode respectively
- `transformed_conv_rate_on_write`: an on demand feature view with `write_to_online_store=True`, whose
transformation runs when data is materialized or written to the online store instead of on every read
* `feature_repo/feature_store.yaml` contains a demo setup configuring where data sources are
* `feature_repo/test_workflow.py` showcases how to run all key Feast commands, including defining, retrieving, and pushing features.

Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
# This is an example feature definition file

from datetime import timedelta
from typing import Any

import pandas as pd

Expand Down Expand Up @@ -89,7 +90,8 @@


# Define an on demand feature view which can generate new features based on
# existing feature views and RequestSource features
# existing feature views and RequestSource features. By default the transformation
# runs in Pandas mode (mode="pandas"): the UDF receives and returns a DataFrame.
@on_demand_feature_view(
sources=[driver_stats_fv, input_request],
schema=[
Expand All @@ -104,6 +106,37 @@ def transformed_conv_rate(inputs: pd.DataFrame) -> pd.DataFrame:
return df


# The same transformation written in native Python mode (mode="python"). The UDF
# receives a dict mapping each input feature name to a list of values (one per
# row) and returns a dict with the same shape. This avoids the Pandas overhead
# for small online requests and is often easier to reason about.
#
# Only the features the UDF needs are selected from the source feature view.
# This is required here: driver_stats_fv also has Map / Struct / Json fields, and
# Python mode feature inference cannot generate sample values for those types.
@on_demand_feature_view(
sources=[driver_stats_fv[["conv_rate"]], input_request],
schema=[
Field(name="conv_rate_plus_val1_python", dtype=Float64),
Field(name="conv_rate_plus_val2_python", dtype=Float64),
],
mode="python",
)
def transformed_conv_rate_python(inputs: dict[str, Any]) -> dict[str, Any]:
return {
"conv_rate_plus_val1_python": [
conv_rate + val_to_add
for conv_rate, val_to_add in zip(inputs["conv_rate"], inputs["val_to_add"])
],
"conv_rate_plus_val2_python": [
conv_rate + val_to_add_2
for conv_rate, val_to_add_2 in zip(
inputs["conv_rate"], inputs["val_to_add_2"]
)
],
}


# This groups features into a model version
driver_activity_v1 = FeatureService(
name="driver_activity_v1",
Expand Down Expand Up @@ -168,6 +201,39 @@ def transformed_conv_rate_fresh(inputs: pd.DataFrame) -> pd.DataFrame:
features=[driver_stats_fresh_fv, transformed_conv_rate_fresh],
)


# The on demand feature views above run their transformation at read time, i.e.
# every time get_online_features() / get_historical_features() is called. Setting
# write_to_online_store=True instead runs the transformation at write time: the
# derived features are computed once when data is materialized or written to the
# online store, and are then served like any other pre-computed feature. This
# trades some ingestion cost for lower online retrieval latency.
#
# Because the results are persisted, the view must declare its entities and can
# only depend on other feature views (not on request-time data).
@on_demand_feature_view(
entities=[driver],
sources=[driver_stats_fv[["conv_rate", "acc_rate", "avg_daily_trips"]]],
schema=[
Field(name="conv_rate_x_acc_rate", dtype=Float64),
Field(name="expected_daily_conversions", dtype=Float64),
],
mode="pandas",
write_to_online_store=True,
)
def transformed_conv_rate_on_write(inputs: pd.DataFrame) -> pd.DataFrame:
df = pd.DataFrame()
# Cast explicitly so the output dtypes match the Float64 fields declared above
# (conv_rate and acc_rate are Float32 in the source feature view).
df["conv_rate_x_acc_rate"] = (inputs["conv_rate"] * inputs["acc_rate"]).astype(
"float64"
)
df["expected_daily_conversions"] = (
inputs["avg_daily_trips"] * inputs["conv_rate"]
).astype("float64")
return df


# --- Label Views ---
# Label views manage mutable human labels for training data, RLHF, and evaluation.
# They use PushSources so labels can be submitted from the UI or external tools.
Expand Down
Loading