Aller au contenu principal

Référence de l'API Feature Views

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.

Contrôle d'accès

Les fonctionnalités sont des objets Unity Catalog gouvernables. L’accès à une fonctionnalité est contrôlé par les privilèges CREATE FEATURE, READ FEATURE et MANAGE Unity Catalog. Pour une description complète, consultez la référence des privilèges de Unity Catalog.

  • CREATE FEATURE — Requis pour créer une fonctionnalité dans un schéma. create_feature et register_feature nécessitent CREATE FEATURE sur le schéma parent. Selon le principe du moindre privilège, accordez CREATE FEATURE au 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_set, et la lecture des données de fonctionnalité matérialisées pour l'entraînement ou la diffusion nécessitent READ FEATURE sur la fonctionnalité. READ FEATURE accordé 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 octrois d'une fonctionnalité. La suppression d'une fonctionnalité avec delete_feature, et la matérialisation d'une fonctionnalité avec materialize_features ou delete_materialized_feature, requièrent MANAGE sur la fonctionnalité.

Toutes les opérations de fonctionnalités nécessitent également USE CATALOG sur le catalogue parent et USE SCHEMA sur le schéma parent. Pour savoir comment MANAGE et READ FEATURE s'appliquent à la matérialisation, consultez Autorisations.

API Feature View

Feature constructeur et register_feature()

L’approche recommandée consiste à construire un objet Feature localement et à utiliser register_feature pour le rendre persistant dans Unity Catalog. Ce workflow en deux étapes vous permet d'expérimenter les fonctionnalités (y compris create_training_set) avant de les enregistrer.

Python
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() enregistre un Feature construit localement dans Unity Catalog.

Python
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
Python
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 enregistre immédiatement une fonctionnalité dans Unity Catalog en une seule étape. Utilisez ceci lorsque vous n'avez pas besoin d'expérimenter la fonctionnalité localement en premier.

Python
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: La source de données utilisée dans le calcul des fonctionnalités (DeltaTableSource, StreamSource ou RequestSource).
  • function: Un AggregationFunction qui regroupe l'opérateur (par exemple, Sum(input="amount")), la colonne d'entrée et la fenêtre temporelle. Ou ColumnSelection("column_name") pour les fonctionnalités de transmission.
  • catalog_name: le nom du catalogue Unity Catalog pour la fonctionnalité.
  • schema_name: Le nom de schéma Unity Catalog pour la fonctionnalité.
  • entity: Liste des noms de colonnes qui définissent l'agrégation ou les clés de recherche (clés primaires). Requis pour tous les types de source, sauf RequestSource. Par exemple, ["user_id"] agrège ou recherche par utilisateur.
  • timeseries_column: La colonne timestamp utilisée pour l’agrégation de fenêtres temporelles ou la sélection de la dernière valeur. Obligatoire pour tous les types de sources, sauf RequestSource.
  • name: Nom de fonctionnalité facultatif. S'il est omis, il est 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é.

Renvoie : Une instance de Feature validée

**Génère :** ValueError si une validation échoue

delete_feature()

Supprime une fonctionnalité de Unity Catalog par son nom entièrement qualifié.

Python
FeatureEngineeringClient.delete_feature(
full_name: str, # Required: '<catalog>.<schema>.<feature_name>'
) -> None
Python
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")

Avant de supprimer une fonctionnalité, supprimez ou mettez à jour tous les modèles ou spécifications de fonctionnalité qui y font référence. Si la fonctionnalité a été matérialisée, supprimez d'abord la fonctionnalité matérialisée. Consultez Comment supprimer une fonctionnalité matérialisée.

Noms générés automatiquement

