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 support CronSchedule or TableTrigger, 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

Utilisez des fonctionnalités de streaming lorsque les valeurs des fonctionnalités doivent être mises à jour en continu plutôt que selon un calendrier de batch. Les fonctionnalités de streaming et de batch utilisent les mêmes constructeurs Feature, fonctions d'agrégation, ainsi que les mêmes workflows d'entraînement et de service.

Les fonctionnalités de streaming ne sont pas matérialisées dans un magasin hors ligne. Pour l'entraînement et l'inférence par batch, Databricks calcule les valeurs des fonctionnalités à partir de la source.

Les fonctionnalités de streaming ont les exigences suivantes :

  • Vous devez fournir un online_config. Les fonctionnalités de streaming ne prennent pas en charge offline_config.
  • Vous ne pouvez pas combiner des fonctionnalités de streaming et de batch dans un seul appel materialize_features. Effectuez un appel distinct pour chaque type de Trigger.
  • transformation_sql n’est pas pris en charge pour les fonctionnalités de streaming.
  • La matérialisation en streaming traite uniquement les enregistrements qui arrivent après le start du pipeline et ne remplit pas les enregistrements historiques. Les agrégats sur fenêtre glissante ne renvoient des résultats complets qu'après l'arrivée de la première fenêtre de données complète.

Choix d’une source de fonctionnalité de streaming

Choisissez une source en fonction de vos exigences de fraîcheur et de votre configuration d'ingestion existante :

  • Utilisez un StreamSource lorsque la fraîcheur inférieure à la seconde est la priorité. Les fonctionnalités StreamSource offrent une latence de bout en bout p99 de 200 millisecondes. Commencez par configurer un Stream, puis référencez-le à l'aide d'un StreamSource. Les sources de Stream prennent en charge Kafka en entrée et maintiennent automatiquement une table Delta d'ingestion en tant que copie historique des données pour l'entraînement.
  • Utilisez un DeltaTableSource lorsque vous disposez déjà d’un chemin d’ingestion à faible latence vers une table Delta. Attendez-vous à une actualité de l’ordre de quelques dizaines de secondes.
  • Utilisez Zerobus pour remplir un DeltaTableSource lorsque vous ne disposez pas déjà d'un chemin d'ingestion à faible latence. L'ingestion Zerobus prend de l'ordre de quelques dizaines de secondes ; attendez-vous donc à une fraîcheur des caractéristiques inférieure à la minute.

Définir une fonctionnalité de streaming à l’aide d’une StreamSource

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)),
),
)

Définir une fonctionnalité de streaming à l'aide d'une DeltaTableSource

Pour matérialiser une fonctionnalité définie sur un DeltaTableSource en tant que fonctionnalité de streaming, transmettez StreamingMode comme trigger à materialize_features. La définition de la fonctionnalité utilise les mêmes APIs qu’une fonctionnalité batch prise en charge par un DeltaTableSource. Les sources de table Delta prennent en charge les fonctionnalités d’agrégation et de sélection de colonnes.

La table Delta source doit avoir le flux de données de modification (CDF) activé en définissant delta.enableChangeDataFeed=true.

L’exemple suivant définit et matérialise une fonctionnalité d’agrégation avec une source de table Delta.

Python
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
DeltaTableSource,
AggregationFunction,
OnlineStoreConfig,
Sum,
RollingWindow,
StreamingMode,
)
from datetime import timedelta

client = FeatureEngineeringClient()

source = DeltaTableSource(
catalog_name="my_catalog",
schema_name="my_schema",
table_name="transactions",
)

feature = client.create_feature(
catalog_name="my_catalog",
schema_name="my_schema",
name="user_purchase_sum",
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(
operator=Sum(input="amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
)

online_config = OnlineStoreConfig(
catalog_name="my_catalog",
schema_name="my_schema",
table_name_prefix="streaming_features",
online_store_name="my_online_store",
)

client.materialize_features(
features=[feature],
online_config=online_config,
trigger=StreamingMode(),
)

Utiliser une table Delta alimentée par Zerobus

Une table Delta alimentée par Zerobus peut servir de source de streaming pour les features. Zerobus ne définit pas delta.enableChangeDataFeed=true automatiquement. Vous devez définir cette propriété manuellement sur la table Delta cible avant de l'utiliser comme source de streaming pour les features.

Conditions de filtrage sur les sources de streaming

Utilisez filter_condition pour filtrer les lignes avant l'agrégation pour StreamSource ou DeltaTableSource.

Python
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)

Sélection de colonnes à partir de sources de streaming

Les fonctionnalités ColumnSelection fonctionnent avec des sources de streaming. La colonne sélectionnée représente la dernière valeur issue de la source pour chaque entité, tout en respectant la précision temporelle (point-in-time).

Les caractéristiques de sélection de colonne n'ont pas de TTL. Pour supprimer une valeur sélectionnée du magasin en ligne, la source doit émettre une valeur nulle pour la colonne sélectionnée.

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éder aux champs imbriqués à partir d’une StreamSource

Pour un StreamSource, vous pouvez accéder aux champs JSON imbriqués en utilisant la notation par points (par exemple, value.nested_field.amount). Au moment du service, la charge utile de la requête et la réponse utilisent les noms des nœuds terminaux (par exemple, amount au lieu de value.amount). Les noms des nœuds terminaux doivent être uniques parmi tous les noms d'entités, de séries temporelles et de fonctionnalités au sein d'un modèle ou d'une spécification de fonctionnalité (Feature Spec), car l'endpoint de service utilise les noms terminaux pour acheminer 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.