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.
Controllo di accesso
Le funzionalità sono oggetti del catalogo Unity gestibili. L'accesso CREATE FEATUREa una funzionalità è controllato dai privilegi del catalogo unity , READ FEATUREe MANAGE . Per le descrizioni complete, vedere Informazioni di riferimento sui privilegi del catalogo Unity.
-
CREATE FEATURE- Obbligatorio per creare una funzionalità in uno schema.create_featureeregister_featurerichiedonoCREATE FEATUREnello schema padre. Seguendo il principio dei privilegi minimi, concedereCREATE FEATUREa livello di schema; è anche possibile concederlo a un catalogo per consentire la creazione di funzionalità in qualsiasi schema in tale catalogo. -
READ FEATURE- Obbligatorio per leggere una funzionalità e i relativi dati.get_feature,create_training_sete la lettura dei dati delle funzionalità materializzati per il training o la gestione richiedonoREAD FEATUREla funzionalità.READ FEATUREconcesso per uno schema o un catalogo si applica a tutte le funzionalità correnti e future contenute. -
MANAGE— Necessario per gestire il ciclo di vita e le concessioni di una funzionalità. L'eliminazione di una funzionalità condelete_featuree la materializzazione di una funzionalità conmaterialize_featuresodelete_materialized_feature, richiedonoMANAGEnella funzionalità .
Tutte le operazioni sulle funzionalità richiedono USE CATALOG anche nel catalogo padre e USE SCHEMA nello schema padre. Per informazioni su come MANAGE e READ FEATURE si applica alla materializzazione, vedere Autorizzazioni.
API Visualizzazione funzionalità
Feature costruttore e register_feature()
L'approccio consigliato consiste nel costruire un Feature oggetto in locale e usarlo register_feature per renderlo persistente in Unity Catalog. Questo flusso di lavoro in due passaggi consente di sperimentare le funzionalità (incluso create_training_set) prima di registrarle.
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() registra un oggetto costruito Feature localmente in 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() convalida, costruisce e registra immediatamente una funzionalità nel catalogo unity in un unico passaggio. Usare questa opzione quando non è necessario provare prima di tutto con la funzionalità in locale.
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
Parametri:
-
source: origine dati usata nel calcolo delle funzionalità (DeltaTableSource,StreamSourceoRequestSource). -
function: cheAggregationFunctionaggrega l'operatore (ad esempio,Sum(input="amount")), la colonna di input e l'intervallo di tempo insieme. OppureColumnSelection("column_name")per le funzionalità pass-through. -
catalog_name: nome del catalogo Unity per la funzionalità. -
schema_name: nome dello schema del catalogo Unity per la funzionalità. -
entity: elenco di nomi di colonna che definiscono le chiavi di aggregazione o di ricerca (chiavi primarie). Obbligatorio per tutti i tipi di origine ad eccezioneRequestSourcedi . Ad esempio,["user_id"]aggrega o cerca per utente. -
timeseries_column: colonna timestamp utilizzata per l'aggregazione dell'intervallo di tempo o per la selezione di valori più recenti. Obbligatorio per tutti i tipi di origine ad eccezioneRequestSourcedi . -
name: nome della funzionalità facoltativo. Se omesso, generato automaticamente dalla colonna di input, dalla funzione e dalla finestra (ad esempio,amount_avg_rolling_7d). -
description: Descrizione facoltativa della funzionalità.
Restituisce: Un'istanza di Feature convalidata
Genera: ValueError se fallisce una convalida
delete_feature()
Elimina una funzionalità dal catalogo unity in base al nome completo.
FeatureEngineeringClient.delete_feature(
full_name: str, # Required: '<catalog>.<schema>.<feature_name>'
) -> None
fe.delete_feature(full_name="main.store.amount_sum_rolling_7d")
Prima di eliminare una funzionalità, rimuovere o aggiornare i modelli o le specifiche di funzionalità che vi fanno riferimento. Se la funzionalità è stata materializzata, eliminare prima la caratteristica materializzata. Vedere Come eliminare una funzionalità materializzata.
Nomi generati automaticamente
Quando name viene omesso, viene generato automaticamente un nome. I nomi generati seguono il modello : {column}_{function}_{window}. Per esempio:
-
price_avg_rolling_1h(prezzo medio di 1 ora) -
transaction_count_rolling_30d_1d(conteggio di 30 giorni della transazione con ritardo 1d dal timestamp dell'evento)
Funzioni supportate
Funzione di aggregazione
Note
Le funzioni di aggregazione vengono disposte in un insieme AggregationFunction a un intervallo di tempo, come descritto nelle finestre temporali. Ogni funzione accetta un input parametro che specifica la colonna di origine da aggregare.
| Function | Description | Esempio di caso d'uso |
|---|---|---|
Sum(input="column") |
Totale dei valori | Utilizzo giornaliero delle app per utente in minuti |
Avg(input="column") |
Media dei valori | Importo medio delle transazioni |
Count(input="column") |
Numero di record | Numero di accessi per utente |
Min(input="column") |
Valore minimo | Frequenza cardiaca più bassa registrata da un dispositivo indossabile |
Max(input="column") |
Valore massimo | Quantità massima di transazioni per sessione |
StddevPop(input="column") |
Deviazione standard della popolazione | Variabilità giornaliera dell'importo delle transazioni in tutti i clienti |
StddevSamp(input="column") |
Deviazione standard di esempio | Variabilità dei tassi di click-through delle campagne pubblicitarie |
VarPop(input="column") |
Varianza della popolazione | Distribuzione delle letture dei sensori per i dispositivi IoT in una fabbrica |
VarSamp(input="column") |
Varianza campione | Distribuzione delle classificazioni dei film su un gruppo campionato |
ApproxCountDistinct(input="column", relativeSD=0.05) |
Conteggio univoco approssimativo | Conteggio distinto degli articoli acquistati |
ApproxPercentile(input="column", percentile=0.95, accuracy=100) |
Percentile approssimativo | Latenza di risposta p95 |
First(input="column") |
Primo valore | Timestamp del primo accesso |
Last(input="column") |
Ultimo valore | Importo dell'acquisto più recente |
FirstN(input="column", n=3) |
I primi n valori come array |
I primi tre prodotti visualizzati in una sessione |
LastN(input="column", n=3) |
Ultimi n valori come array |
Tre stati più recenti dei casi di supporto |
FirstDistinct(input="column", n=3) |
Innanzitutto n , valori distinti come array |
Prime tre categorie di prodotto distinte viste |
LastDistinct(input="column", n=3) |
Ultimi n valori distinti come array |
Le tre categorie di mercanti più recenti e distinte |
Note
First, Last, FirstN, LastN, FirstDistinct, e LastDistinct includono valori nulli di default. Per ignorare i valori Null, aggiungere un oggetto filter_condition che esclude in modo esplicito le colonne di input null.
FirstN, LastN, , e LastDistinct usare le timeseries_column feature per ordinare le righe di input e restituire un array contenente fino a n valori. FirstDistinct Il parametro n deve essere un intero positivo.
FirstN e FirstDistinct selezionare valori dal più antico all'ultimo livello.
LastN e LastDistinct selezionare i valori dall'ultimo al più antico, poi restituire i valori selezionati in ordine temporale.
FirstDistinct e LastDistinct rimuovere i valori duplicati mentre si selezionano i valori in quella direzione.
Ad esempio, se le righe sorgente di un'entità sono ordinate da event_time , ["A", "A", "B", "C", "B", "B"]le seguenti funzioni restituiscono:
| Function | Result |
|---|---|
FirstN(input="event_type", n=3) |
["A", "A", "B"] |
LastN(input="event_type", n=3) |
["C", "B", "B"] |
FirstDistinct(input="event_type", n=3) |
["A", "B", "C"] |
LastDistinct(input="event_type", n=3) |
["A", "C", "B"] |
FirstN, LastN, FirstDistinct, e LastDistinct richiedono databricks-feature-engineering la versione 0.17.0 o successiva.
ColumnSelection (passaggio)
ColumnSelection seleziona una singola colonna da un'origine senza applicare alcuna aggregazione. Viene eseguito il wrapping direttamente nel function parametro (non all'interno AggregationFunctiondi ). Il tipo restituito viene dedotto dallo schema di origine.
| Function | Description | Esempio di caso d'uso |
|---|---|---|
ColumnSelection("col") |
Valore più recente di una colonna (nessuna aggregazione) | Categoria fornitore più recente, pass-through di un campo richiesta |
ColumnSelection può essere usato con qualsiasi origine dati:
-
DeltaTableSource: restituisce il valore più recente per ogni chiave di entità tramite un join temporizzato (nessuna aggregazione della finestra di lookback). -
StreamSource: Restituisce l'ultimo valore per chiave di entità dal Stream (senza aggregazione a finestra di lookback). -
RequestSource: passa attraverso il valore fornito in fase di inferenza (o estratto dal dataframe etichettato in fase di training).
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",
)
Esempio: funzionalità di aggregazione e selezione di colonne
L'esempio seguente mostra le funzionalità definite sulla stessa origine dati.
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",
)
Funzionalità con condizioni di filtro
Il filter_condition parametro consente di filtrare le righe dalla tabella di origine prima di calcolare le aggregazioni. Questa funzione funge da clausola SQL WHERE applicata prima del raggruppamento e dell'aggregazione dei dati.
Note
filter_condition filtra le righe prima dell'aggregazione, ad esempio una clausola SQL WHERE applicata prima GROUP BYdi . Non modifica la granularità, che viene sempre definita dalla entity definizione della funzionalità.
I filtri sono utili quando si lavora con tabelle di origine di grandi dimensioni che includono un superset di dati necessari per il calcolo delle funzionalità e ridurre al minimo la necessità di creare viste separate sopra queste tabelle.
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))),
)
Origini dati
DeltaTableSource
DeltaTableSource è un oggetto Python temporaneo usato per definire il modo in cui le funzionalità vengono calcolate da una tabella di origine. Non crea una nuova tabella. Specifica la configurazione per la lettura dei dati e l'aggregazione delle funzionalità.
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
lateness: Optional[SourceLateness] = None, # Optional: Expected source settling time
)
Parametri:
-
catalog_name,schema_name,table_name: identificare la tabella Delta di origine nel catalogo unity. -
filter_condition: clausola SQLWHEREapplicata prima dell'aggregazione. Esempio:"status = 'completed'". -
transformation_sql: espressione SQLSELECTapplicata alla tabella di origine. Usare questa opzione per rinominare colonne, tipi di cast o calcolare colonne derivate prima dell'aggregazione. Se omesso, vengono selezionate tutte le colonne (*). Esempio:"user_id, CAST(amount AS DOUBLE) AS amount, event_time". -
dataframe_schema: schema del dataframe risultante dopo le trasformazioni, in formato JSON Spark StructType (dadf.schema.json()). Obbligatorio setransformation_sqlspecificato. Questo indica al sistema i nomi e i tipi di colonna risultanti dalla trasformazione. -
lateness: UnSourceLatenessoggetto che descrive quanto tempo la sorgente normalmente impiega a diventare completa in tempo di evento. Se omesso, la fonte è considerata immediatamente completa.
filter_condition
transformation_sql Quando e sono impostati, la query risultante è : SELECT {transformation_sql} FROM {table} WHERE {filter_condition}.
SourceLateness.settling_delay è il modo raccomandato per simulare durante l'addestramento un ritardo ETL costante che influisce sulla materializzazione online. Azure Databricks sposta indietro il tempo di valutazione dell'addestramento idoneo di questa durata, così che un esempio di addestramento non utilizzi dati che sarebbero stati ancora in transito online. Durante la materializzazione, Azure Databricks attende la stessa durata prima di pubblicare una finestra completata e serve l'ultima finestra completata durante il periodo intermedio.
Ad esempio, supponiamo che un lavoro ETL giornaliero si completi 8 ore dopo mezzanotte in un fuso orario locale dove la mezzanotte corrisponde alle 07:00 UTC. Usa un ritardo di assestamento di 8 ore e uno spostamento della finestra di 7 ore:
from datetime import timedelta
from databricks.feature_engineering.entities import (
DeltaTableSource,
SourceLateness,
TumblingWindow,
)
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="events",
lateness=SourceLateness(settling_delay=timedelta(hours=8)),
)
window = TumblingWindow(
window_duration=timedelta(days=1),
offset=timedelta(hours=7),
)
Note
Il timeseries_column deve essere di tipo TimestampType o TimestampNTZType.
DateType non è supportata per le serie temporali; cast la colonna in TimestampType first (ad esempio, con transformation_sql).
Esempio: uso di transformation_sql per le trasformazioni delle colonne
source = DeltaTableSource(
catalog_name="main",
schema_name="analytics",
table_name="raw_events",
transformation_sql="user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time",
filter_condition="event_type = 'purchase'",
dataframe_schema=spark.sql(
"SELECT user_id, CAST(price_cents AS DOUBLE) / 100 AS price, event_time FROM main.analytics.raw_events LIMIT 0"
).schema.json(),
)
Esempio: derivazione e transformation_sqldataframe_schema da un dataframe PySpark
È possibile scrivere la trasformazione come query PySpark, quindi estrarre lo schema dal dataframe risultante:
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(),
)
Espressioni supportate transformation_sql
Le stesse regole si applicano a transformation_sql su DeltaTableSource e StreamSource.
transformation_sql supporta qualsiasi espressione riga a riga; le operazioni valutate indipendentemente per ogni riga. Non cambiano il numero di righe né la corrispondenza uno a uno con la fonte. Le espressioni riga per riga includono rinomi di colonne, cast, operazioni aritmetiche e altro ancora.
Le operazioni che modificano la forma o il conteggio delle righe non sono supportate, come aggregazioni come SUM() o COUNT(). Usare AggregationFunction invece nella definizione di funzionalità.
DeltaTableSource.from_sql()
Per praticità, è possibile creare un oggetto DeltaTableSource da una query SQL. Il metodo analizza la query per estrarre automaticamente il nome della tabella, transformation_sql, e filter_condition.
DeltaTableSource.from_sql(
sql: str, # Required: SQL SELECT query
spark: SparkSession, # Required: active SparkSession (for schema inference)
) -> DeltaTableSource
Sono supportate solo query semplici SELECT ... FROM ... [WHERE ...] . L'istruzione SQL complessa (JOINs, sottoquery, CTEs, UNIONs) viene rifiutata. Per le query complesse, costruire DeltaTableSource direttamente con transformation_sql e 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",
)
Eseguire l'iterazione con to_dataframe()
Usare source.to_dataframe() per visualizzare in anteprima i dati che verranno usati per il calcolo delle funzionalità. Ciò è utile per eseguire filter_condition l'iterazione e transformation_sql fino a quando non producono i risultati previsti.
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()
Informazioni sulle entità
Le colonne di entità definiscono il livello di aggregazione per le funzionalità. Vengono specificati nella Feature definizione, non in DeltaTableSource. Le entità determinano:
-
Modalità di raggruppamento dei dati: le funzionalità vengono aggregate per ogni combinazione univoca di valori di entità (simile a
GROUP BYin SQL) - La struttura della chiave primaria: ogni combinazione di entità univoca restituisce una riga di funzionalità calcolate
Esempio: Funzionalità a livello di cliente
Il codice seguente aggrega le funzionalità a livello di cliente (una riga per cliente):
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))),
)
Esempio: funzionalità a livello di Customer Store
Per aggregare le funzionalità a un livello più dettagliato (una riga per ogni combinazione di customer-store), usare più colonne di 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))),
)
Quando sono necessarie funzionalità a diversi livelli di aggregazione (ad esempio, a livello di cliente e a livello di punto vendita del cliente), usare valori diversi entity nelle definizioni delle funzionalità. Lo stesso DeltaTableSource può essere condiviso tra le funzionalità con configurazioni di entità diverse.
StreamSource
StreamSource fa riferimento a un oggetto Stream. Stream contiene la configurazione di connessione, autenticazione, schema e inserimento per l'origine di streaming. Per Kafka, i riferimenti alle colonne nelle definizioni di funzionalità devono essere preceduti value. da o key. per indicare la parte del messaggio da leggere.
StreamSource(
full_name: str, # Required: Three-part Stream name (catalog.schema.stream)
filter_condition: Optional[str] = None, # Optional: SQL WHERE clause applied before aggregation
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
lateness: Optional[SourceLateness] = None, # Optional: Expected source settling time
)
Parametri:
-
full_name: nome completo in tre parti di un oggetto Stream (ad esempio,"my_catalog.my_schema.my_stream"). -
filter_condition(facoltativo): clausola SQLWHEREapplicata ai dati di flusso prima dell'aggregazione, usando riferimenti a colonne con prefisso punto (ad esempio,"value.event_type = 'purchase'"). -
transformation_sql(opzionale): un'espressione SQLSELECTapplicata prima dell'aggregazione o della selezione delle colonne, utilizzando riferimenti con prefisso a punti allekeystruct e.valueSupporta le stesse espressioni riga per riga diDeltaTableSource. Se omesso, la sorgente utilizza tutte le colonne (*). -
dataframe_schema: Lo schema JSON SparkStructTypedell'output proiettato. Obbligatorio se impostitransformation_sql. -
lateness: UnSourceLatenessoggetto che descrive quanto tempo il flusso normalmente impiega per completarsi in tempo di evento. VedeteSourceLateness.settling_delay.
from databricks.feature_engineering.entities import StreamSource
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
filter_condition="value.event_type = 'purchase'",
)
Si deriva dataframe_schema eseguendo la proiezione contro la tabella di ingestione del Stream, che espone le key strutture e value .
transformation_sql = (
"value.amount * value.conversion_rate AS converted_amount, "
"struct(value.user_id AS user_id, value.event_time AS time) AS event"
)
ingestion_table = "my_catalog.my_schema.events_ingestion"
dataframe_schema = spark.sql(
f"SELECT {transformation_sql} FROM {ingestion_table} LIMIT 0"
).schema.json()
stream_source = StreamSource(
full_name="my_catalog.my_schema.my_stream",
transformation_sql=transformation_sql,
dataframe_schema=dataframe_schema,
)
RequestSource
RequestSource definisce uno schema per i dati forniti in fase di inferenza nel payload della richiesta anziché ricercati da una tabella pre materializzata. Durante il training, queste colonne vengono estratte dal dataframe etichettato passato a create_training_set. Durante la gestione del modello, il chiamante deve includerli nel payload della richiesta HTTP.
RequestSource viene usato con ColumnSelection (per passare direttamente un valore). Non supporta funzioni di aggregazione o finestre temporali.
Definizione dello schema
Definire lo schema come elenco di oggetti, ognuno dei quali specifica un nome di FieldDefinition colonna e un ScalarDataTypeoggetto :
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),
]
)
Tipi di dati supportati
RequestSourcesupporta i tipi scalari definiti in ScalarDataType: INTEGER, FLOAT, BOOLEAN, STRINGDOUBLE, LONG. TIMESTAMPDATESHORT I tipi complessi come matrici, mappe e struct non sono supportati.
Come vengono idratati i dati della richiesta
| Contesto | Behavior |
|---|---|
Formazione (create_training_set) |
Le colonne vengono estratte dal dataframe etichettato. I tipi vengono convalidati rispetto allo schema dichiarato. Le mancate corrispondenze generano un errore (nessun cast implicito). |
| Gestione (endpoint modello) | Le colonne vengono estratte da dataframe_records o dataframe_split nella richiesta HTTP. I valori JSON vengono cast ai tipi dichiarati (ad esempio, numero JSON → DOUBLE). |
Firma del modello
Quando un modello viene registrato usando log_model con un set di training che include RequestSource funzionalità, le RequestSource colonne vengono aggiunte alla firma del modello MLflow come input necessari. Ciò significa che lo schema API dell'endpoint di servizio riflette i campi che i chiamanti devono fornire in fase di inferenza.
API di training e inferenza
create_training_set e score_batch calcolare i valori di funzionalità corretti temporizzato su richiesta dai dati di origine. Per le funzionalità che supportano la materializzazione offline, ad esempio le aggregazioni di finestre scorrevoli nelle origini di tabelle delta, la materializzazione delle funzionalità prima in un archivio offline migliora le prestazioni di entrambe le operazioni. Quando sono disponibili funzionalità offline materializzate, le operazioni leggono i dati offline precompilate anziché ricompilare i valori delle funzionalità dall'origine. Vedere Materialze Feature Views (Materialze Feature Views ) per materializzare le funzionalità in un negozio offline.
create_training_set()
Crea un set di dati di training con il calcolo delle funzionalità corretto temporizzato. Per informazioni dettagliate, vedere Eseguire il training di modelli con visualizzazioni delle funzionalità.
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()
Registra un modello con metadati delle funzionalità per il rilevamento della derivazione e la ricerca automatica delle funzionalità durante l'inferenza. Per informazioni dettagliate, vedere Eseguire il training di modelli con visualizzazioni delle funzionalità.
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()
Esegue l'inferenza batch offline con la ricerca automatica delle funzionalità. Usa i metadati delle funzionalità archiviati con il modello per calcolare le funzionalità corrette temporizzato, garantendo la coerenza con il training.
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
Il dataframe di input deve contenere le colonne di entità e timeeries usate durante il training. Le funzionalità vengono calcolate automaticamente dai dati di origine.
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()
Intervalli di tempo
Le Feature Views supportano quattro tipi di finestre per controllare il comportamento di lookback per aggregazioni basate su finestre temporali. I tipi di finestre disponibili dipendono dalla fonte della caratteristica:
Le funzionalità di streaming possono utilizzare finestre a rotoli e a dente di sega.
Le funzionalità a batch source possono utilizzare finestre a rotola, a rotazione e scorrevoli.
Il rollback delle finestre dall'ora dell'evento. La durata e il ritardo vengono definiti in modo esplicito.
Le finestre scorrevoli sono finestre temporali fisse e non sovrapposte. Ogni punto dati appartiene esattamente a una finestra.
Le finestre scorrevoli sono finestre temporali sovrapposte, mobili, con un intervallo di spostamento configurabile.
Le finestre Sawtooth mantengono fresca una lunga finestra di retrospettazione su una sorgente in streaming utilizzando un percorso ibrido batch e streaming. Vedi la finestra Sawtooth.
L'illustrazione seguente mostra i tipi di finestre a rotola, scivolamento, rotolante e a dente di sega.
Temporizzazione della finestra temporale
Usare delay per valutare una finestra in un momento analitico precedente. Ad esempio, una finestra di 30 giorni con un ritardo di 7 giorni calcola un valore di 30 giorni a partire da una settimana prima del periodo di valutazione.
delay è indipendente dal tempo di arrivo della fonte. Per modellare il tempo che impiega l'arrivo dei dati di origine, configura SourceLateness.settling_delay invece.
Quando entrambe le impostazioni sono presenti, compongono. Azure Databricks considera la finestra completa dopo il ritardo di assestato della sorgente e la valuta usando il ritardo analitico.
Usare offset per cambiare l'allineamento dei confini fissi delle finestre. Per impostazione predefinita, le finestre a rotazione e le finestre scorrevoli sono allineate a mezzanotte UTC. Ad esempio, uno spostamento di 22 ore allinea un limite giornaliero alle 22:00 UTC. Per approssimare i confini in un fuso orario locale, si configura uno spostamento statico rispetto all'UTC. Lo spostamento non si aggiusta per l'ora legale, non sposta il tempo di valutazione né modella i dati in ritardo.
La seguente tabella riassume il supporto per questi campi:
| Campo | Windows supportati | Constraint |
|---|---|---|
delay |
Rotola, rotolare e scivolare | Deve essere non negativo datetime.timedelta |
offset |
Rotoli e scivolamenti | Deve essere non negativo e più breve del periodo* |
SourceLateness.settling_delay |
Caratteristiche rotolanti, rotolanti e scivolanti | Deve essere non negativo datetime.timedelta |
*Punto: Per una finestra a rotazione, lo spostamento deve essere più breve di window_duration. Per una finestra scorrevole, deve essere più corta di slide_duration.
Finestra mobile
Note
RollingWindow era precedentemente denominato ContinuousWindow. Se si esegue la migrazione da una versione precedente dell'SDK, aggiornare le importazioni di conseguenza.
Le finestre in sequenza sono up-to-date e aggregazioni in tempo reale, in genere usate sui dati di streaming. Nelle pipeline di streaming, la finestra in sequenza genera una nuova riga solo quando il contenuto della finestra a lunghezza fissa cambia, ad esempio quando un evento entra o esce. Quando viene usata una funzionalità finestra in sequenza nelle pipeline di training, viene eseguito un calcolo accurato delle funzionalità temporizzato sui dati di origine usando la durata della finestra a lunghezza fissa immediatamente precedente al timestamp di un evento specifico. Ciò consente di evitare l'asimmetria online-offline o la perdita di dati. Caratteristiche al tempo T aggregare eventi da [T − durata, T).
class RollingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
Nella tabella seguente sono elencati i parametri per una finestra in sequenza. L'ora di inizio e di fine della finestra si basa su questi parametri come indicato di seguito:
- Ora di inizio:
evaluation_time - window_duration - delay(inclusivo) - Ora di fine:
evaluation_time - delay(esclusivo)
| Parametro | Vincoli |
|---|---|
delay (facoltativo) |
Deve essere ≥ 0. Sposta la finestra analitica all'indietro rispetto al timestamp della valutazione. Usalo SourceLateness.settling_delay per modellare una base coerente per il ritardo di arrivo della sorgente nel tuo flusso. |
window_duration |
Deve essere > 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))
Definire una finestra in sequenza con ritardo usando il codice seguente.
# Compute a 7-day value as of one day before the evaluation time
window = RollingWindow(
window_duration=timedelta(days=7),
delay=timedelta(days=1)
)
Esempi di finestre in sequenza
window_duration=timedelta(days=7): crea una finestra di osservazione retrospettiva di 7 giorni che termina al momento attuale della valutazione. Per un evento alle 2:00 del giorno 7, sono inclusi tutti gli eventi dalle 2:00 del giorno 0 fino (ma non incluso) alle 2:00 del giorno 7.window_duration=timedelta(hours=1), delay=timedelta(minutes=30): crea una finestra di lookback di 1 ora che termina 30 minuti prima del tempo di valutazione. Per un evento alle 15:00, sono inclusi tutti gli eventi dalle 13:30 fino alle 14:30, ma escluso le 14:30.
Finestra a cascata
Per le funzionalità definite tramite finestre a scorrimento, le aggregazioni vengono calcolate su una finestra a lunghezza fissa pre-determinata che avanza in base a un intervallo di scorrimento, producendo finestre non sovrapposte che partizionano completamente il tempo. Di conseguenza, ogni evento nell'origine contribuisce esattamente a una finestra. Le funzionalità al momento t aggregano dati dalle finestre che terminano a o prima di t in modo esclusivo. Windows inizia dall'epoca Unix.
class TumblingWindow(TimeWindow):
window_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
offset: Optional[datetime.timedelta] = None
Nella tabella seguente sono elencati i parametri per una finestra a cascata.
| Parametro | Vincoli |
|---|---|
window_duration |
Deve essere > 0 |
delay (facoltativo) |
Deve essere ≥ 0. Sposta la finestra analitica all'indietro rispetto al timestamp della valutazione. |
offset (facoltativo) |
Deve essere ≥ 0 e più breve di window_duration. Sposta i confini della finestra da mezzanotte UTC. |
from databricks.feature_engineering.entities import TumblingWindow
from datetime import timedelta
window = TumblingWindow(
window_duration=timedelta(days=1),
delay=timedelta(hours=2),
offset=timedelta(hours=22),
)
Esempio di finestra a cascata
-
window_duration=timedelta(days=5): crea finestre a lunghezza fissa predefinita di 5 giorni ciascuno. Esempio: la finestra n. 1 si estende da Giorno 0 a Giorno 4, Window #2 si estende da Giorno 5 a Giorno 9, Finestra #3 si estende da Giorno 10 a Giorno 14 e così via. In particolare, Window #1 include tutti gli eventi con timestamp a partire dal00:00:00.00giorno 0 fino a (ma non incluso) qualsiasi evento con timestamp00:00:00.00il giorno 5. Ogni evento appartiene esattamente a una finestra.
Finestra scorrevole
Per le caratteristiche definite tramite finestre scorrevoli, le aggregazioni vengono calcolate su una finestra che avanza di un intervallo di scorrimento. Una finestra scorrevole può avere una durata fissa o una durata di tutta la vita. Le finestre a durata fissa si sovrappongono, quindi ogni evento sorgente può contribuire all'aggregazione di funzionalità per più finestre. Una finestra di vita include tutti gli eventi di sorgente prima della fine della finestra. Le funzionalità al momento t aggregano dati dalle finestre che terminano a o prima di t in modo esclusivo. Windows sono allineati all'epoca Unix.
class SlidingWindow(TimeWindow):
window_duration: Optional[datetime.timedelta]
slide_duration: datetime.timedelta
delay: Optional[datetime.timedelta] = None
offset: Optional[datetime.timedelta] = None
Nella tabella seguente sono elencati i parametri per una finestra scorrevole.
| Parametro | Vincoli |
|---|---|
window_duration |
Deve essere positivo per una finestra di durata fissa. Impostato per None una finestra a vita. |
slide_duration |
Deve essere positivo. Per una finestra di durata fissa, deve anche essere più breve di window_duration. |
delay (facoltativo) |
Deve essere ≥ 0. Sposta la finestra analitica all'indietro rispetto al timestamp della valutazione. |
offset (facoltativo) |
Deve essere ≥ 0 e più breve di slide_duration. Sposta i confini della finestra da mezzanotte UTC. |
from databricks.feature_engineering.entities import SlidingWindow
from datetime import timedelta
window = SlidingWindow(
window_duration=timedelta(days=7),
slide_duration=timedelta(days=1),
delay=timedelta(hours=2),
offset=timedelta(hours=22),
)
Esempio di finestra scorrevole
-
window_duration=timedelta(days=5), slide_duration=timedelta(days=1): crea finestre sovrapposte di 5 giorni che avanzano di 1 giorno ogni volta. Esempio: la finestra n. 1 si estende da Giorno 0 a Giorno 4, Window #2 si estende da Giorno 1 a Giorno 5, Finestra #3 si estende da Giorno 2 a Giorno 6 e così via. Ogni finestra include gli eventi dal00:00:00.00giorno di inizio fino al (ma non incluso)00:00:00.00giorno di fine. Poiché le finestre si sovrappongono, un singolo evento può appartenere a più finestre (in questo esempio ogni evento appartiene a un massimo di 5 finestre diverse).
Finestra di vita
Imposta window_duration=None per creare una finestra di vita intera. A ogni confine di slide, la caratteristica aggrega tutti gli eventi sorgente per l'entità con timestamp precedenti a quel confine. Ad esempio, una diapositiva di un giorno produce un valore cumulativo una volta al giorno.
Le finestre di vita sono supportate solo da SlidingWindow.
RollingWindow e TumblingWindow richiedono un finito window_duration.
Note
Le finestre a vita richiedono una databricks-feature-engineering versione client che supporti window_duration=None l'abilitazione dello spazio di lavoro. Le versioni client precedenti non supportano questa sintassi.
from datetime import timedelta
from databricks.feature_engineering.entities import (
AggregationFunction,
DeltaTableSource,
Feature,
SlidingWindow,
Sum,
)
lifetime_spend = Feature(
source=DeltaTableSource(
catalog_name="main",
schema_name="store",
table_name="transactions",
),
entity=["user_id"],
timeseries_column="transaction_time",
function=AggregationFunction(
Sum(input="amount"),
SlidingWindow(
window_duration=None,
slide_duration=timedelta(days=1),
),
),
name="lifetime_spend",
)
Finestra a dente di sega
Importante
SawtoothWindow è in Beta.
Una finestra a dente di sega è un'aggregazione che supporta aggiornamenti molto recenti per eventi recenti, insieme alla compattazione giornaliera dei dati storici. Il suo bordo posteriore (più vecchio) avanza in passi fissi giornalieri mentre il bordo anteriore (recente) rimane aggiornato con gli ultimi eventi, quindi la lunghezza effettiva della finestra "sega" nel corso di ogni giorno. La maggior parte della finestra viene servita dai dati nella tabella di ingestione dello Stream, e solo i due giorni più recenti provengono dalla diretta stream. Questo è un compromesso che calcola in modo efficiente finestre di lunga durata (scalabili fino ad anni) pur rimanendo reattivo agli aggiornamenti recenti.
Le finestre a dente di sega si materializzano su un percorso ibrido batch e streaming. Una pipeline batch mantiene la maggior parte della finestra, mentre una pipeline streaming mantiene freschi i dati più recenti in tempo reale. I due vengono uniti in lettura, quindi per il modello o il consumatore che serve è una singola caratteristica.
Poiché la parte storica della finestra è calcolata dalla condotta batch, una caratteristica a dente di sega è pronta a servire poco dopo l'inizio della materializzazione, anche quando la finestra si estende su mesi o anni. Una finestra a rotelle è completa solo dopo che la durata completa della finestra è scaduta. Il minimo window_duration deve essere superiore a due giorni (il limite inferiore imposto), ma si raccomandano finestre a dente di sega per durate superiori a 7 giorni; per finestre più brevi, si utilizza invece una finestra a rotelle .
Note
Una caratteristica a dente di sega si basa su una storia già presente. La tabella di ingestione del Stream deve contenere dati che coprono almeno l'intera durata della finestra, altrimenti la finestra calcolata è incompleta. Prima che trascorrono 2 giorni completi, la caratteristica riflette solo i dati materializzati finora. Non è consigliato offrire il film in produzione fino a quando non sono trascorsi 2 giorni interi. Un'aggregazione su una finestra vuota restituisce 0 per Sum e Count, e nullo per Avg, Min, Max, First, LastVarPop, , VarSamp, , StddevPop, , e StddevSamp.
Per sapere se una funzione a dente di sega è pronta, apri la Feature View in Catalog Explorer. Nella sezione Caratteristiche materializzate, il backfill batch è completato una volta che il tempo di materializzazione dell'ultimo feature avanza e il suo stato mostra successo. La parte di streaming è materializzata da un gasdotto dichiarativo Lakeflow. Dopo che la Feature View supera la validazione, la funzionalità materializzata si collega a quella pipeline, dove puoi monitorarne lo stato dell'esecuzione.
Le finestre a dente di sega richiedono un StreamSource e sono materializzate con StreamingMode.
class SawtoothWindow(TimeWindow):
window_duration: datetime.timedelta
I bordi di una finestra a dente di sega si muovono in modo diverso rispetto a quelli di una finestra a rotola: il bordo d'attacco segue l'ultimo evento, mentre quello posteriore avanza una volta al giorno invece che in modo continuo. Ogni giorno, a un limite fisso delle 18:00 UTC, il bordo d'uscita avanza fino al confine UTC-mezzanotte di quel giorno. Di conseguenza, la finestra effettiva è leggermente più lunga e window_duration cresce nel corso del giorno prima di tornare un giorno al termine successivo. Addestramento e servizio utilizzano lo stesso limite delle 18:00 UTC, quindi la formazione offline e il servizio online rimangono costanti.
| Parametro | Vincoli |
|---|---|
window_duration |
Devono essere più di due giorni. È consentita una durata che non sia un intero numero di giorni (ad esempio, timedelta(days=3, minutes=15)), ma la finestra viene comunque aggiornata con granularità giornaliera. |
Le finestre a dente di sega supportano le Sumfunzioni di aggregazione , Avg, Count, Min, StddevSampVarPopMaxLastVarSampFirstStddevPope aggregazione.
Esempio di finestra a dente di sega
Il seguente esempio mostra un conteggio di 7 giorni delle transazioni di un utente. Il bordo d'attacco segue l'evento corrente mentre il bordo d'uscita avanza un giorno alla volta. Per gli eventi del 10 marzo, la finestra risale a circa il 3 marzo. Con il passare del 10 marzo, il bordo d'attacco continua ad avanzare mentre il bordo d'uscita regge, quindi la campata coperta aumenta. Poi, all'inizio dell'11 marzo, il bordo d'uscita passa a circa il 4 marzo. La finestra effettiva è sempre un po' più lunga di sette giorni. I due giorni più recenti sono serviti dalla diretta streaming, mentre i giorni precedenti dalla tabella di ingestione della Stream.
from databricks.feature_engineering.entities import SawtoothWindow
from datetime import timedelta
# 7-day window kept continuously fresh with streaming data
window = SawtoothWindow(window_duration=timedelta(days=7))
Limitazioni delle finestre a dente di sega
- Il
delayparametro non è supportato. -
SourceLateness.settling_delaynon è supportato. - Le funzioni di aggregazione diverse da
Sum,Avg,Count,MinMax,First,Last,VarPop,VarSamp, ,StddevPop, , eStddevSampnon sono supportate (ad esempio,ApproxCountDistinct,ApproxPercentile,FirstN,LastN,FirstDistinct, eLastDistinct). - Le finestre a dente di sega richiedono un
StreamSource. ADeltaTableSourcenon è supportato.
Trigger di materializzazione
Attiva il controllo quando viene eseguita una pipeline di materializzazione. Il tipo di trigger dipende dal tipo di funzionalità.
CronSchedule
Da usare CronSchedule per le funzionalità di aggregazione batch. Di default, Azure Databricks deriva un programma dalla finestra di aggregazione. Un programma derivato tiene conto del periodo della finestra, della finestra delay e offset, e della sorgente settling_delay , in modo che una run non pubblichi una finestra prima che i dati sorgente siano completi. I programmi derivati supportano finestre di tumbling e scorrevoli.
Per richiedere un programma derivato, si omette l'espressione cron.
CronSchedule() e la forma esplicita CronSchedule(mode=CronScheduleMode.DERIVED) sono equivalenti:
from databricks.feature_engineering.entities import (
CronSchedule,
CronScheduleMode,
)
trigger = CronSchedule(mode=CronScheduleMode.DERIVED)
Non impostare quartz_cron_expression con CronScheduleMode.DERIVED. Quando recuperi la funzionalità materializzata, l'orario restituito può contenere l'espressione cron calcolata da Azure Databricks.
Per controllare direttamente il programma, fornire un'espressione cron di Quarzo.
CronScheduleMode.MANUAL si deduce quando si fornisce un'espressione:
from databricks.feature_engineering.entities import CronSchedule
trigger = CronSchedule(
quartz_cron_expression="0 0 * * * ?", # Hourly
timezone_id="UTC",
)
TableTrigger
Uso TableTrigger per caratteristiche o funzionalità aggregate (ColumnSelection) supportate da un AggregationFunctionDeltaTableSource . La pipeline viene eseguita ogni volta che la tabella Delta upstream riceve un nuovo commit.
Per le funzionalità di aggregazione, la pipeline è limitata in modo da non funzionare su ogni commit. La pipeline si attiva al massimo una volta ogni metà della lunghezza della finestra della feature, ma mai più spesso di ogni 5 minuti. Ad esempio, una funzione con una finestra di rotazione di 1 ora viene eseguita al massimo una volta ogni 30 minuti, oppure una funzione con una finestra di 8 ore al massimo una volta ogni 4 ore. Il pavimento dei 5 minuti si applica quando metà della finestra è più piccola, quindi finestre di 10 minuti o meno funzionano al massimo una volta ogni 5 minuti. Le funzionalità di aggregazione la cui finestra è inferiore a 5 minuti non possono usare TableTrigger, usano invece un trigger di streaming.
from databricks.feature_engineering.entities import TableTrigger
trigger = TableTrigger()
StreamingMode
Usare StreamingMode per le funzionalità supportate da un oggetto StreamSource. La pipeline viene eseguita come pipeline di streaming continuo.
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(),
)
Scelta di un trigger
Ogni caratteristica utilizza un solo trigger; Le opzioni per tipo di caratteristica sono:
| Tipo di funzionalità | Trigger | Quando viene eseguito |
|---|---|---|
Aggregazione (AggregationFunction) da DeltaTableSource |
CronSchedule |
Su un programma derivato o manuale |
Aggregazione (AggregationFunction) da DeltaTableSource |
TableTrigger |
In ogni commit della tabella di origine |
ColumnSelection (da DeltaTableSource) |
TableTrigger |
In ogni commit della tabella di origine |
Funzionalità di StreamSource |
StreamingMode |
Streaming continuo |
Non è possibile materializzare le funzionalità che richiedono tipi di trigger diversi in una singola materialize_features chiamata. In alternativa, eseguire chiamate separate.