Remarque
L’accès à cette page requiert une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page requiert une autorisation. Vous pouvez essayer de modifier des répertoires.
Important
Cette fonctionnalité est disponible en préversion publique. Les administrateurs d’espace de travail peuvent contrôler l’accès à cette fonctionnalité à partir de la page Aperçus . Consultez Gérer les préversions d’Azure Databricks.
Les vues de fonctionnalités vous permettent de définir et de calculer des fonctionnalités à partir de sources de données. Les fonctionnalités peuvent être définies à l’aide de diverses sources (table Delta, Flux Kafka et données au moment de la requête) et de calculs (agrégations à fenêtres temporelles, sélections de colonnes simples, etc.). Ce guide couvre les flux de travail suivants :
-
Workflow de développement de fonctionnalités
- Utilisez
create_featurepour définir des objets de fonctionnalités du catalogue Unity pouvant être utilisés dans l'entraînement de modèles et les flux de travaux. - Vous pouvez également construire
Featuredes objets localement et les utiliserregister_featurepour les rendre persistants dans le catalogue Unity ultérieurement. Les fonctionnalités construites localement peuvent être utilisées aveccreate_training_setavant l’enregistrement.
- Utilisez
-
Flux de travail d’entraînement de modèle
- Utilisez
create_training_setpour calculer des caractéristiques agrégées à un point donné dans le temps pour l'apprentissage automatique. Pour obtenir une documentation détaillée sur l’apprentissage avec des vues de fonctionnalités, consultez Entraîner des modèles avec des vues de fonctionnalités.
- Utilisez
-
Matérialisation des caractéristiques et flux de travail de service
- 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 offline pour une réutilisation efficace, ou dans un magasin en ligne pour la distribution en ligne. - Utilisez
create_training_setavec la vue matérialisée pour préparer un jeu de données d'entraînement par lots hors ligne.
- Après avoir défini une fonctionnalité avec
Pour plus de détails sur l’API, consultez la documentation de référence de l’API Feature Views.
Exigences
Un calcul serverless ou un cluster de calcul classique exécutant Databricks Runtime 17.0 ML ou version ultérieure.
Vous devez installer le package Python personnalisé. Exécutez les lignes de code suivantes chaque fois que vous exécutez un notebook :
%pip install databricks-feature-engineering>=0.16.0 dbutils.library.restartPython()
Exemple de démarrage rapide
Pour obtenir un notebook de démarrage rapide exécutable, consultez Example 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 d'ordinateur portable
Bloc-notes de démarrage rapide Vues de fonctionnalités
Obtenir un ordinateur portable
Fonctionnalités de diffusion en continu
En plus des fonctionnalités de traitement par lots à partir de tables Delta, vous pouvez définir des fonctionnalités à partir de sources de diffusion en continu pour les cas d’usage en temps réel. Les fonctionnalités de streaming utilisent la même classe Feature que les fonctionnalités batch — les mêmes constructeurs Feature, les mêmes fonctions d’agrégation, les mêmes workflows d’entraînement et d’inférence — de sorte que le passage du batch au temps réel nécessite très peu de modifications du code. Une fois matérialisées, les fonctionnalités de diffusion en continu fournissent une fraîcheur de bout en bout inférieure à la seconde (latence p99 de 200 ms) directement à vos points de terminaison de service de modèle.
Pour utiliser des fonctionnalités de streaming, commencez par configurer un flux, puis référencez-le à l’aide d’un StreamSource. Les sources de flux prennent en charge Kafka comme entrée et conservent automatiquement une table d’ingestion (Delta) en tant que copie historique des données pour l’entraînement.
Définir une fonctionnalité de diffusion en continu
Une StreamSource référence à un flux par son nom en trois parties (catalog.schema.stream_name). Un flux n’est pas un objet sécurisable du catalogue Unity, mais il est limité à un schéma de catalogue Unity et l’accès est régi par la table d’ingestion de Stream. Les références de colonne dans les définitions d’entité, de séries temporelles et de fonction doivent être préfixées par value. ou key. pour indiquer de 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 filtre sur StreamSource
Permet filter_condition de filtrer les lignes du flux avant l’agrégation, comme sur DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Sélection de colonnes à partir de flux
ColumnSelection les fonctionnalités fonctionnent avec les sources de diffusion en continu. La colonne sélectionnée représente la valeur la plus récente du flux pour chaque entité tout en respectant la précision à un point dans le temps.
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
Vous pouvez accéder aux champs JSON imbriqués à l’aide de la notation par points (par exemple, value.nested_field.amount). Lors de l’exécution, le corps de la requête et la réponse utilisent les noms des nœuds feuilles (par exemple, amount au lieu de value.amount). Les noms de nœud feuille doivent être uniques dans toutes les colonnes d’entité, de série chronologique et de sortie de caractéristique au sein d’un modèle ou d’une spécification de fonctionnalité, car le point de terminaison de service utilise des noms feuilles pour router les valeurs.
Fenêtres de temps pour les fonctionnalités de diffusion en continu
Les fonctionnalités de diffusion en continu prennent uniquement RollingWindow en charge les agrégations. Les fenêtres glissantes recalculent en continu à partir des données les plus récentes, ce qui correspond à la nature en temps réel des sources de données en continu.
TumblingWindow et SlidingWindow sont conçus pour le calcul par lots sur des intervalles historiques fixes.
Exemple de carnet sur les fonctionnalités de diffusion en continu
Notebook de démarrage rapide des vues de fonctionnalités en streaming
Obtenir un ordinateur portable
Formation et inférence du modèle
Pour entraîner des modèles et exécuter l’inférence par lots avec des vues de fonctionnalités, notamment log_model(), score_batch()et create_training_set(), consultez Entraîner des modèles avec des vues de fonctionnalité.
Matérialisation des fonctionnalités
Après avoir défini des fonctionnalités, vous pouvez les matérialiser dans des magasins hors connexion ou en ligne pour une réutilisation efficace dans l’apprentissage et le service des flux de travail. Après avoir matérialisé des fonctionnalités, vous pouvez servir des modèles avec un serveur de modèles utilisant le processeur. Pour plus d’informations, consultez Materialize Feature Views.
Bonnes pratiques
Nommage des fonctionnalités
- Utilisez des noms descriptifs pour les fonctionnalités critiques pour l’entreprise.
- Suivez les conventions d’affectation de noms cohérentes entre les équipes.
- Utilisez des noms générés automatiquement lorsque vous commencez à développer des fonctionnalités.
Fenêtres Délai
- Aligner les limites des fenêtres avec les cycles d’activité (quotidiens, hebdomadaires).
- Les fenêtres plus courtes capturent les tendances récentes, mais peuvent être bruyantes. Les fenêtres plus longues produisent des distributions de fonctionnalités plus stables, mais peuvent manquer des changements de comportement récents. Choisissez en fonction de la rapidité avec laquelle le signal sous-jacent change pour votre cas d’usage. Par exemple, une fenêtre de 7 jours lisse les fluctuations quotidiennes et produit des entrées de modèle cohérentes, tandis qu’une fenêtre de 1 heure réagit rapidement aux changements comportementaux, mais peut 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 périodiques et glissantes sont plus évolutives que les fenêtres roulantes. Commencez par les fenêtres glissantes pour la plupart des cas d’usage.
Performance
- Matérialisez les caractéristiques de la même source de données dans un appel unique
materialize_featurespour réduire les analyses de données. - Utilisez la même granularité (par exemple, toutes les durées de diapositives de 1 heure ou de 1 jour) pour les fonctionnalités de la même source de données afin d’améliorer le regroupement pendant la matérialisation.
Colonnes d’entité et conditions de filtre
Utilisez ce guide de décision lors de l’utilisation des fonctionnalités de la même table source :
Utilisez entity (on 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-marchand (plusieurs lignes par client) :
entity=["customer_id", "merchant_id"] -
Différents niveaux d’agrégation peuvent partager le même
DeltaTableSource: spécifier des valeurs différentesentitysur chaque définition de fonctionnalité
Utilisez filter_condition (on DeltaTableSource) quand vous devez filtrer des lignes au même niveau d’agrégation :
-
Transactions à valeur élevée uniquement :
filter_condition="amount > 100"(toujours agrégée par client) -
Commandes terminées uniquement :
filter_condition="status = 'completed'"(toujours agrégé par client)
Règle générale: Si votre modification entraînerait un nombre différent de lignes par valeur d’entité, utilisez différentes entity valeurs sur vos définitions de fonctionnalités. Si vous filtrez simplement les lignes qui contribuent à la même agrégation, utilisez-la filter_condition sur la source.
Modèles courants
Analyse des clients
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 de 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érie chronologique doivent correspondre entre le jeu de données d’entraînement (étiqueté) et les définitions de fonctionnalités lorsqu’elles sont utilisées dans l’API
create_training_set. - Le nom de colonne utilisé comme colonne
labeldans le jeu de données d’entraînement ne devrait 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 les fonctions prises en charge. - Les colonnes d’entité ne peuvent pas être de type
DATEouTIMESTAMP. -
RequestSourceprend uniquement en charge les types de données scalaires définis dansScalarDataType(INTEGER,FLOAT,BOOLEAN,STRINGDOUBLE,LONG,TIMESTAMP, ,DATE, ).SHORTLes types complexes tels que les tableaux, les cartes et les structs ne sont pas pris en charge. -
RequestSourcene prend pas en charge les fonctions d’agrégation ou les fenêtres de temps. Seules les fonctionsColumnSelectionpeuvent être utilisées. - L’ensemble des noms de colonnes d’entité, des noms de colonnes de séries temporelles et des noms de colonnes de caractéristiques de requête doit être unique à l’échelle globale sur toutes les sources d’un jeu d’entraînement ou un point de terminaison de service.
-
score_batchrisque de ne pas fonctionner dans un environnement de calcul serverless. Contourner ce problème à l’aide d’un cluster de calcul classique exécutant Databricks Runtime 17.0 ML ou version ultérieure.
Pour connaître les limitations spécifiques à la matérialisation, consultez Limitations.