Lorsque name est omis, un nom est automatiquement généré. Les noms générés suivent le modèle : {column}_{function}_{window}. Par exemple :

  • price_avg_rolling_1h (Prix moyen sur 1 heure)
  • transaction_count_rolling_30d_1d (Nombre de transactions sur 30 jours avec un délai de 1 jour à partir du timestamp de l'événement)

Fonctions prises en charge

Fonctions d'agrégation

remarque

Les fonctions d’agrégation sont encapsulées dans un AggregationFunction avec une fenêtre temporelle, comme décrit dans les fenêtres temporelles. Chaque fonction prend un paramètre input spécifiant la colonne source à agréger.

Fonction

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 des transactions

Count(input="column")

Nombre d'enregistrements

Nombre de connexions par utilisateur.

Min(input="column")

Valeur minimale

Fréquence cardiaque la plus basse enregistrée par un appareil portable

Max(input="column")

Valeur maximale

Montant maximal des transactions par session

StddevPop(input="column")

Écart-type de la population

Variabilité quotidienne du montant des transactions pour tous les clients

StddevSamp(input="column")

Écart-type échantillon

Variabilité des taux de clics des campagnes publicitaires

VarPop(input="column")

Variance de la population

Répartition des relevés 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 approximatif unique

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

Première Timestamp de connexion

Last(input="column")

Dernière valeur

Montant du dernier achat

Fonction

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 des transactions

Count(input="column")

Nombre d'enregistrements

Nombre de connexions par utilisateur.

Min(input="column")

Valeur minimale

Fréquence cardiaque la plus basse enregistrée par un appareil portable

Max(input="column")

Valeur maximale

Montant maximal des transactions par session

StddevPop(input="column")

Écart-type de la population

Variabilité quotidienne du montant des transactions pour tous les clients

StddevSamp(input="column")

Écart-type échantillon

Variabilité des taux de clics des campagnes publicitaires

VarPop(input="column")

Variance de la population

Répartition des relevés 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 approximatif unique

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

Première Timestamp de connexion

Last(input="column")

Dernière valeur

Montant du dernier achat

remarque

First et Last incluent des valeurs nulles by default. Pour ignorer les valeurs nulles, ajoutez un filter_condition qui exclut explicitement les colonnes d'entrée qui sont nulles.

ColumnSelection (transmission)

ColumnSelection sélectionne une seule colonne d'une source sans appliquer d'agrégation. Il est directement inclus dans le paramètre function (et non dans AggregationFunction). Le type de retour est inféré du schéma source.

Fonction

Description

Exemple de cas d'usage

ColumnSelection("col")

Dernière valeur d'une colonne (pas d'agrégation)

Catégorie de fournisseur la plus récente, transmission directe d'un champ de requête

Fonction

Description

Exemple de cas d'usage

ColumnSelection("col")

Dernière valeur d'une colonne (pas d'agrégation)

Catégorie de fournisseur la plus récente, transmission directe d'un champ de requête

ColumnSelection peut être utilisé avec n'importe quelle source de données :

  • DeltaTableSource : Renvoie la dernière valeur par clé d'entité via une jointure à un instant T (pas d'agrégation de fenêtre de recherche).
  • StreamSource : Retourne la dernière valeur par clé d'entité du Stream (aucune agrégation de fenêtre rétrospective).
  • RequestSource ** ** : Transmet la valeur fournie au moment de l'inférence (ou extraite du DataFrame étiqueté au moment de l'entraînement).
Python
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 de sélection d'agrégation et de colonne

L'exemple suivant montre les fonctionnalités définies sur la même source de données.

Python
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 conditions de filtre

Le filter_condition paramètre vous permet de filtrer les lignes de la table source avant de calculer les agrégations. Ceci fonctionne comme une clause SQL WHERE qui est appliquée avant le regroupement et l'agrégation des données.

remarque

filter_condition filtre les lignes avant l'agrégation, comme une clause WHERE SQL appliquée avant GROUP BY. Cela ne modifie pas la granularité, qui est toujours définie par entity sur la définition de la fonctionnalité.

Les filtres sont utiles lorsque l’on travaille avec de grandes tables source qui incluent un sur-ensemble de données nécessaires au calcul de fonctionnalités, et minimisent le besoin de créer des vues séparées sur ces tables.

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

