Référence de l'API Feature Views
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_featureetregister_featurenécessitentCREATE FEATUREsur le schéma parent. Selon le principe du moindre privilège, accordezCREATE 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_set, et la lecture des données de fonctionnalité matérialisées pour l'entraînement ou la diffusion nécessitentREAD 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 octrois d'une fonctionnalité. La suppression d'une fonctionnalité avecdelete_feature, et la matérialisation d'une fonctionnalité avecmaterialize_featuresoudelete_materialized_feature, requièrentMANAGEsur 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.
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.
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 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.
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,StreamSourceouRequestSource).function: UnAggregationFunctionqui regroupe l'opérateur (par exemple,Sum(input="amount")), la colonne d'entrée et la fenêtre temporelle. OuColumnSelection("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, saufRequestSource. 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, saufRequestSource.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é.
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 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
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 |
|---|---|---|
| Total des valeurs | Utilisation quotidienne de l'application par utilisateur en minutes. |
| Moyenne des valeurs | Montant moyen des transactions |
| Nombre d'enregistrements | Nombre de connexions par utilisateur. |
| Valeur minimale | Fréquence cardiaque la plus basse enregistrée par un appareil portable |
| Valeur maximale | Montant maximal des transactions par session |
| Écart-type de la population | Variabilité quotidienne du montant des transactions pour tous les clients |
| Écart-type échantillon | Variabilité des taux de clics des campagnes publicitaires |
| Variance de la population | Répartition des relevés de capteurs pour les appareils IoT dans une usine |
| Variance d'échantillon | Répartition des évaluations de films sur un groupe échantillonné |
| Nombre approximatif unique | Nombre distinct d'articles achetés |
| Percentile approximatif | latence de réponse p95 |
| Première valeur | Première Timestamp de connexion |
| Dernière valeur | Montant du dernier achat |
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 |
|---|---|---|
| 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).
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.
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.
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.
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.
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 SQLWHEREappliquée avant l'agrégation. Exemple :"status = 'completed'".transformation_sql: Une expression SQLSELECTappliqué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 dedf.schema.json()). Obligatoire sitransformation_sqlest 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}.
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
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 :
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(),
)
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.
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.
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.
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 BYdans 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) :
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é :
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.
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 SQLWHEREappliqué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'").
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:
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 ( | 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 |
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.
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.
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.
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.
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être glissante
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).
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 |
|---|---|
| Doit être ≥ 0 (déplace la fenêtre vers l'arrière dans le temps à partir du timestamp d'évaluation). Utilisez |
| 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 glissante 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 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.
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta
Le tableau suivant répertorie les paramètres d'une fenêtre basculante.
parameter | Contraintes |
|---|---|
| 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 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.00le jour 0 jusqu'à (mais sans inclure) les événements avec le Timestamp00:00:00.00le 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.
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 |
|---|---|
| Doit être > 0 |
| Doit être > 0 et < |
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 de00:00:00.00le start day, jusqu'à (mais non 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).
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.
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.
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.
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 ( |
| Sur un calendrier cron fixe. |
|
| Lors de chaque commit de table source |
Fonctionnalités de |
| 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.
# 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.
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,
)
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
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)