Vues de fonctionnalités
Aperçu
Cette fonctionnalité est en Aperçu public. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité à partir de la page Previews . Consultez Gérer les aperçus Databricks.
Les vues de fonctionnalités vous permettent de définir et de compute des fonctionnalités à partir de sources de données. Les fonctionnalités peuvent être définies à l'aide de diverses sources (table Delta, Kafka Stream et données au moment de la requête) et de calculs (agrégations par fenêtre temporelle, sélections de colonnes simples, et plus encore). Ce guide couvre les workflows suivants :
-
**Workflow de développement de fonctionnalités**
- Utilisez
create_featurepour définir des objets de fonctionnalité Unity Catalog qui peuvent être utilisés dans les workflows d'entraînement et de service de modèles. - Vous pouvez également construire localement des objets
Featureet utiliserregister_featurepour les conserver ultérieurement dans Unity Catalog. Les fonctionnalités construites localement peuvent être utilisées aveccreate_training_setavant l'enregistrement.
- Utilisez
-
Workflow de Model training
- Utilisez
create_training_setpour calculer des fonctionnalités agrégées ponctuelles pour le Machine Learning. Pour une documentation détaillée sur l'entraînement avec Feature Views, consultez Entraîner des modèles avec Feature Views.
- Utilisez
-
Workflow de matérialisation et de service des fonctionnalités
- Après avoir défini une fonctionnalité avec
create_featureou l'avoir récupérée à l'aide deget_feature, vous pouvez utilisermaterialize_featurespour matérialiser la fonctionnalité ou l'ensemble de fonctionnalités dans un magasin hors ligne pour une réutilisation efficace, ou dans un magasin en ligne pour la diffusion en ligne. - Utilisez
create_training_setavec la vue matérialisée pour préparer un dataset d'entraînement batch hors ligne.
- Après avoir défini une fonctionnalité avec
Pour les détails de l'API, consultez la référence de l'API des vues de fonctionnalités.
Exigences
-
Serverless compute ou un cluster de calcul classique exécutant Databricks Runtime 17.0 ML ou supérieur.
-
Vous devez installer le package Python personnalisé. Exécutez les lignes de code suivantes chaque fois que vous exécutez un Notebook :
Python%pip install databricks-feature-engineering>=0.16.0
dbutils.library.restartPython()
Exemple de démarrage rapide
Pour un notebook de démarrage rapide exécutable, consultez Exemple de notebook.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
CronSchedule, DeltaTableSource, Feature, AggregationFunction,
Sum, Avg, ColumnSelection, TableTrigger,
TumblingWindow, SlidingWindow,
OfflineStoreConfig, OnlineStoreConfig,
)
from datetime import timedelta
CATALOG_NAME = "main"
SCHEMA_NAME = "feature_store"
TABLE_NAME = "transactions"
# 1. Create data source
source = DeltaTableSource(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name=TABLE_NAME,
)
# 2. Define features locally (no catalog/schema needed yet)
avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), TumblingWindow(window_duration=timedelta(days=30))),
name="avg_transaction_30d",
)
sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), SlidingWindow(window_duration=timedelta(days=7), slide_duration=timedelta(days=1))),
# name auto-generated: "amount_sum_sliding_7d_1d"
)
fe = FeatureEngineeringClient()
# 3. Explore features with compute_features
feature_df = fe.compute_features(features=[avg_feature, sum_feature])
feature_df.display()
# 4. Create training set using local features
# `labeled_df` should have columns "user_id", "transaction_time", and "target".
training_set = fe.create_training_set(
df=labeled_df,
features=[avg_feature, sum_feature],
label="target",
)
training_set.load_df().display()
# 5. Register features in Unity Catalog
avg_feature = fe.register_feature(
feature=avg_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)
sum_feature = fe.register_feature(
feature=sum_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)
# 6. Or use create_feature for a one-step define-and-register workflow
latest_amount = fe.create_feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
name="latest_amount",
)
# 7. Train model
with mlflow.start_run():
training_df = training_set.load_df()
# training code
fe.log_model(
model=model,
artifact_path="recommendation_model",
flavor=mlflow.sklearn,
training_set=training_set,
registered_model_name=f"{CATALOG_NAME}.{SCHEMA_NAME}.recommendation_model",
)
# 8. (Optional) Materialize features for serving
# Features must be registered in UC before calling materialize_features
online_config = OnlineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features_serving",
online_store_name="customer_features_store",
)
# Aggregation features use CronSchedule and support both offline and online configs
fe.materialize_features(
features=[avg_feature, sum_feature],
offline_config=OfflineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features",
),
online_config=online_config,
trigger=CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
),
)
# ColumnSelection features use TableTrigger and only support online config
fe.materialize_features(
features=[latest_amount],
online_config=online_config,
trigger=TableTrigger(),
)
Exemple de Notebook
Notebook de démarrage rapide de Feature Views
Fonctionnalités de streaming
Outre les fonctionnalités de batch des tables Delta, vous pouvez définir des fonctionnalités à partir de sources de streaming pour les cas d'utilisation en temps réel. Les fonctionnalités de streaming utilisent la même classe de caractéristiques que les fonctionnalités de batch — mêmes constructeurs Feature, mêmes fonctions d'agrégation, mêmes workflows d'entraînement et de service — de sorte que la mise à niveau du traitement par batch vers le temps réel nécessite des modifications de code minimales. Une fois matérialisées, les fonctionnalités de streaming offrent une actualisation de bout en bout en moins d'une seconde (latence P99 de 200 ms) directement à vos endpoints de service de modèle.
Pour utiliser les fonctionnalités de streaming, commencez par configurer un Stream, puis référencez-le à l'aide d'un StreamSource. Les sources de Stream prennent en charge Kafka comme entrée et maintiennent automatiquement une table d'ingestion (Delta) en tant que copie historique des données pour la formation.
Définir une fonctionnalité de streaming
Un StreamSource référence un Stream par son nom en trois parties (catalog.schema.stream_name). Un Stream n'est pas un objet sécurisable d'Unity Catalog, mais il est associé à un schéma Unity Catalog et son accès est régi par la table d'ingestion du Stream. Les références de colonne dans les définitions d'entité, de série temporelle et de fonction doivent être préfixées par value. ou key. pour indiquer quelle partie du message Kafka lire. Les champs imbriqués sont pris en charge à l'aide de la notation par points (par exemple, value.user.address.city).
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
StreamSource,
Feature,
AggregationFunction,
Sum,
RollingWindow,
)
from datetime import timedelta
client = FeatureEngineeringClient()
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
)
feature = Feature(
name="user_purchase_sum",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(
operator=Sum(input="value.amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
)
Conditions de filtrage sur StreamSource
Utilisez filter_condition pour filtrer les lignes du stream avant l'agrégation, tout comme sur DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Sélection de colonne à partir de Stream
Les fonctionnalités de ColumnSelection fonctionnent avec les sources de streaming. La colonne sélectionnée représente la dernière valeur du Stream pour chaque entité tout en respectant la précision ponctuelle.
from databricks.feature_engineering.entities import ColumnSelection
passenger_count = Feature(
name="passenger_count",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=ColumnSelection(column="value.passenger_count"),
)
Accédez aux champs imbriqués
Vous pouvez accéder aux champs JSON imbriqués à l'aide de la notation par points (par exemple, value.nested_field.amount). Au moment de la diffusion, la charge utile et la réponse de la requête utilisent les noms des nœuds feuilles (par exemple, amount au lieu de value.amount). Les noms des nœuds feuilles doivent être uniques dans toutes les colonnes de sortie d'entité, de séries temporelles et de fonctionnalités au sein d'un modèle ou d'une spécification de fonctionnalité, car l'endpoint de service utilise les noms de feuilles pour router les valeurs.
Fenêtres temporelles pour les fonctionnalités de streaming
Les fonctionnalités de streaming ne prennent en charge que RollingWindow pour les agrégations. Les fenêtres glissantes recalculent en continu sur les données les plus récentes, ce qui s'aligne avec la nature en temps réel des sources de streaming. TumblingWindow et SlidingWindow sont conçus pour le calcul par batch sur des intervalles historiques fixes.
Exemple de notebook des fonctionnalités de streaming
Notebook de démarrage rapide des vues de fonctionnalités en streaming
Entraînement de modèles et inférence
Pour entraîner des modèles et exécuter l'inférence par batch avec des vues de fonctionnalités, y compris log_model(), score_batch() et create_training_set(), veuillez consulter Entraîner des modèles avec des vues de fonctionnalités.
Matérialisation de fonctionnalité
Après avoir défini les fonctionnalités, vous pouvez les matérialiser dans des magasins hors ligne ou en ligne pour une réutilisation efficace dans les workflows d'entraînement et de diffusion. Après la matérialisation des fonctionnalités, vous pouvez servir des modèles à l'aide du service de modèles sur CPU. Pour plus de détails, consultez Vues de fonctionnalités matérialisées.
Bonnes pratiques
Nommage des fonctionnalités
- Utilisez des noms descriptifs pour les fonctionnalités critiques pour l'entreprise.
- Suivez des conventions de nommage cohérentes entre les équipes.
- Utilisez des noms générés automatiquement au fur et à mesure que vous développez des fonctionnalités.
Fenêtres temporelles
- Alignez les limites de fenêtre avec les cycles commerciaux (quotidiens, hebdomadaires).
- Des fenêtres plus courtes capturent les tendances récentes mais peuvent être bruyantes. Des fenêtres plus longues produisent des distributions de caractéristiques plus stables mais pourraient manquer les récents changements de comportement. Choisissez en fonction de la rapidité à laquelle le signal sous-jacent change pour votre cas d'utilisation. Par exemple, une fenêtre de 7 jours atténue les fluctuations quotidiennes et produit des entrées de modèle cohérentes, tandis qu'une fenêtre d'une heure réagit rapidement aux changements de comportement mais pourrait introduire une variance qui dégrade les performances du modèle. Si la précision de votre modèle se dégrade lorsque la distribution change, utilisez une fenêtre plus longue pour stabiliser les entrées.
- Les fenêtres basculantes et glissantes sont plus évolutives que les fenêtres glissantes (continues). start avec sliding windows pour la plupart des cas d'utilisation.
Performance
- Matérialisez les fonctionnalités de la même source de données en un seul appel
materialize_featuresafin de minimiser les analyses de données. - Utilisez la même granularité (par exemple, toutes les durées de diapositives d'une heure ou d'un jour) pour les fonctionnalités sur la même source de données afin de permettre un meilleur regroupement pendant la matérialisation.
Colonnes d'entité vs. conditions de filtre
Utilisez ce guide de décision lorsque vous travaillez avec des fonctionnalités issues de la même table source :
Utilisez entity (sur create_feature) lorsque vous avez besoin de différents niveaux d'agrégation :
- Fonctionnalités au niveau du client (une ligne par client) :
entity=["customer_id"] - Fonctionnalités client-commerçant (plusieurs lignes par client) :
entity=["customer_id", "merchant_id"] - Différents niveaux d'agrégation peuvent partager le même
DeltaTableSource: spécifiez différentes valeursentitypour chaque définition de fonctionnalité.
Utilisez filter_condition (sur DeltaTableSource) lorsque vous devez filtrer des lignes au même niveau d'agrégation :
- Transactions de grande valeur uniquement :
filter_condition="amount > 100"(toujours agrégées par client) - **Commandes terminées uniquement** :
filter_condition="status = 'completed'"(toujours agrégées par client)
Règle générale : Si votre modification entraînait un nombre de lignes différent par valeur d'entité, utilisez différentes valeurs entity sur vos définitions de fonctionnalités. Si vous filtrez uniquement les lignes qui contribuent à la même agrégation, utilisez filter_condition sur la source.
Modèles courants
Analytique client
from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow
fe = FeatureEngineeringClient()
features = [
# Recency: Number of transactions in the last day
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=1)))),
# Frequency: transaction count over the last 90 days
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=90)))),
# Monetary: total spend in the last month
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30)))),
]
Analyse des tendances
# Compare recent vs. historical behavior
fe = FeatureEngineeringClient()
recent_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
historical_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7), delay=timedelta(days=7))),
)
Modèles saisonniers
# Same day of week, 4 weeks ago
fe = FeatureEngineeringClient()
weekly_pattern = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=1), delay=timedelta(weeks=4))),
)
Limitations
- Les noms des colonnes d'entité et de séries chronologiques doivent correspondre entre le dataset d'entraînement (étiqueté) et les définitions de fonctionnalités lorsqu'ils sont utilisés dans l'API
create_training_set. - Le nom de colonne utilisé comme colonne
labeldans le dataset d'entraînement ne doit pas exister dans les tables sources utilisées pour définir lesFeature. - Une liste limitée de fonctions (UDAFs) est prise en charge dans l'API
create_feature. Consultez Fonctions prises en charge. - Les colonnes d’entité ne peuvent pas être de type
DATEouTIMESTAMP. RequestSourceprend en charge uniquement les types de données scalaires définis dansScalarDataType(INTEGER,FLOAT,BOOLEAN,STRING,DOUBLE,LONG,TIMESTAMP,DATE,SHORT). Les types complexes tels que les tableaux, les cartes et les structures ne sont pas pris en charge.RequestSourcene prend pas en charge les fonctions d'agrégation ou les fenêtres temporelles. Seules les fonctionsColumnSelectionpeuvent être utilisées.- L'ensemble des noms de colonnes d'entité, des noms de colonnes de séries chronologiques et des noms de colonnes de caractéristiques de requête doit être globalement unique pour toutes les sources d'un ensemble d'entraînement ou d'un endpoint de service.
Pour les limitations spécifiques à la matérialisation, voir Limitations.