source de données, origine de données

DeltaTableSource

DeltaTableSource est un objet Python éphémère utilisé pour définir comment les fonctionnalités sont calculées à partir d'une table source. Il ne crée pas de nouvelle table. Il spécifie la configuration pour la lecture des données et l'agrégation des fonctionnalités.

Python
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_name: Identifiez la table Delta source dans Unity Catalog.
  • filter_condition: une clause SQL WHERE appliquée avant l'agrégation. Exemple : "status = 'completed'".
  • transformation_sql: Une expression SQL SELECT appliquée à la table source. Utilisez ceci pour renommer des colonnes, convertir des types ou compute des colonnes dérivées avant l'agrégation. Si omises, toutes les colonnes sont sélectionnées (*). Exemple : "user_id, CAST(amount AS DOUBLE) AS amount, event_time".
  • dataframe_schema: Le schéma du DataFrame résultant après transformations, au format JSON Spark StructType (issu de df.schema.json()). Obligatoire si transformation_sql est fourni. Ceci indique au système les noms et les types de colonnes qui résultent de votre transformation.

Lorsque filter_condition et transformation_sql sont définis, la query résultante est : SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.

remarque

Le timeseries_column (spécifié dans la définition de la fonctionnalité, et non dans DeltaTableSource) doit être de type TimestampType ou DateType. Les types entiers peuvent fonctionner mais entraînent une perte de précision pour les agrégats de fenêtres temporelles.

Exemple : utilisation de transformation_sql pour les transformations de colonnes

Python
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 de transformation_sql et dataframe_schema à partir d'un DataFrame PySpark

Vous pouvez écrire votre transformation sous forme de query PySpark, puis extraire le schéma du DataFrame résultant :

Python
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(),
)
remarque

transformation_sql prend en charge uniquement les expressions par ligne (renommages de colonnes, conversions de type, opérations arithmétiques). Les fonctions d'agrégation telles que COUNT(*) ou SUM() ne sont pas prises en charge. Utilisez AggregationFunction sur la définition de la fonctionnalité à la place.

DeltaTableSource.from_sql()

Par commodité, vous pouvez créer DeltaTableSource à partir d’une requête SQL. La méthode analyse la query pour extraire automatiquement le nom de la table, transformation_sql et filter_condition.

Python
DeltaTableSource.from_sql(
sql: str, # Required: SQL SELECT query
spark: SparkSession, # Required: active SparkSession (for schema inference)
) -> DeltaTableSource

Seules les SELECT ... FROM ... [WHERE ...] query simples sont prises en charge. Le SQL complexe (JOINs, sous-requêtes, CTEs, UNIONs) est rejeté. Pour les requêtes complexes, construisez DeltaTableSource directement avec transformation_sql et filter_condition.

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

Utilisez source.to_dataframe() pour prévisualiser les données qui seront utilisées pour le calcul de fonctionnalités. C'est utile pour itérer sur filter_condition et transformation_sql jusqu'à ce qu'ils produisent les résultats attendus.

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

Comprendre les entités

Les colonnes d'entité définissent le niveau d'agrégation de vos fonctionnalités. Ils sont spécifiés sur la définition Feature, 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é (similaire à GROUP BY dans SQL)
  • La structure de la clé primaire : chaque combinaison d'entités unique donne 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) :

Python
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 des clients et des magasins

Pour agréger les fonctionnalités à un niveau plus détaillé (une ligne par combinaison client-magasin), utilisez plusieurs colonnes d'entité :

Python
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 clients et au niveau du magasin du clients), utilisez différentes valeurs entity dans vos définitions de fonctionnalités. Le même DeltaTableSource peut être partagé entre des fonctionnalités avec différentes configurations d'entités.

StreamSource

StreamSource fait référence à un Stream. Le Stream contient la configuration de connexion, d'authentification, de schéma et d'ingestion pour la source de streaming. Pour Kafka, les références de colonne dans les définitions de fonctionnalités doivent être préfixées par value. ou key. pour indiquer quelle partie du message lire.

