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.
Contrôle d’accès
Les fonctionnalités sont des objets catalogue Unity pouvant être régis. L’accès à une fonctionnalité est contrôlé par les CREATE FEATUREprivilèges de catalogue , READ FEATUREet MANAGE Unity. Pour obtenir des descriptions complètes, consultez informations de référence sur les privilèges du catalogue Unity.
-
CREATE FEATURE— Requis pour créer une fonctionnalité dans un schéma.create_featureetregister_featureexigerCREATE FEATUREsur le schéma parent. En suivant le principe du privilège minimum, accordez-leCREATE FEATUREau niveau du schéma ; vous pouvez également l’accorder sur un catalogue pour permettre la création de fonctionnalités dans n’importe quel schéma de ce catalogue. -
READ FEATURE— Requis pour lire une fonctionnalité et ses données.get_feature,create_training_setet lecture des données de caractéristique matérialisées pour l’entraînement ou le service requisREAD FEATUREsur la fonctionnalité.READ FEATUREaccordé sur un schéma ou un catalogue s’applique à toutes les fonctionnalités actuelles et futures qu’il contient. -
MANAGE— Requis pour gérer le cycle de vie et les subventions d’une fonctionnalité. La suppression d’une fonctionnalité avecdelete_feature, et la matérialisation d’une fonctionnalité avecmaterialize_featuresoudelete_materialized_feature, nécessiteMANAGEla fonctionnalité.
Toutes les opérations de fonctionnalité nécessitent USE CATALOG également sur le catalogue parent et USE SCHEMA sur le schéma parent. Pour savoir comment MANAGE et READ FEATURE s’appliquer à la matérialisation, consultez Autorisations.
API d’affichage des fonctionnalités
Feature constructeur et register_feature()
L’approche recommandée consiste à construire un Feature objet localement et à l’utiliser register_feature pour le conserver dans le catalogue Unity. Ce flux de travail en deux étapes vous permet d’expérimenter des fonctionnalités (y compris create_training_set) avant de les inscrire.
Feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
)
FeatureEngineeringClient.register_feature() inscrit une construction Feature locale dans le catalogue Unity.
FeatureEngineeringClient.register_feature(
feature: Feature, # Required: A Feature instance (not already registered)
catalog_name: str, # Required: UC catalog name
schema_name: str, # Required: UC schema name
) -> Feature
from databricks.feature_engineering.entities import Feature, DeltaTableSource, AggregationFunction, Sum, RollingWindow
from datetime import timedelta
# Step 1: Construct the feature locally
feature = Feature(
source=DeltaTableSource(catalog_name="main", schema_name="store", table_name="transactions"),
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
# Step 2: Register in Unity Catalog
fe = FeatureEngineeringClient()
registered_feature = fe.register_feature(
feature=feature,
catalog_name="main",
schema_name="store",
)
create_feature()
FeatureEngineeringClient.create_feature() valide, construit et inscrit immédiatement une fonctionnalité dans le catalogue Unity en une seule étape. Utilisez cette option lorsque vous n’avez pas besoin d’expérimenter la fonctionnalité localement.
FeatureEngineeringClient.create_feature(
source: DataSource, # Required: DeltaTableSource, StreamSource, or RequestSource
function: Union[AggregationFunction, ColumnSelection], # Required: Aggregation or column selection
catalog_name: str, # Required: The catalog name for the feature
schema_name: str, # Required: The schema name for the feature
entity: Optional[List[str]] = None, # Required for all sources except RequestSource: entity columns
timeseries_column: Optional[str] = None, # Required for all sources except RequestSource: timestamp column
name: Optional[str] = None, # Optional: Feature name (auto-generated if omitted)
description: Optional[str] = None, # Optional: Feature description
) -> Feature
Paramètres :
-
source: source de données utilisée dans le calcul des fonctionnalités (DeltaTableSourceouStreamSourceRequestSource). -
function: quiAggregationFunctionregroupe l’opérateur (par exemple,Sum(input="amount")), la colonne d’entrée et la fenêtre de temps. OuColumnSelection("column_name")pour les fonctionnalités directes. -
catalog_name: nom du catalogue Du catalogue Unity pour la fonctionnalité. -
schema_name: nom du schéma du catalogue Unity pour la fonctionnalité. -
entity: liste des noms de colonnes qui définissent les clés d’agrégation ou de recherche (clés primaires). Obligatoire pour tous les types sources, à l’exceptionRequestSourcede . Par exemple,["user_id"]les agrégats ou les recherches par utilisateur. -
timeseries_column: colonne timestamp utilisée pour l’agrégation de fenêtre de temps ou la sélection de la dernière valeur. Obligatoire pour tous les types sources, à l’exceptionRequestSourcede . -
name: nom de fonctionnalité facultatif. S’il est omis, généré automatiquement à partir de la colonne d’entrée, de la fonction et de la fenêtre (par exemple,amount_avg_rolling_7d). -
description: description facultative de la fonctionnalité.
Retourne : Une instance de fonction validée
Soulève: ValueError si une validation échoue
delete_feature()
Supprime une fonctionnalité du catalogue Unity par son nom complet.
FeatureEngineeringClient.delete_feature(
full_name: str, # Required: '<catalog>.<schema>.<feature_name>'
) -> None
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")
Avant de supprimer une fonctionnalité, supprimez ou mettez à jour des modèles ou des spécifications de fonctionnalités qui le référencent. Si la fonctionnalité a été matérialisée, supprimez d’abord la fonctionnalité matérialisée. Découvrez comment supprimer une fonctionnalité matérialisée.
Noms générés automatiquement
Lorsqu’il name est omis, un nom est généré automatiquement. Les noms générés suivent le modèle : {column}_{function}_{window}. Par exemple:
-
price_avg_rolling_1h(prix moyen de 1 heure) -
transaction_count_rolling_30d_1d(Nombre de transactions de 30 jours avec un délai de 1d à partir de l’horodatage de l’événement)
Fonctions prises en charge
Fonctions d’agrégation
Note
Les fonctions d’agrégation sont encapsulées dans une AggregationFunction fenêtre de temps, comme décrit dans les fenêtres de temps. Chaque fonction prend un input paramètre spécifiant la colonne source à agréger.
| Function | Description | Exemple de cas d’usage |
|---|---|---|
Sum(input="column") |
Total des valeurs | Utilisation quotidienne de l’application par utilisateur en minutes |
Avg(input="column") |
Moyenne des valeurs | Montant moyen de la transaction |
Count(input="column") |
Nombre d’enregistrements | Nombre de connexions par utilisateur |
Min(input="column") |
Valeur minimale | Fréquence cardiaque la plus faible enregistrée par un appareil portable |
Max(input="column") |
Valeur maximale | Montant de transaction le plus élevé par session |
StddevPop(input="column") |
Écart type de population | Variabilité quotidienne des transactions entre tous les clients |
StddevSamp(input="column") |
Exemple d’écart type | Variabilité des taux de clics de campagne publicitaire |
VarPop(input="column") |
Variance de la population | Propagation des lectures de capteurs pour les appareils IoT dans une usine |
VarSamp(input="column") |
Variance d'échantillon | Répartition des évaluations de films sur un groupe échantillonné |
ApproxCountDistinct(input="column", relativeSD=0.05) |
Nombre unique approximatif | Nombre distinct d’articles achetés |
ApproxPercentile(input="column", percentile=0.95, accuracy=100) |
Percentile approximatif | Latence de réponse p95 |
First(input="column") |
Première valeur | Horodatage de première connexion |
Last(input="column") |
Dernière valeur | Montant d’achat le plus récent |
Note
First et Last inclure des valeurs null par défaut. Pour ignorer les valeurs Null, ajoutez un filter_condition élément qui exclut explicitement les colonnes d’entrée qui sont null.
ColumnSelection (pass-through)
ColumnSelection sélectionne une seule colonne à partir d’une source sans appliquer d’agrégation. Il est encapsulé directement dans le paramètre (pas à l’intérieur functionAggregationFunction). Le type de retour est déduit du schéma source.
| Function | Description | Exemple de cas d’usage |
|---|---|---|
ColumnSelection("col") |
Dernière valeur d’une colonne (aucune agrégation) | Catégorie de fournisseur la plus récente, pass-through d’un champ de demande |
ColumnSelection peut être utilisé avec n’importe quelle source de données :
-
DeltaTableSource: retourne la valeur la plus récente par clé d’entité via une jointure à un point dans le temps (aucune agrégation de fenêtre de recherche). -
StreamSource: retourne la valeur la plus récente par clé d’entité à partir du flux (aucune agrégation de fenêtre de recherche). -
RequestSource: passe par la valeur fournie au moment de l’inférence (ou extraite du DataFrame étiqueté au moment de l’entraînement).
from databricks.feature_engineering.entities import (
ColumnSelection, DeltaTableSource, Feature, FieldDefinition,
RequestSource, ScalarDataType,
)
delta_source = DeltaTableSource(
catalog_name="main", schema_name="feature_store", table_name="transactions",
)
request_source = RequestSource(
schema=[
FieldDefinition(name="session_duration", data_type=ScalarDataType.DOUBLE),
]
)
# ColumnSelection from a Delta table
latest_amount = Feature(
source=delta_source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
name="latest_transaction_amount",
)
# ColumnSelection from a RequestSource
session_feature = Feature(
source=request_source,
function=ColumnSelection("session_duration"),
name="session_duration",
)
Exemple : fonctionnalités d’agrégation et de sélection de colonnes
L’exemple suivant montre les fonctionnalités définies sur la même source de données.
from databricks.feature_engineering.entities import (
AggregationFunction, Feature, Sum, Avg, ApproxCountDistinct,
ColumnSelection, RollingWindow,
)
from datetime import timedelta
window = RollingWindow(window_duration=timedelta(days=7))
sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Sum(input="amount"), window),
)
avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(Avg(input="amount"), window),
)
distinct_count = Feature(
source=source,
entity=["user_id"],
timeseries_column="event_time",
function=AggregationFunction(ApproxCountDistinct(input="product_id", relativeSD=0.01), window),
)
# Column selection (no aggregation, no time window)
latest_amount = Feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="event_time",
name="latest_amount",
)
Fonctionnalités avec des conditions de filtre
Le filter_condition paramètre vous permet de filtrer les lignes de la table source avant l’informatique des agrégations. Cela fonctionne comme une clause SQL WHERE appliquée avant le regroupement et l’agrégation des données.
Note
filter_condition filtre les lignes avant l’agrégation, comme une clause SQL WHERE appliquée avant GROUP BY. Elle ne modifie pas la granularité, qui est toujours définie par entity la définition de fonctionnalité.
Les filtres sont utiles lors de l’utilisation de tables sources volumineuses qui incluent un super-ensemble de données nécessaire pour le calcul des fonctionnalités et réduisent le besoin de créer des vues distinctes sur ces tables.
from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow, DeltaTableSource
from datetime import timedelta
# Source with filter applied at the source level
high_value_transactions = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="transactions",
filter_condition="amount > 100", # Only transactions over $100
)
high_value_sales = Feature(
source=high_value_transactions,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30))),
)
# Multiple conditions
completed_orders_source = DeltaTableSource(
catalog_name="main",
schema_name="ecommerce",
table_name="orders",
filter_condition="status = 'completed' AND payment_method = 'credit_card'",
)
completed_orders = Feature(
source=completed_orders_source,
entity=["user_id"],
timeseries_column="order_time",
function=AggregationFunction(Count(input="order_id"), RollingWindow(window_duration=timedelta(days=7))),
)
# Filter on a StreamSource
from databricks.feature_engineering.entities import StreamSource
purchase_stream = StreamSource(
full_name="main.ecommerce.transactions_stream",
filter_condition="value.event_type = 'purchase'",
)
purchase_total = Feature(
source=purchase_stream,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(Sum(input="value.amount"), RollingWindow(window_duration=timedelta(hours=1))),
)
Sources de données
DeltaTableSource
DeltaTableSource est un objet Python éphémère utilisé pour définir la façon dont les fonctionnalités sont calculées à partir d’une table source. Elle ne crée pas de table. Il spécifie la configuration pour lire les données et agréger les fonctionnalités.
DeltaTableSource(
catalog_name: str, # Required: Catalog name
schema_name: str, # Required: Schema name
table_name: str, # Required: Table name
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause to filter source data
transformation_sql: Optional[str] = None, # Optional: SQL SELECT expression for column transformations
dataframe_schema: Optional[str] = None, # Required if transformation_sql is set: schema of the resulting DataFrame
)
Paramètres :
-
catalog_name, ,schema_name:table_nameidentifier la table Delta source dans le catalogue Unity. -
filter_condition: clause SQLWHEREappliquée avant l’agrégation. Exemple :"status = 'completed'". -
transformation_sql: expression SQLSELECTappliquée à la table source. Utilisez-la pour renommer des colonnes, des types de cast ou des colonnes dérivées de calcul avant l’agrégation. En cas d’omission, toutes les colonnes sont sélectionnées (*). Exemple :"user_id, CAST(amount AS DOUBLE) AS amount, event_time". -
dataframe_schema: schéma du DataFrame résultant après les transformations, au format JSON Spark StructType (à partir dedf.schema.json()). Obligatoire sitransformation_sqlest fourni. Cela indique au système les noms et les types de colonnes résultant de votre transformation.
Lorsque les deux filter_condition et transformation_sql sont définis, la requête résultante est la suivante : SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.
Note
( timeseries_column spécifié sur la définition de fonctionnalité, et non activé DeltaTableSource) doit être de type TimestampType ou DateType. Les types entiers peuvent fonctionner, mais entraîner une perte de précision pour les agrégats de fenêtre de temps.
Exemple : Utilisation transformation_sql pour les transformations de colonne
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="raw_events",
transformation_sql="user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time",
filter_condition="event_type = 'purchase'",
dataframe_schema=spark.sql(
"SELECT user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time FROM main.analytics.raw_events LIMIT 0"
).schema.json(),
)
Exemple : dérivation transformation_sql et dataframe_schema utilisation d’un DataFrame PySpark
Vous pouvez écrire votre transformation en tant que requête PySpark, puis extraire le schéma du DataFrame résultant :
df = spark.sql(f"""
SELECT user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time
FROM main.analytics.events
WHERE event_date >= date_sub(current_date(), 7)
LIMIT 0
""")
# Use df.schema.json() as the dataframe_schema
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
transformation_sql="user_id, CAST(amount AS DOUBLE) / 100 AS amount_dollars, event_time",
filter_condition="event_date >= date_sub(current_date(), 7)",
dataframe_schema=df.schema.json(),
)
Note
transformation_sql prend uniquement en charge les expressions en ligne (renommages de colonne, casts, arithmétiques). Les fonctions d’agrégation comme COUNT(*) ou SUM() non prises en charge. Utilisez AggregationFunction plutôt la définition de fonctionnalité.
DeltaTableSource.from_sql()
En guise de commodité, vous pouvez créer une DeltaTableSource requête SQL. La méthode analyse la requête pour extraire automatiquement le nom de la table, transformation_sqlet filter_condition.
DeltaTableSource.from_sql(
sql: str, # Required: SQL SELECT query
spark: SparkSession, # Required: active SparkSession (for schema inference)
) -> DeltaTableSource
Seules les requêtes simples SELECT ... FROM ... [WHERE ...] sont prises en charge. Le code SQL complexe (JOIN, sous-requêtes, CTEs, UNION) est rejeté. Pour les requêtes complexes, créez DeltaTableSource directement avec transformation_sql et filter_condition.
from databricks.feature_engineering.entities import (
AggregationFunction,
DeltaTableSource,
Feature,
Sum,
TumblingWindow,
)
source = DeltaTableSource.from_sql(
spark=spark,
sql=f"SELECT customer_id, event_ts, amount * 2 AS doubled_amount, amount FROM {CATALOG}.{SCHEMA}.{TABLE}",
)
feature = Feature(
source=source,
function=AggregationFunction(Sum(input="doubled_amount"), time_window=TumblingWindow(window_duration=timedelta(days=7))),
entity=["customer_id"], timeseries_column="event_ts",
)
Itérer avec to_dataframe()
Permet source.to_dataframe() d’afficher un aperçu des données qui seront utilisées pour le calcul des fonctionnalités. Cela est utile pour itérer filter_condition et transformation_sql jusqu’à ce qu’ils produisent les résultats attendus.
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
filter_condition="event_type = 'purchase'",
)
# Preview the filtered source data
source.to_dataframe().display()
Présentation des entités
Les colonnes d’entité définissent le niveau d’agrégation de vos fonctionnalités. Elles sont spécifiées sur la Feature définition, et non sur DeltaTableSource. Les entités déterminent :
-
Comment les données sont regroupées : les fonctionnalités sont agrégées par combinaison unique de valeurs d’entité (comme
GROUP BYdans SQL) - Structure de clé primaire : chaque combinaison d’entités unique entraîne une ligne de fonctionnalités calculées
Exemple : fonctionnalités au niveau du client
Le code suivant agrège les fonctionnalités au niveau du client (une ligne par client) :
from databricks.feature_engineering.entities import DeltaTableSource
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="user_events",
)
Feature(
source=source,
entity=["user_id"], # Features aggregated per user
timeseries_column="event_time", # Timestamp for time windows
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
Exemple : fonctionnalités au niveau du Customer-Store
Pour agréger des fonctionnalités à un niveau plus détaillé (une ligne par combinaison de magasin de clients), utilisez plusieurs colonnes d’entité :
source = DeltaTableSource(
catalog_name="main",
schema_name="retail",
table_name="transactions",
)
Feature(
source=source,
entity=["user_id", "store_id"], # Features aggregated per user-store pair
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
Lorsque vous avez besoin de fonctionnalités à différents niveaux d’agrégation (par exemple, au niveau du client et au niveau du magasin client), utilisez différentes entity valeurs dans vos définitions de fonctionnalités. La même DeltaTableSource chose peut être partagée entre les fonctionnalités avec différentes configurations d’entité.
StreamSource
StreamSource fait référence à un flux. Le flux contient la configuration de connexion, d’authentification, de schéma et d’ingestion pour la source de diffusion en continu. Pour Kafka, les références de colonne dans les définitions de fonctionnalités doivent être préfixées value. ou key. indiquer la partie du message à lire.
StreamSource(
full_name: str, # Required: Three-part Stream name (catalog.schema.stream)
filter_condition: Optional[str], # Optional: SQL WHERE clause applied before aggregation
)
Paramètres :
-
full_name: nom complet en trois parties d’un flux (par exemple,"my_catalog.my_schema.my_stream"). -
filter_condition(facultatif) : clause SQLWHEREappliquée au flux de données avant l’agrégation, à l’aide de références de colonnes avec préfixe point (par exemple)."value.event_type = 'purchase'"
from databricks.feature_engineering.entities import StreamSource
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
RequestSource
RequestSource définit un schéma pour les données fournies au moment de l’inférence dans la charge utile de la requête plutôt qu’à partir d’une table pré-matérialisée. Pendant l’entraînement, ces colonnes sont extraites du DataFrame étiqueté passé à create_training_set. Pendant le service de modèle, l’appelant doit les inclure dans la charge utile de requête HTTP.
RequestSource est utilisé avec ColumnSelection (pour passer directement par une valeur). Il ne prend pas en charge les fonctions d’agrégation ou les fenêtres de temps.
Définition du schéma
Définissez le schéma comme une liste d’objets, chacun spécifiant un nom de FieldDefinition colonne et un ScalarDataType:
from databricks.feature_engineering.entities import (
FieldDefinition, RequestSource, ScalarDataType,
)
request_source = RequestSource(
schema=[
FieldDefinition(name="transaction_amount", data_type=ScalarDataType.DOUBLE),
FieldDefinition(name="vendor_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_id", data_type=ScalarDataType.STRING),
FieldDefinition(name="transaction_time", data_type=ScalarDataType.DATE),
]
)
Types de données prises en charge
RequestSourceprend en charge les types scalaires définis dans ScalarDataType: INTEGER, , FLOATBOOLEANSTRING, DOUBLE, LONG, TIMESTAMP, , DATE. SHORT Les types complexes tels que les tableaux, les cartes et les structs ne sont pas pris en charge.
Comment les données de requête sont hydratées
| Context | Behavior |
|---|---|
Formation (create_training_set) |
Les colonnes sont extraites du DataFrame étiqueté. Les types sont validés par rapport au schéma déclaré. Les incompatibilités déclenchent une erreur (aucune conversion implicite). |
| Service (point de terminaison de modèle) | Les colonnes sont extraites ou dataframe_recordsdataframe_split dans la requête HTTP. Les valeurs JSON sont converties en types déclarés (par exemple, le nombre JSON → DOUBLE). |
Signature de modèle
Lorsqu’un modèle est enregistré à l’aide log_model d’un jeu d’entraînement qui inclut des RequestSource fonctionnalités, les RequestSource colonnes sont ajoutées à la signature du modèle MLflow en fonction des entrées requises. Cela signifie que le schéma d’API du point de terminaison de service reflète les champs que les appelants doivent fournir au moment de l’inférence.
API d’apprentissage et d’inférence
create_training_set et score_batch calculer des valeurs de fonctionnalité correctes à la demande à partir des données sources. Pour les fonctionnalités qui prennent en charge la matérialisation hors connexion, telles que les agrégations de fenêtre glissantes sur les sources de table delta, la matérialisation des fonctionnalités d’abord dans un magasin hors connexion améliore les performances des deux opérations. Lorsque des fonctionnalités hors connexion matérialisées sont disponibles, les opérations lisent les données hors connexion précomputées au lieu de recomputer des valeurs de fonctionnalités à partir de la source. Consultez Matérialiser les vues des fonctionnalités pour matérialiser les fonctionnalités dans un magasin hors connexion.
create_training_set()
Crée un jeu de données d’entraînement avec un calcul de fonctionnalité correct dans le temps. Pour plus d’informations, consultez Entraîner des modèles avec des vues de fonctionnalités.
FeatureEngineeringClient.create_training_set(
df: DataFrame, # DataFrame with training data
features: Optional[List[Feature]], # List of Feature objects
label: Union[str, List[str], None], # Label column name(s)
exclude_columns: Optional[List[str]] = None, # Optional: columns to exclude
) -> TrainingSet
log_model()
Consigne un modèle avec des métadonnées de fonctionnalité pour le suivi de la traçabilité et la recherche automatique des fonctionnalités pendant l’inférence. Pour plus d’informations, consultez Entraîner des modèles avec des vues de fonctionnalités.
FeatureEngineeringClient.log_model(
model, # Trained model object
artifact_path: str, # Path to store model artifact
flavor: ModuleType, # MLflow flavor module (e.g., mlflow.sklearn)
training_set: TrainingSet, # TrainingSet used for training
registered_model_name: Optional[str], # Optional: register model in Unity Catalog
)
score_batch()
Effectue l’inférence de lots hors connexion avec la recherche automatique des fonctionnalités. Utilise les métadonnées de fonctionnalité stockées avec le modèle pour calculer des fonctionnalités correctes dans le temps, garantissant ainsi la cohérence avec l’entraînement.
FeatureEngineeringClient.score_batch(
model_uri: str, # URI of logged model (e.g., "models:/catalog.schema.model/1")
df: DataFrame, # DataFrame with entity keys and timestamps
) -> DataFrame
Le DataFrame d’entrée doit contenir les colonnes d’entité et de série chronologique utilisées pendant l’entraînement. Les fonctionnalités sont automatiquement calculées à partir des données sources.
fe = FeatureEngineeringClient()
# Batch scoring with automatic feature lookup
predictions = fe.score_batch(
model_uri="models:/main.ecommerce.fraud_model/1",
df=inference_df,
)
predictions.display()
Fenêtres Délai
Les vues de fonctionnalités prennent en charge trois types de fenêtres différents pour contrôler le comportement de recherche en fonction de la fenêtre de temps pour les agrégations basées sur la fenêtre de temps : propagé, bascule et glissant.
- La restauration des fenêtres revient à partir de l’heure de l’événement. La durée et le délai sont définis explicitement.
- Les fenêtres tumbling sont fixes et ne se chevauchent pas. Chaque point de données appartient exactement à une fenêtre.
- Les fenêtres glissantes sont des fenêtres temporelles mobiles qui se chevauchent, avec un intervalle de glissement configurable.
L’illustration suivante montre comment elles fonctionnent.
Fenêtre propagée
Note
RollingWindow a été précédemment nommé ContinuousWindow. Si vous migrez à partir d’une version antérieure du SDK, mettez à jour vos importations en conséquence.
Les fenêtres propagées sont up-to-date et agrégats en temps réel, généralement utilisés sur les données de streaming. Dans les pipelines de diffusion en continu, la fenêtre propagée émet une nouvelle ligne uniquement lorsque le contenu de la fenêtre de longueur fixe change, par exemple lorsqu’un événement entre ou quitte. Lorsqu’une fonctionnalité de fenêtre propagée est utilisée dans les pipelines d’apprentissage, un calcul précis de fonctionnalité à un point dans le temps est effectué sur les données sources à l’aide de la durée de la fenêtre de longueur fixe qui précède immédiatement l’horodatage d’un événement spécifique. Cela permet d’empêcher le déséquilibre en ligne-hors ligne ou les fuites de données. Les fonctionnalités au moment T agrègent les événements à partir de [T − la durée, T).
class RollingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
Le tableau suivant répertorie les paramètres d’une fenêtre propagée. Les heures de début et de fin de la fenêtre sont basées sur ces paramètres comme suit :
- Heure de début :
evaluation_time - window_duration - delay(inclusive) - Heure de fin :
evaluation_time - delay(exclusif)
| Paramètre | Contraintes |
|---|---|
delay (facultatif) |
Doit être ≥ 0 (déplace la fenêtre vers l’arrière dans le temps à partir de l’horodatage d’évaluation). Utilisez delay pour tenir compte de tout délai système entre le moment où l’événement est créé et l’horodatage de l’événement afin de prévenir les futures fuites d’événements dans les ensembles de données d'entraînement. Par exemple, s’il y a un délai d’une minute entre le moment où les événements sont créés et que ces événements sont finalement atterris dans une table source où ils sont affectés à un horodatage, le délai serait timedelta(minutes=1). |
window_duration |
Doit être > 0 |
from databricks.feature_engineering.entities import RollingWindow
from datetime import timedelta
# Look back 7 days from evaluation time
window = RollingWindow(window_duration=timedelta(days=7))
Définissez une fenêtre propagée avec un délai à l’aide du code ci-dessous.
# Look back 7 days, offset by 1 minute to account for data ingestion delay
window = RollingWindow(
window_duration=timedelta(days=7),
delay=timedelta(minutes=1)
)
Exemples de fenêtres propagées
window_duration=timedelta(days=7): cela crée une fenêtre de rétrospection de 7 jours se terminant au moment d'évaluation actuel. Pour un événement à 2h00 le jour 7, cela inclut tous les événements de 2h00 le jour 0 jusqu’à (mais pas compris) 2h00 le jour 7.window_duration=timedelta(hours=1), delay=timedelta(minutes=30): Cela crée une fenêtre de consultation de 1 heure se terminant 30 minutes avant l’heure d’évaluation. Pour un événement à 13h00, cela inclut tous les événements de 13h30 jusqu’à (mais pas inclus) 2h30. Cela est utile pour prendre en compte les retards d’ingestion des données.
Fenêtre glissante fixe
Pour les fonctionnalités définies à l'aide de fenêtres à bascule, les agrégations sont calculées sur une fenêtre de longueur fixe prédéfinie qui avance par un intervalle de glissement, produisant des fenêtres non-chevauchantes qui partitionnent entièrement le temps. Par conséquent, chaque événement de la source contribue exactement à une fenêtre. Les caractéristiques au moment t agréger les données des fenêtres terminant à ou avant t (exclusif). Windows commence à l’époque Unix.
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta
Le tableau suivant répertorie les paramètres d'une fenêtre de basculement.
| Paramètre | Contraintes |
|---|---|
window_duration |
Doit être > 0 |
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta
window = TumblingWindow(
window_duration=timedelta(days=7)
)
Exemple de fenêtre bascule
-
window_duration=timedelta(days=5): cela crée des fenêtres de longueur fixe prédéfinies de 5 jours chacune. Exemple : La fenêtre #1 s’étend sur le jour 0 à jour 4, la fenêtre #2 s’étend sur le jour 5 à le jour 9, la fenêtre #3 s’étend sur le jour 10 à la journée 14, et ainsi de suite. Plus précisément, la fenêtre n°1 inclut tous les événements avec horodatages commençant au00:00:00.00jour 0 jusqu’à (mais pas inclus) les événements avec horodatage00:00:00.00le jour 5. Chaque événement appartient exactement à une fenêtre.
Fenêtre glissante
Pour les caractéristiques définies à l'aide de fenêtres glissantes, les agrégations sont calculées sur une fenêtre de longueur fixe prédéfinie qui progresse par un intervalle de glissement, produisant des fenêtres qui se chevauchent. Chaque événement de la source peut contribuer à l’agrégation de fonctionnalités pour plusieurs fenêtres. Les caractéristiques au moment t agréger les données des fenêtres terminant à ou avant t (exclusif). Windows commence à l’époque Unix.
class SlidingWindow(TimeWindow):
window_duration: datetime.timedelta
slide_duration: datetime.timedelta
Le tableau suivant répertorie les paramètres d’une fenêtre glissante.
| Paramètre | Contraintes |
|---|---|
window_duration |
Doit être > 0 |
slide_duration |
Doit être > 0 et <window_duration |
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta
window = SlidingWindow(
window_duration=timedelta(days=7),
slide_duration=timedelta(days=1)
)
Exemple de fenêtre glissante
-
window_duration=timedelta(days=5), slide_duration=timedelta(days=1): cela crée des fenêtres superposées de 5 jours qui avancent de 1 jour chaque fois. Exemple : La fenêtre n°1 s’étend sur le jour 0 à jour 4, la fenêtre #2 s’étend sur le jour 1 à jour 5, la fenêtre #3 s’étend sur le jour 2 à la journée 6, et ainsi de suite. Chaque fenêtre inclut les événements du00:00:00.00jour de début jusqu’à (mais pas compris)00:00:00.00le jour de fin. Étant donné que les fenêtres se chevauchent, un événement unique peut appartenir à plusieurs fenêtres (dans cet exemple, chaque événement appartient à jusqu’à 5 fenêtres différentes).
Déclencheurs de matérialisation
Déclenche le contrôle lorsqu’un pipeline de matérialisation s’exécute. Le type de déclencheur dépend du type de fonctionnalité.
CronSchedule
Utiliser CronSchedule pour les fonctionnalités d’agrégation (AggregationFunction). Le pipeline s’exécute selon une planification fixe définie par une expression cron de quartz.
from databricks.feature_engineering.entities import CronSchedule
trigger = CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
)
TableTrigger
Utiliser TableTrigger pour ColumnSelection les fonctionnalités sauvegardées par un DeltaTableSource. Le pipeline s’exécute chaque fois que la table Delta en amont reçoit une nouvelle validation.
from databricks.feature_engineering.entities import TableTrigger
trigger = TableTrigger()
StreamingMode
Utiliser StreamingMode pour les fonctionnalités sauvegardées par un StreamSource. Le pipeline s’exécute en tant que pipeline de streaming continu.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
StreamSource, Feature, AggregationFunction, Sum,
RollingWindow, OnlineStoreConfig, StreamingMode,
)
from datetime import timedelta
fe = FeatureEngineeringClient()
stream_source = StreamSource(full_name="my_catalog.my_schema.my_stream")
streaming_feature = fe.create_feature(
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)),
),
catalog_name="my_catalog",
schema_name="my_schema",
name="user_purchase_sum",
)
fe.materialize_features(
features=[streaming_feature],
online_config=OnlineStoreConfig(
catalog_name="my_catalog",
schema_name="my_schema",
table_name_prefix="streaming_features_serving",
online_store_name="feature_store_online",
),
trigger=StreamingMode(),
)
Choix d’un déclencheur
| Type de fonctionnalité | Trigger | Quand il s’exécute |
|---|---|---|
Agrégation (AggregationFunction) à partir de DeltaTableSource |
CronSchedule |
Selon une planification cron fixe |
ColumnSelection (de DeltaTableSource) |
TableTrigger |
Sur chaque validation de table source |
Fonctionnalités de StreamSource |
StreamingMode |
Streaming continu |
Vous ne pouvez pas matérialiser les fonctionnalités qui nécessitent différents types de déclencheurs dans un seul materialize_features appel. Émettre des appels distincts à la place.
Migrer des fonctionnalités bêta vers la préversion publique
La préversion publique des vues de fonctionnalités introduit les entités de fonctionnalité de première classe dans le catalogue Unity, régies par les CREATE FEATURE privilèges et READ FEATURE les privilèges, et nécessite databricks-feature-engineering la version 0.16.0 ou ultérieure. Les fonctionnalités créées pendant la version bêta (avec la version 0.15.0) sont stockées en tant que fonctions de catalogue Unity et ne prennent pas en charge toutes les fonctionnalités de préversion publique. Pour obtenir la prise en charge de la préversion publique à long terme, recréez vos fonctionnalités bêta avec la version 0.16.0. Les fonctionnalités doivent être supprimées et recréées, pas seulement recréées.
Pour plus d’informations sur les fonctionnalités, consultez Vues des fonctionnalités.
Ce que vous devez faire
- Effectuez une mise à niveau vers la version 0.16.0. Il s’agit de la version cliente requise pour les fonctionnalités en préversion publique (traitement par lots et diffusion en continu).
- Recréez vos fonctionnalités. Les vues de fonctionnalité bêta doivent être supprimées et recréées, pas matérialisées, car elles ne prennent pas en charge toutes les fonctionnalités d’évaluation publique.
- Migrer avant la fermeture de la fenêtre. Les fonctionnalités bêta existantes doivent être migrées avant le 22 juillet 2026.
Identifier les fonctionnalités bêta et préversion publique
Les fonctionnalités en préversion publique apparaissent sous la forme d’un objet Feature dans le catalogue Unity, par exemple dans l’Explorateur de catalogues. Les fonctionnalités bêta apparaissent en tant que fonction avec une définition YAML. Toute fonctionnalité représentée en tant que fonction est une fonctionnalité bêta que vous devez migrer.
Migrer des fonctionnalités bêta
La migration d’une fonctionnalité bêta comporte trois parties :
- Recréez la fonctionnalité en tant que fonctionnalité en préversion publique.
- Matérialisez à nouveau la fonctionnalité, de sorte que ses tables hors connexion et en ligne sont reconstruites sous la nouvelle fonctionnalité.
- Après avoir vérifié les fonctionnalités migrées, supprimez les fonctionnalités bêta et leurs matérialisations.
Recréer les fonctionnalités
Permet list_beta_feature_views de rechercher vos fonctionnalités bêta, Feature.clone() de créer une copie non inscrite et register_feature de réinscrire chaque copie en tant que fonctionnalité en préversion publique. Le clonage efface l’inscription, le catalogue et le schéma afin que la fonctionnalité puisse être réinscrite.
Pour éviter les collisions de noms, inscrivez les fonctionnalités migrées avec un nom différent ou dans un schéma différent des fonctionnalités bêta. L’exemple suivant réinscrit chaque fonctionnalité dans son schéma d’origine avec un _migrated suffixe de nom.
# Update this to the catalog whose beta Feature Views you want to migrate.
CATALOG_TO_MIGRATE = "main"
from databricks.feature_engineering import FeatureEngineeringClient
fe = FeatureEngineeringClient()
# 1. Find every beta Feature View in the catalog. Returns Feature objects,
# scanned across all schemas in the catalog.
beta_features = fe.list_beta_feature_views(catalog_name=CATALOG_TO_MIGRATE)
# Keep each beta feature paired with its migrated counterpart for the next steps.
migrations = []
for beta_feature in beta_features:
catalog_name, schema_name, leaf_name = beta_feature.full_name.split(".")
# 2. Clone the feature as an unregistered copy, renamed with a "_migrated" suffix.
cloned = beta_feature.clone(new_name=f"{leaf_name}_migrated")
# 3. Re-register the clone as a Public Preview feature.
migrated = fe.register_feature(
feature=cloned,
catalog_name=catalog_name,
schema_name=schema_name,
)
migrations.append((beta_feature, migrated))
Matérialiser à nouveau les fonctionnalités migrées
Si une fonctionnalité bêta a été matérialisée, matérialisez à nouveau son équivalent en préversion publique afin que ses tables hors connexion et en ligne soient reconstruites sous la nouvelle fonctionnalité. Fournissez les configurations de magasin hors connexion et en ligne de la fonctionnalité migrée et reconstruisez le déclencheur à partir de la matérialisation existante de la fonctionnalité bêta.
from databricks.feature_engineering.entities import (
CronSchedule,
OfflineStoreConfig,
OnlineStoreConfig,
TableTrigger,
)
for beta_feature, migrated in migrations:
# Inspect the beta feature's existing materializations to see what to rebuild and
# to reconstruct the same trigger.
trigger = None
needs_offline = needs_online = False
for mf in fe.list_materialized_features(feature_name=beta_feature.full_name):
needs_online = needs_online or bool(mf.is_online)
needs_offline = needs_offline or not mf.is_online
# Rebuild the trigger from the materialized feature.
if mf.cron_schedule_trigger is not None:
trigger = CronSchedule(
quartz_cron_expression=mf.cron_schedule_trigger.cron_expression,
timezone_id="UTC", # Materialized schedules run in UTC.
)
elif mf.table_trigger is not None:
trigger = TableTrigger()
elif mf.streaming_mode is not None:
# Streaming features use StreamingMode, which can be reused as-is.
trigger = mf.streaming_mode
if not (needs_offline or needs_online):
continue # The beta feature was never materialized.
catalog_name, schema_name, _ = migrated.full_name.split(".")
fe.materialize_features(
features=[migrated],
offline_config=OfflineStoreConfig(
catalog_name=catalog_name,
schema_name=schema_name,
table_name_prefix="migrated_features",
)
if needs_offline
else None,
online_config=OnlineStoreConfig(
catalog_name=catalog_name,
schema_name=schema_name,
table_name_prefix="migrated_features",
online_store_name="my_online_store",
)
if needs_online
else None,
trigger=trigger,
)
Note
La matérialisation de chaque fonctionnalité dans son propre materialize_features appel crée un pipeline distinct. Pour réduire le coût de calcul, les fonctionnalités de groupe qui partagent une destination hors connexion et en ligne et se déclenchent en un seul materialize_features appel en les featurestransmettant ensemble.
Supprimer les fonctionnalités bêta
Warning
Supprimez les fonctionnalités bêta et leurs matérialisations uniquement après avoir vérifié que les fonctionnalités migrées et leurs données matérialisées sont correctes. La suppression est irréversible.
Après avoir vérifié les fonctionnalités migrées, supprimez les matérialisations de chaque fonctionnalité bêta, puis la fonctionnalité bêta elle-même.
for beta_feature, _ in migrations:
# Delete the beta feature's materializations first.
mfs = list(fe.list_materialized_features(feature_name=beta_feature.full_name))
offline_mfs = [mf for mf in mfs if not mf.is_online]
if offline_mfs:
# Aggregation features pair an offline and online table; deleting the offline
# materialized feature removes its paired online table too.
for mf in offline_mfs:
fe.delete_materialized_feature(materialized_feature=mf)
else:
# Online-only features (ColumnSelection, streaming) have no offline pair; delete
# the online materialized feature directly.
for mf in mfs:
fe.delete_materialized_feature(materialized_feature=mf)
# Then delete the beta feature definition.
fe.delete_feature(full_name=beta_feature.full_name)