Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite 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
- Workflow de matérialisation et mise en service des caractéristiques
- 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 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 d'ordinateur portable
Notebook de démarrage rapide des vues de fonctionnalités
Obtenir un ordinateur portable
Fonctionnalités de diffusion en continu
Utilisez des fonctionnalités de streaming lorsque les valeurs des fonctionnalités doivent être mises à jour en continu plutôt que sur un calendrier batch. Les fonctionnalités de streaming et de batch utilisent les mêmes Feature constructeurs, fonctions d’agrégation, ainsi que les workflows de formation et de service.
Les fonctionnalités de streaming ne se concrétisent pas dans une boutique hors ligne. Pour l’entraînement et l’inférence en lot, Databricks calcule les valeurs des caractéristiques à partir de la source.
Les fonctionnalités de streaming ont les exigences suivantes :
- Vous devez fournir une valeur
online_config. Les fonctionnalités de streaming ne prennent pas en chargeoffline_config. - Vous ne pouvez pas combiner des fonctionnalités de streaming et de batch en un seul
materialize_featuresappel. Faites un appel séparé pour chaque type de déclencheur. -
transformation_sqln’est pas prise en charge pour les fonctionnalités de streaming. - La matérialisation en streaming traite uniquement les enregistrements qui arrivent après le démarrage du pipeline et ne remplit pas rétroactivement les enregistrements historiques. Les agrégats à fenêtre roulante ne rendent des résultats complets qu’après l’arrivée de la première fenêtre complète de données.
Choisir une source de fonctionnalités de streaming
Choisissez une source en fonction de vos besoins en fraîcheur et de votre installation d’ingestion existante :
- Utilisez une
StreamSourcelorsque la priorité est une fraîcheur inférieure à une seconde.StreamSourceLes fonctionnalités offrent une latence de bout en bout P99 de 200 millisecondes. D’abord mettez en place un flux, puis référencez-le à l’aide deStreamSource. Les sources de flux prennent en charge Kafka comme entrée et maintiennent automatiquement une table Delta d’ingestion en tant que copie historique des données pour l’entraînement. - Utilisez un
DeltaTableSourcelorsque vous avez déjà un chemin d’ingestion à faible latence vers une table Delta. Attendez-vous à une fraîcheur de l’ordre de dizaines de secondes. - Utilisez Zerobus pour alimenter un
DeltaTableSourcelorsque 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. Vous pouvez donc vous attendre à une fraîcheur des fonctionnalités inférieure à une minute.
Définir une fonctionnalité de streaming à l’aide d’un StreamSource
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)),
),
)
Définir une fonctionnalité de streaming à l’aide d’une DeltaTableSource
Pour matérialiser une caractéristique définie sur a DeltaTableSource comme une fonction de streaming, passez StreamingMode comme déclencheur à materialize_features. La définition de la fonctionnalité utilise les mêmes API qu’une fonctionnalité par lot reposant sur un DeltaTableSource. Les sources de table Delta supportent l’agrégation et la sélection de colonnes.
La table Delta source doit avoir activé le flux de données de changement (CDF) en définissant delta.enableChangeDataFeed=true.
L’exemple suivant définit et matérialise une caractéristique d’agrégation avec une source de table Delta.
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(),
)
Utilisez une table Delta remplie par Zerobus
Une table Delta remplie par Zerobus peut servir de source de fonctionnalités de streaming. 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 fonctionnalité de streaming.
Conditions de filtrage sur les sources en flux
Utilisez filter_condition pour filtrer les lignes avant l’agrégation pour un StreamSource ou DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Sélection des colonnes à partir de sources en streaming
ColumnSelection Pour chaque clé d’entité, les variables transmettent la valeur la plus récente provenant de la source, sans agrégation. Lors de l’entraînement, les valeurs des variables respectent l’exactitude temporelle.
Les fonctionnalités de sélection de colonnes ne possèdent pas de TTL. Pour retirer une valeur sélectionnée de la boutique en ligne, la source doit émettre une valeur nulle pour la colonne sélectionnée.
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 depuis un StreamSource
Pour un StreamSource, vous pouvez accéder à des champs JSON imbriqués en utilisant la notation à points (par exemple, value.nested_field.amount). Au moment du service, les données utiles 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 des nœuds feuilles doivent être uniques parmi tous les noms d’entités, de séries chronologiques et de fonctionnalités au sein d’un modèle ou d’une spécification de fonctionnalités, car le point de terminaison de service utilise les noms des feuilles pour acheminer les valeurs.
Fenêtres de temps pour les fonctionnalités de diffusion en continu
Les fonctionnalités de diffusion prennent uniquement en charge RollingWindow pour 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 mettre en service des modèles avec la mise en service de modèles CPU. 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’apprentissage ou un point de terminaison de mise en 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.