Python
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: le nom complet en trois parties d'un Stream (par exemple, "my_catalog.my_schema.my_stream").
  • filter_condition (facultatif) : Une clause SQL WHERE appliquée aux données de stream avant l'agrégation, en utilisant des références de colonne précédées d'un point (par exemple, "value.event_type = 'purchase'").
Python
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 qui sont fournies au moment de l'inférence dans la charge utile de la requête plutôt que recherchées à partir d'une table pré-matérialisée. Pendant l'entraînement, ces colonnes sont extraites du DataFrame étiqueté transmis à create_training_set. Lors du service de modèles, l'appelant doit les inclure dans la charge utile de la requête HTTP.

RequestSource est utilisé avec ColumnSelection (pour transmettre une valeur directement). Il ne prend pas en charge les fonctions d’agrégation ou les fenêtres temporelles.

Définition du schéma

Définissez le schéma comme une liste d'objets FieldDefinition, chacun spécifiant un nom de colonne et un ScalarDataType:

Python
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 pris en charge

RequestSource prend en charge les types 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.

Comment les données de requête sont hydratées

Contexte

Comportement

Formation (create_training_set)

Les colonnes sont extraites du DataFrame étiqueté. Les types sont validés par rapport au schéma déclaré. Les incohérences entraînent une erreur (pas de conversion implicite).

Déploiement (Endpoint de modèle)

Les colonnes sont extraites de dataframe_records ou dataframe_split dans la requête HTTP. Les valeurs JSON sont transtypées vers les types déclarés (par ex. : nombre JSON → DOUBLE).

Contexte

Comportement

Formation (create_training_set)

Les colonnes sont extraites du DataFrame étiqueté. Les types sont validés par rapport au schéma déclaré. Les incohérences entraînent une erreur (pas de conversion implicite).

Déploiement (Endpoint de modèle)

Les colonnes sont extraites de dataframe_records ou dataframe_split dans la requête HTTP. Les valeurs JSON sont transtypées vers les types déclarés (par ex. : nombre JSON → DOUBLE).

Signature du modèle

Lorsqu'un modèle est enregistré à l'aide de log_model avec un ensemble d'entraînement qui inclut RequestSource fonctionnalités, les colonnes RequestSource sont ajoutées à la signature du modèle MLflow comme entrées requises. Cela signifie que le schéma de l'API de l'Endpoint de diffusion reflète les champs que les appelants doivent fournir au moment de l'inférence.

API de formation et d’inférence

create_training_set et score_batch compute les valeurs de fonctionnalité correctes à un moment précis à la demande à partir des données sources. Pour les fonctionnalités qui prennent en charge la matérialisation hors ligne, telles que les agrégations de fenêtres glissantes sur les sources de tables Delta, la matérialisation des fonctionnalités d'abord vers un magasin hors ligne améliore les performances des deux opérations. Lorsque les fonctionnalités hors ligne matérialisées sont disponibles, les opérations lisent les données hors ligne précalculées au lieu de recalculer les valeurs des fonctionnalités à partir de la source. Consultez Matérialiser les vues de fonctionnalités pour matérialiser les fonctionnalités vers un magasin hors ligne.

create_training_set()

Crée un dataset d'entraînement avec un calcul de fonctionnalités exact à un instant T. Pour plus de détails, consultez Former des modèles avec des vues de fonctionnalités.

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

Logs un modèle avec les métadonnées de fonctionnalité pour le suivi de la traçabilité et la recherche automatique de fonctionnalités pendant l’inférence. Pour plus de détails, consultez Former des modèles avec des vues de fonctionnalités.

Python
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 batch hors ligne avec recherche automatique de fonctionnalités. Utilise les métadonnées de fonctionnalité stockées avec le modèle pour calculer des fonctionnalités correctes à un instant donné, garantissant la cohérence avec l’entraînement.

