Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
Importante
Questa funzionalità è in Anteprima Pubblica. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.
Le visualizzazioni delle funzionalità consentono di definire e calcolare le funzionalità dalle origini dati. Le funzionalità possono essere definite usando un'ampia gamma di origini (tabella Delta, flusso Kafka e dati in fase di richiesta) e calcoli (aggregazioni con intervallo di tempo, selezioni di colonne semplici e altro ancora). Questa guida illustra i flussi di lavoro seguenti:
-
Flusso di lavoro di sviluppo delle funzionalità
- Usare
create_featureper definire oggetti funzionalità del catalogo Unity che possono essere usati nel training del modello e nella gestione dei flussi di lavoro. - In alternativa, costruire
Featureoggetti in locale e usarliregister_featureper renderli persistenti in Unity Catalog in un secondo momento. Le funzionalità costruite localmente possono essere usate concreate_training_setprima della registrazione.
- Usare
-
Flusso di lavoro di addestramento del modello
- Usare
create_training_setper calcolare le funzionalità aggregate a un momento specifico per il Machine Learning. Per informazioni dettagliate sul training con visualizzazioni delle funzionalità, vedere Eseguire il training di modelli con visualizzazioni delle funzionalità.
- Usare
-
Materializzazione e gestione delle funzionalità del flusso di lavoro
- Dopo aver definito una funzionalità con
create_featureo recuperarla usandoget_feature, è possibile usarematerialize_featuresper materializzare la funzionalità o il set di funzionalità in un archivio offline per un riutilizzo efficiente o in un negozio online per la gestione online. - Usare
create_training_setcon la vista materializzata per preparare un set di dati di training batch offline.
- Dopo aver definito una funzionalità con
Per informazioni dettagliate sull'API, vedere Informazioni di riferimento sulle API delle visualizzazioni delle funzionalità.
Requisiti
Calcolo serverless o un cluster di calcolo classico che esegue Databricks Runtime 17.0 ML o versione successiva.
È necessario installare il pacchetto Python personalizzato. Eseguire le linee di codice seguenti ogni volta che si apre un notebook:
%pip install databricks-feature-engineering>=0.16.0 dbutils.library.restartPython()
Esempio di avvio rapido
Per un notebook di avvio rapido eseguibile, vedere Notebook di esempio.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
CronSchedule, DeltaTableSource, Feature, AggregationFunction,
Sum, Avg, ColumnSelection, TableTrigger,
TumblingWindow, SlidingWindow,
OfflineStoreConfig, OnlineStoreConfig,
)
from datetime import timedelta
CATALOG_NAME = "main"
SCHEMA_NAME = "feature_store"
TABLE_NAME = "transactions"
# 1. Create data source
source = DeltaTableSource(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name=TABLE_NAME,
)
# 2. Define features locally (no catalog/schema needed yet)
avg_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), TumblingWindow(window_duration=timedelta(days=30))),
name="avg_transaction_30d",
)
sum_feature = Feature(
source=source,
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), SlidingWindow(window_duration=timedelta(days=7), slide_duration=timedelta(days=1))),
# name auto-generated: "amount_sum_sliding_7d_1d"
)
fe = FeatureEngineeringClient()
# 3. Explore features with compute_features
feature_df = fe.compute_features(features=[avg_feature, sum_feature])
feature_df.display()
# 4. Create training set using local features
# `labeled_df` should have columns "user_id", "transaction_time", and "target".
training_set = fe.create_training_set(
df=labeled_df,
features=[avg_feature, sum_feature],
label="target",
)
training_set.load_df().display()
# 5. Register features in Unity Catalog
avg_feature = fe.register_feature(
feature=avg_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)
sum_feature = fe.register_feature(
feature=sum_feature,
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
)
# 6. Or use create_feature for a one-step define-and-register workflow
latest_amount = fe.create_feature(
source=source,
function=ColumnSelection("amount"),
entity=["user_id"],
timeseries_column="transaction_time",
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
name="latest_amount",
)
# 7. Train model
with mlflow.start_run():
training_df = training_set.load_df()
# training code
fe.log_model(
model=model,
artifact_path="recommendation_model",
flavor=mlflow.sklearn,
training_set=training_set,
registered_model_name=f"{CATALOG_NAME}.{SCHEMA_NAME}.recommendation_model",
)
# 8. (Optional) Materialize features for serving
# Features must be registered in UC before calling materialize_features
online_config = OnlineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features_serving",
online_store_name="customer_features_store",
)
# Aggregation features use CronSchedule and support both offline and online configs
fe.materialize_features(
features=[avg_feature, sum_feature],
offline_config=OfflineStoreConfig(
catalog_name=CATALOG_NAME,
schema_name=SCHEMA_NAME,
table_name_prefix="customer_features",
),
online_config=online_config,
trigger=CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
),
)
# ColumnSelection features use TableTrigger and only support online config
fe.materialize_features(
features=[latest_amount],
online_config=online_config,
trigger=TableTrigger(),
)
Notebook di esempio
Notebook introduttivo rapido di Feature Views
Funzionalità di streaming
Oltre alle funzionalità batch delle tabelle Delta, è possibile definire funzionalità da origini di streaming per casi d'uso in tempo reale. Le funzionalità di streaming usano la stessa classe Feature delle funzionalità batch, ovvero gli stessi Feature costruttori, le stesse funzioni di aggregazione, gli stessi flussi di lavoro di training e gestione, quindi l'aggiornamento da batch a in tempo reale richiede modifiche minime al codice. Una volta materializzate, le funzionalità di streaming garantiscono un aggiornamento end-to-end inferiore al secondo (latenza p99 di 200 ms) direttamente ai tuoi endpoint di serving del modello.
Per usare le funzionalità di streaming, configurare prima un oggetto Stream, quindi farvi riferimento usando un oggetto StreamSource. Le sorgenti di flusso supportano Kafka come input e mantengono automaticamente una tabella di ingestione (Delta) come copia storica dei dati per l'addestramento.
Definire una funzionalità di streaming
Un StreamSource fa riferimento a uno Stream tramite il suo nome in tre parti (catalog.schema.stream_name). Stream non è un oggetto a protezione diretta di Unity Catalog, ma ha come ambito uno schema del catalogo Unity e l'accesso è regolato dalla tabella di inserimento di Stream. I riferimenti alle colonne nelle definizioni di entità, serie temporali e funzioni devono essere preceduti da value. o key. per indicare quale parte del messaggio Kafka leggere. I campi annidati sono supportati usando la notazione punto (ad esempio, value.user.address.city).
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
StreamSource,
Feature,
AggregationFunction,
Sum,
RollingWindow,
)
from datetime import timedelta
client = FeatureEngineeringClient()
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
)
feature = Feature(
name="user_purchase_sum",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=AggregationFunction(
operator=Sum(input="value.amount"),
time_window=RollingWindow(window_duration=timedelta(hours=1)),
),
)
Condizioni di filtro in StreamSource
Usare filter_condition per filtrare le righe dal flusso prima dell'aggregazione, proprio come in DeltaTableSource.
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Selezione di colonne da flussi
ColumnSelection le funzionalità funzionano con le origini di streaming. La colonna selezionata rappresenta il valore più recente dello Stream per ogni entità, mantenendo l'accuratezza riferita a uno specifico momento temporale.
from databricks.feature_engineering.entities import ColumnSelection
passenger_count = Feature(
name="passenger_count",
source=stream_source,
entity=["value.user_id"],
timeseries_column="value.event_time",
function=ColumnSelection(column="value.passenger_count"),
)
Accedere ai campi annidati
È possibile accedere ai campi JSON annidati usando la notazione punto (ad esempio, value.nested_field.amount). In fase di servizio, il payload della richiesta e la risposta usano i nomi dei nodi foglia (ad esempio, amount anziché value.amount). I nomi dei nodi foglia devono essere univoci tra tutte le colonne di output di entità, serie temporali e caratteristiche all'interno di un modello o di una Specifica di caratteristiche, perché l'endpoint di serving usa i nomi dei nodi foglia per instradare i valori.
Finestre temporali per le funzionalità di streaming
Le funzionalità di streaming supportano RollingWindow solo per le aggregazioni. Le finestre mobili ricalcolano continuamente in base ai dati più recenti, il che si allinea alla natura in tempo reale delle origini dati in streaming.
TumblingWindow e SlidingWindow sono progettati per il calcolo batch su intervalli cronologici fissi.
Notebook di esempio di funzionalità di streaming
Notebook di avvio rapido sulle visualizzazioni delle funzionalità di streaming
Addestramento e inferenza del modello
Per eseguire il training dei modelli ed eseguire l'inferenza batch con visualizzazioni funzionalità, tra cui log_model(), score_batch()e create_training_set(), vedere Eseguire il training di modelli con visualizzazioni delle funzionalità.
Materializzazione delle funzionalità
Dopo aver definito le funzionalità, è possibile materializzarle in negozi offline o online per un riutilizzo efficiente nel training e nella gestione dei flussi di lavoro. Dopo aver materializzato le funzionalità, è possibile gestire i modelli usando la gestione del modello cpu. Per ulteriori dettagli, vedere Materialize Feature Views.
Procedure consigliate
Denominazione delle funzionalità
- Usare nomi descrittivi per le funzionalità business critical.
- Seguire convenzioni di denominazione coerenti tra i team.
- Usare i nomi generati automaticamente durante lo sviluppo di funzionalità.
Intervalli di tempo
- Allineare i limiti delle finestre ai cicli di business (giornalieri, settimanali).
- Le finestre più brevi acquisisce tendenze recenti, ma possono essere rumorose. Le finestre più lunghe producono distribuzioni di funzionalità più stabili, ma potrebbero non superare i recenti cambiamenti comportamentali. Scegliere in base alla velocità con cui cambia il segnale sottostante per il caso d'uso. Ad esempio, una finestra di 7 giorni riduce le fluttuazioni giornaliere e produce input del modello coerenti, mentre una finestra di 1 ora reagisce rapidamente alle modifiche comportamentali, ma potrebbe introdurre varianza che riduce le prestazioni del modello. Se l'accuratezza del modello si riduce quando si sposta la distribuzione, usare una finestra più lunga per stabilizzare gli input.
- Le finestre a cascata e a slittamento sono più scalabili rispetto alle finestre rotanti (continue). Iniziare con le finestre scorrevoli per la maggior parte dei casi d'uso.
Performance
- Materializzare le caratteristiche dalla stessa origine dati in una singola
materialize_featureschiamata per ridurre al minimo le scansioni dei dati. - Usa la stessa granularità (ad esempio, tutte le durate di 1 ora o tutte le durate di 1 giorno) per le funzionalità della stessa origine dati per consentire un miglior raggruppamento durante la materializzazione.
Colonne di entità e condizioni di filtro
Usare questa guida decisionale quando si usano le funzionalità della stessa tabella di origine:
Usare entity (in create_feature) quando sono necessari livelli di aggregazione diversi:
-
Funzionalità a livello di cliente (una riga per cliente):
entity=["customer_id"] -
Funzionalità cliente-mercante (più righe per cliente):
entity=["customer_id", "merchant_id"] -
Diversi livelli di aggregazione possono condividere lo stesso
DeltaTableSource: specificare valori diversientityin ogni definizione di funzionalità
Usare filter_condition (in DeltaTableSource) quando è necessario filtrare le righe allo stesso livello di aggregazione:
-
Solo transazioni di valore elevato:
filter_condition="amount > 100"(ancora aggregato per cliente) -
Solo ordini completati:
filter_condition="status = 'completed'"(ancora aggregato per cliente)
Regola generale: Se la modifica genera un numero diverso di righe per valore di entità, usare valori diversi entity nelle definizioni delle funzionalità. Se si filtrano semplicemente le righe che contribuiscono alla stessa aggregazione, usare filter_condition nell'origine.
Modelli comuni
Analisi dei clienti
from databricks.feature_engineering.entities import AggregationFunction, Sum, Count, RollingWindow
fe = FeatureEngineeringClient()
features = [
# Recency: Number of transactions in the last day
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=1)))),
# Frequency: transaction count over the last 90 days
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Count(input="transaction_id"), RollingWindow(window_duration=timedelta(days=90)))),
# Monetary: total spend in the last month
fe.create_feature(catalog_name="main", schema_name="ecommerce", source=transactions,
entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Sum(input="amount"), RollingWindow(window_duration=timedelta(days=30)))),
]
Analisi delle tendenze
# Compare recent vs. historical behavior
fe = FeatureEngineeringClient()
recent_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7))),
)
historical_avg = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=7), delay=timedelta(days=7))),
)
Modelli stagionali
# Same day of week, 4 weeks ago
fe = FeatureEngineeringClient()
weekly_pattern = fe.create_feature(
catalog_name="main", schema_name="ecommerce",
source=transactions, entity=["user_id"], timeseries_column="transaction_time",
function=AggregationFunction(Avg(input="amount"), RollingWindow(window_duration=timedelta(days=1), delay=timedelta(weeks=4))),
)
Limitations
- I nomi delle colonne di entità e serie temporali devono corrispondere tra il set di dati di training (etichettato) e le definizioni di caratteristiche quando utilizzate nell'API
create_training_set. - Il nome della colonna
labelusato come colonna nel set di dati di addestramento non deve esistere nelle tabelle di origine usate per la definizione diFeature. - Nell'API
create_featureè supportato un elenco limitato di funzioni (UDAFs). Vedere Funzioni supportate. - Le colonne di entità non possono essere di tipo
DATEoTIMESTAMP. -
RequestSourcesupporta solo i tipi di dati scalari definiti inScalarDataType(INTEGER,FLOAT,BOOLEANSTRING,DOUBLE,LONGTIMESTAMP,DATE).SHORTI tipi complessi, ad esempio matrici, mappe e struct, non sono supportati. -
RequestSourcenon supporta funzioni di aggregazione o finestre temporali. Solo le funzioniColumnSelectionpossono essere usate. - Il set di nomi di colonne di entità, nomi di colonne delle serie temporali e nomi di colonne delle funzionalità di richiesta devono essere univoci a livello globale in tutte le origini in un set di addestramento o in un endpoint di servizio.
-
score_batchpotrebbe non riuscire nell'ambiente di calcolo serverless. Risolvere questo problema usando un cluster di calcolo classico che esegue Databricks Runtime 17.0 ML o versione successiva.
Per le limitazioni specifiche della materializzazione, vedere Limitazioni.