Aller au contenu principal

Vues de fonctionnalités

info

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_feature pour 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 Feature et utiliser register_feature pour les conserver ultérieurement dans Unity Catalog. Les fonctionnalités construites localement peuvent être utilisées avec create_training_set avant l'enregistrement.
  • Workflow de Model training

    • Utilisez create_training_set pour 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.
  • Workflow de matérialisation et de service des fonctionnalités

    • Après avoir défini une fonctionnalité avec create_feature ou l'avoir récupérée à l'aide de get_feature, vous pouvez utiliser materialize_features pour 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_set avec la vue matérialisée pour préparer un dataset d'entraînement batch hors ligne.

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.

Python
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).

Python
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.

Python
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.

Python
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_features afin 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 valeurs entity pour 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

Python
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

Python
# 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

Python
# 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 label dans le dataset d'entraînement ne doit pas exister dans les tables sources utilisées pour définir les Feature.
  • 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 DATE ou TIMESTAMP.
  • RequestSource prend en charge uniquement les types de données scalaires définis dans ScalarDataType (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.
  • RequestSource ne prend pas en charge les fonctions d'agrégation ou les fenêtres temporelles. Seules les fonctions ColumnSelection peuvent ê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.