Python
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éries temporelles utilisées pendant l'entraînement. Les fonctionnalités sont automatiquement calculées à partir des données source.

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

Les Feature Views prennent en charge trois types de fenêtres différents pour contrôler le comportement de rétrospection pour les agrégations basées sur des fenêtres temporelles : glissantes, chevauchantes et coulissantes.

  • Les fenêtres glissantes prennent en compte l’historique à partir de l’heure de l’événement. La durée et le délai sont explicitement définis.
  • Les fenêtres bascules sont des fenêtres de temps fixes et non chevauchantes. Chaque point de données appartient à exactement une fenêtre.
  • Les fenêtres glissantes sont des fenêtres temporelles chevauchantes et continues avec un intervalle de glissement configurable.

Le schéma suivant montre comment ils fonctionnent.

Fenêtres de rétrospection glissantes, évolutives et défilantes.

Fenêtre glissante

remarque

RollingWindow a été précédemment nommé ContinuousWindow. Si vous migrez depuis une version antérieure du SDK, mettez à jour vos importations en conséquence.

Les fenêtres glissantes sont des agrégats à jour et en temps réel, généralement utilisés sur les données en streaming. Dans les pipelines de streaming, la fenêtre glissante émet une nouvelle ligne uniquement lorsque le contenu de la fenêtre de longueur fixe change, par exemple lorsqu'un événement entre ou sort. Lorsqu'une fonctionnalité de fenêtre glissante est utilisée dans les pipelines d'entraînement, un calcul précis des fonctionnalités à un instant T est effectué sur les données sources en utilisant la durée de la fenêtre de longueur fixe immédiatement précédant le Timestamp d'un événement spécifique. Cela permet d'éviter l'asymétrie en ligne-hors ligne ou la fuite de données. Les fonctionnalités à l'heure T agrègent les événements de [T − durée, T).

Python
class RollingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None

Le tableau suivant répertorie les parameters d’une fenêtre glissante. Les heures de start et de fin de fenêtre sont basées sur ces parameters comme suit :

  • Heure de start : evaluation_time - window_duration - delay (inclusive)
  • Heure de fin : evaluation_time - delay (exclusive)

parameter

Contraintes

delay (facultatif)

Doit être ≥ 0 (déplace la fenêtre vers l'arrière dans le temps à partir du timestamp d'évaluation). Utilisez delay pour tenir compte de tout délai système entre le moment où l’événement est créé et son timestamp afin d’éviter toute fuite d’événements futurs dans les datasets 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 le moment où ces événements arrivent dans une table source où ils se voient attribuer un Timestamp, alors le délai serait de timedelta(minutes=1).

window_duration

Doit être > 0

parameter

Contraintes

delay (facultatif)

Doit être ≥ 0 (déplace la fenêtre vers l'arrière dans le temps à partir du timestamp d'évaluation). Utilisez delay pour tenir compte de tout délai système entre le moment où l’événement est créé et son timestamp afin d’éviter toute fuite d’événements futurs dans les datasets 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 le moment où ces événements arrivent dans une table source où ils se voient attribuer un Timestamp, alors le délai serait de timedelta(minutes=1).

window_duration

Doit être > 0

Python
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 glissante avec un délai à l'aide du code ci-dessous.

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

  • window_duration=timedelta(days=7): cela crée une fenêtre rétrospective de 7 jours se terminant à l'heure d'évaluation actuelle. Pour un événement à 14 h 00 le Jour 7, cela inclut tous les événements de 14 h 00 le Jour 0 jusqu'à (mais sans l'inclure) 14 h 00 le Jour 7.

  • window_duration=timedelta(hours=1), delay=timedelta(minutes=30): Cela crée une fenêtre rétrospective d'une heure se terminant 30 minutes avant l'heure d'évaluation. Pour un événement à 15 h 00, cela inclut tous les événements de 13 h 30 jusqu'à (mais non compris) 14 h 30. C'est utile pour tenir compte des délais d'ingestion des données.

Fenêtre basculante

Pour les fonctionnalités définies à l'aide de fenêtres basculantes, les agrégations sont calculées sur une fenêtre de longueur fixe prédéterminée qui avance par un intervalle de glissement, produisant des fenêtres non chevauchantes qui partitionnent entièrement le temps. En conséquence, chaque événement dans la source contribue à une seule fenêtre. Les fonctionnalités au temps t agrègent les données des fenêtres se terminant à ou avant t (exclue). Windows commencent à l'époque Unix.

Python
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta

Le tableau suivant répertorie les paramètres d'une fenêtre basculante.

parameter

Contraintes

window_duration

Doit être > 0

parameter

Contraintes

window_duration

Doit être > 0

Python
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta

window = TumblingWindow(
window_duration=timedelta(days=7)
)

Exemple de fenêtre glissante

  • window_duration=timedelta(days=5): Cela crée des fenêtres de longueur fixe prédéterminées de 5 jours chacune. Exemple : la fenêtre n°1 s'étend du jour 0 au jour 4, la fenêtre n°2 du jour 5 au jour 9, la fenêtre n°3 du jour 10 au jour 14, et ainsi de suite. Plus précisément, la fenêtre n° 1 inclut tous les événements avec des Timestamp commençant à 00:00:00.00 le jour 0 jusqu'à (mais sans inclure) les événements avec le Timestamp 00:00:00.00 le jour 5. Chaque événement appartient à une seule fenêtre.

Fenêtre glissante

Pour les fonctionnalités 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 avance par un intervalle de glissement, produisant des fenêtres qui se chevauchent. Chaque événement dans la source peut contribuer à l'agrégation des fonctionnalités pour plusieurs fenêtres. Les fonctionnalités au temps t agrègent les données des fenêtres se terminant à ou avant t (exclue). Windows commencent à l'époque Unix.

Python
class SlidingWindow(TimeWindow):
window_duration: datetime.timedelta
slide_duration: datetime.timedelta

Le tableau suivant répertorie les paramètres d'une fenêtre glissante.

parameter

Contraintes

window_duration

Doit être > 0

slide_duration

Doit être > 0 et < window_duration

parameter

Contraintes

window_duration

Doit être > 0

slide_duration

Doit être > 0 et < window_duration

Python
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 de 5 jours qui se chevauchent et avancent d'1 jour à chaque fois. Exemple : la fenêtre n° 1 s'étend du jour 0 au jour 4, la fenêtre n° 2 s'étend du jour 1 au jour 5, la fenêtre n° 3 s'étend du jour 2 au jour 6, et ainsi de suite. Chaque fenêtre inclut les événements à partir de 00:00:00.00 le start day, jusqu'à (mais non compris) 00:00:00.00 le 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).

Materialization triggers

Les déclencheurs déterminent le moment où un pipeline de matérialisation s'exécute. Le type de Trigger dépend du type de fonctionnalité.

CronSchedule

Utilisez CronSchedule pour les fonctionnalités d'agrégation (AggregationFunction). Le pipeline s'exécute selon un calendrier fixe défini par une expression cron Quartz.

Python
from databricks.feature_engineering.entities import CronSchedule

trigger = CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
)

TableTrigger

Utilisez TableTrigger pour les fonctionnalités ColumnSelection prises en charge par un DeltaTableSource. Le pipeline s'exécute chaque fois que la table Delta en amont reçoit un nouveau commit.

Python
from databricks.feature_engineering.entities import TableTrigger

trigger = TableTrigger()

StreamingMode

Utilisez StreamingMode pour les fonctionnalités prises en charge par un StreamSource. Le pipeline s'exécute en tant que pipeline de streaming continu.

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

Choosing a Trigger

Type de fonctionnalité

Déclencheur

Lors de l'exécution

Agrégation (AggregationFunction) de DeltaTableSource

CronSchedule

Sur un calendrier cron fixe.

ColumnSelection (à partir de DeltaTableSource)

TableTrigger

Lors de chaque commit de table source

Fonctionnalités de StreamSource

StreamingMode

Streaming continu

Type de fonctionnalité

Déclencheur

Lors de l'exécution

Agrégation (AggregationFunction) de DeltaTableSource

CronSchedule

Sur un calendrier cron fixe.

ColumnSelection (à partir de DeltaTableSource)

TableTrigger

Lors de chaque commit de table source

Fonctionnalités de StreamSource

StreamingMode

Streaming continu

Vous ne pouvez pas matérialiser de fonctionnalités qui nécessitent différents types de Trigger dans un seul appel materialize_features. Lancez plutôt des appels séparés.

Migrer les fonctionnalités bêta vers l'Aperçu public

L'aperçu public de Feature Views introduit des entités Feature de premier ordre dans Unity Catalog, régies par les privilèges CREATE FEATURE et READ FEATURE, et nécessite la version databricks-feature-engineering 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 Unity Catalog et ne prennent pas en charge toutes les fonctionnalités de l'aperçu public. Pour obtenir une prise en charge à long terme de l'aperçu public, recréez vos fonctionnalités bêta avec la version 0,16,0. Les fonctionnalités doivent être supprimées et recréées, et non simplement rematérialisées.

Pour plus d'information sur les fonctionnalités, consultez les Vues de fonctionnalités.

Ce que vous devez faire

  • Mettre à niveau vers la version 0.16.0. Il s'agit de la version client requise pour les fonctionnalités en Public Preview (batch et streaming).
  • Recréer vos fonctionnalités. Les vues de fonctionnalités bêta doivent être supprimées et recréées, et non rematérialisées, car elles ne prennent pas en charge toutes les fonctionnalités de prévisualisation publique.
  • Migrez 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 de prévisualisation publique

Les fonctionnalités en Public Preview apparaissent comme un objet Feature dans Unity Catalog, par exemple dans l’Explorateur de catalogues. Les fonctionnalités bêta apparaissent comme une fonction avec une définition YAML. Toute fonctionnalité représentée sous forme de fonction est une fonctionnalité en version bêta que vous devez migrer.

Migrer les fonctionnalités bêta

La migration d’une fonctionnalité bêta comporte trois parties :

  • Recréez la fonctionnalité en tant que fonctionnalité de préversion publique.
  • Rematérialisez la fonctionnalité, afin que ses tables hors ligne et en ligne soient 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éez les fonctionnalités.

Utilisez list_beta_feature_views pour trouver vos fonctionnalités bêta, Feature.clone() pour créer une copie non enregistrée, et register_feature pour réenregistrer chaque copie en tant que fonctionnalité d'aperçu public. Le clonage efface l'enregistrement, le catalogue et le schéma afin que la fonctionnalité puisse être réenregistrée.

Pour éviter les conflits de noms, enregistrez les fonctionnalités migrées avec un nom différent ou dans un schéma différent de celui des fonctionnalités bêta. L'exemple suivant réenregistre chaque fonctionnalité dans son schéma d'origine avec un suffixe de nom _migrated.

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

Re-matérialiser les fonctionnalités migrées

Si une fonctionnalité bêta a été matérialisée, re-matérialisez son équivalent en préversion publique afin que ses tables hors ligne et en ligne soient reconstruites sous la nouvelle fonctionnalité. Fournissez les configurations de stockage hors ligne et en ligne pour la fonctionnalité migrée, et reconstituez le trigger à partir de la matérialisation existante de la fonctionnalité bêta.

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

La matérialisation de chaque fonctionnalité dans son propre appel materialize_features crée un pipeline distinct. Pour réduire les coûts de compute, regroupez les fonctionnalités qui partagent une destination hors ligne et en ligne et un trigger en un seul appel materialize_features en les transmettant ensemble dans features.

Supprimer les fonctionnalités bêta

attention

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.

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