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
La creazione e l'aggiornamento di viste materializzate autonome e tabelle di streaming da un notebook in un ambiente di calcolo generale serverless sono disponibili in beta e disponibili in aree selezionate. Vedere Notebook.
È possibile creare e aggiornare viste materializzate autonome e tabelle di streaming da un notebook usando Python. In questo modo è possibile gestire pipeline autonome insieme agli altri flussi di lavoro basati su notebook basati su Python.
A questo scopo è possibile procedere in due modi:
- Definisci la tabella con i
pyspark.pipelinesdecoratori,@dp.materialized_viewe@dp.table. Usalo quando la logica è più facile da esprimere come codice DataFrame. Vedi Definire tabelle con i decoratori delle pipeline. - Invia le stesse istruzioni SQL che un SQL warehouse di Databricks esegue passandole a
spark.sql(). In questo modo ottieni la vista materializzata autonoma completa e l’interfaccia SQL delle tabelle in streaming, incluse le istruzioniREFRESHe le pianificazioni di aggiornamento. Vedi Invia istruzioni SQL conspark.sql().
Il codice sorgente Python per le pipeline standalone richiede un notebook collegato a calcolo generico serverless. Non è possibile usare Python per creare o aggiornare pipeline autonome da un databricks SQL Warehouse, perché un warehouse esegue istruzioni SQL, non Python notebook. Per usare invece un SQL warehouse, consulta Usare viste materializzate autonome e Usare tabelle autonome in streaming.
Requisiti
Per creare e aggiornare pipeline autonome con Python, è necessario un notebook collegato al calcolo serverless per uso generale su Databricks Runtime 18.1 o superiore. Per l'elenco completo dei requisiti, tra cui disponibilità e autorizzazioni a livello di area, vedere Notebook.
Definisci tabelle con i decoratori di pipeline
Puoi definire una visualizzazione materializzata autonoma o una tabella streaming con gli stessi decoratori che usi in una pipeline Lakeflow. Ogni funzione decorata definisce una tabella. Quando esegui una cella, Azure Databricks crea la tabella ed esegue una pipeline serverless per popolarla. La cella ritorna quando l'aggiornamento termina.
Avvertimento
I decoratori delle pipeline richiedono la versione 5 o successiva dell'ambiente serverless.
Definisci una visione materializzata
Usa @dp.materialized_view su una funzione che restituisce un DataFrame batch. Il seguente esempio crea la vista materializzata daily_booking_revenue dalla tabella bookings nel set di dati di esempio Wanderbricks:
from pyspark import pipelines as dp
from pyspark.sql import functions as F
@dp.materialized_view(name="main.default.daily_booking_revenue")
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
Per definire una tabella a partire da una lettura in streaming, usa @dp.table invece.
Definisci una tabella di streaming
Usa @dp.table su una funzione che restituisce un DataFrame in streaming. L'esempio seguente crea la tabella di streaming bookings_raw da una lettura in streaming della stessa tabella bookings:
from pyspark import pipelines as dp
@dp.table(name="main.default.bookings_raw")
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")
Se la funzione restituisce un DataFrame batch, @dp.table crea invece una vista materializzata. L'unica eccezione è replace_where, che porta sempre a una tabella di streaming. Il seguente esempio mantiene aggiornati i ricavi giornalieri per i check-in in o dopo il 1° luglio 2025, senza ricalcolare le date precedenti:
from pyspark import pipelines as dp
from pyspark.sql import functions as F
@dp.table(
name="main.default.booking_revenue_rw",
replace_where=F.col("check_in") >= F.to_date(F.lit("2025-07-01")),
)
def booking_revenue_rw():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
A ogni esecuzione vengono eliminate le righe corrispondenti al predicato e viene ricalcolato solo tale intervallo. Vedi Elaborazione batch con flussi REPLACE WHERE.
Aggiorna una tabella
Per aggiornare una tabella definita con un decoratore, esegui di nuovo il codice che la define, ad esempio rieseguendo la cella del notebook, eseguendo l'intero notebook o eseguendo il notebook come un job. Ogni run crea la tabella se non esiste e la aggiorna se esiste.
Per rielaborare tutti i dati disponibili nella fonte, passa full_refresh=True a uno dei due decoratori:
@dp.table(name="main.default.bookings_raw", full_refresh=True)
def bookings_raw():
return spark.readStream.table("samples.wanderbricks.bookings").select("booking_id", "property_id", "total_amount")
Non puoi usare un'affermazione REFRESH su una tabella definita con un decoratore, né aggiornare l'agenda con SCHEDULE o TRIGGER ON UPDATE. Per aggiornare in modo pianificato, definisci la tabella in SQL oppure pianifica il notebook come processo. Consulta Attività di Lakeflow.
Configura la tabella
I decoratori accettano gli stessi parametri comuni del dataset che accettano all'interno di una pipeline, inclusi comment, table_properties, partition_cols, cluster_by, schema, e spark_conf:
@dp.materialized_view(
name="main.default.daily_booking_revenue",
comment="Daily booking revenue.",
table_properties={"quality": "gold"},
cluster_by=["check_in"],
)
def daily_booking_revenue():
return (
spark.read.table("samples.wanderbricks.bookings")
.groupBy("check_in")
.agg(F.sum("total_amount").alias("total_revenue"))
)
Per l'elenco dei parametri, vedi materialized_view e tabella.
private=True non è supportata, perché una tabella privata può essere letta solo da altri dataset nella stessa pipeline.
API non supportate
Una tabella standalone è un singolo dataset con un unico flusso, quindi le API che descrivono le relazioni tra dataset non sono disponibili. I seguenti generano un errore al di fuori di una pipeline:
-
@dp.temporary_viewedp.create_streaming_table -
@dp.append_flowe altri flussi aggiuntivi -
dp.create_auto_cdc_flowedp.create_auto_cdc_from_snapshot_flow -
@dp.replace_flowe il parametroreplace_using, che definiscono i flussi REPLACE USING. Consultare Sostituzione parziale dello snapshot con i flussi REPLACE USING. dp.create_sink- Aspettative, come
@dp.expecte@dp.expect_or_fail
Per usarli, crea invece una pipeline Lakeflow. Vedere Sviluppare codice della pipeline con Python.
Invia istruzioni SQL con spark.sql()
In un notebook Python, passa a spark.sql() le stesse istruzioni che eseguiresti da un warehouse SQL di Databricks. La sintassi della vista materializzata autonoma e della tabella di streaming è identica; solo il modo in cui si invia l'istruzione è diversa. Come in un data warehouse, ogni istruzione CREATE o REFRESH esegue una pipeline serverless per elaborare l'operazione.
La spark sessione è disponibile per impostazione predefinita nei notebook Azure Databricks, quindi non è necessaria alcuna importazione.
Creare una vista materializzata
L'esempio seguente crea la vista materializzata mv1 dalla tabella di base base_table1.
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW mv1
AS SELECT
date,
sum(sales) AS sum_of_sales
FROM base_table1
GROUP BY date
""")
Per informazioni dettagliate complete CREATE MATERIALIZED VIEW , ad esempio aggiornamenti pianificati e attivati, vedere Creare una vista materializzata.
Creare una tabella di streaming
L'esempio seguente crea la tabella di streaming sales dalla tabella raw_data:
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")
Per informazioni dettagliate CREATE STREAMING TABLE , incluso il caricamento di file con il caricatore automatico e la pianificazione, vedere Usare tabelle di streaming autonome.
Aggiornare una vista materializzata o una tabella di streaming
Usare un'istruzione REFRESH per aggiornare una tabella autonoma con i dati più recenti dell'origine:
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")
Nel calcolo generico serverless, gli aggiornamenti sono sincroni. Gli aggiornamenti asincroni (parola ASYNC chiave) non sono supportati. Vedi calcolo generico serverless.
Parametrizzare le istruzioni
Per passare valori dal codice Python in un'istruzione invece di codificarli direttamente, usa marcatori di parametro con nome nell'SQL e fornisci i relativi valori tramite l'argomento args di spark.sql(). Usa un marcatore come :min_sales direttamente per i valori letterali. Racchiudere il marcatore in IDENTIFIER() solo quando il parametro è il nome di un oggetto, come una tabella, una vista o uno schema, perché gli identificatori non possono essere sostituiti come semplici valori stringa.
Nell'esempio seguente vengono parametrizzati sia il nome della vista materializzata che un valore di filtro:
mv_name = "main.sales.regional_sales"
min_sales = 1000
spark.sql("""
CREATE OR REPLACE MATERIALIZED VIEW IDENTIFIER(:mv)
AS SELECT
region,
sum(sales) AS sum_of_sales
FROM base_table1
WHERE sales > :min_sales
GROUP BY region
""", args={
"mv": mv_name,
"min_sales": min_sales,
})
Per altre informazioni, vedere Parametri marcatori e IDENTIFIER clausola.
Esegui altre istruzioni
È possibile eseguire qualsiasi vista materializzata autonoma o istruzione di tabella di streaming da un notebook di Python passandola a spark.sql(), incluse le istruzioni per pianificare gli aggiornamenti, modificare una tabella o eliminare una tabella. Per informazioni su come usare viste materializzate e tabelle di streaming, inclusa la sintassi SQL, vedere Usare viste materializzate autonome e Usare tabelle di streaming autonome.
Limitazioni
Le viste materializzate autonome e le tabelle di streaming create nel calcolo generale serverless presentano limitazioni aggiuntive, ad esempio nessun supporto per gli aggiornamenti asincroni e nessuna attribuzione dei costi per tabella. Per l'elenco completo, vedere Calcolo generale serverless.
Poiché queste pipeline vengono eseguite su elaborazione general-purpose serverless anziché su un data warehouse SQL, non ereditano i tag personalizzati dal data warehouse di riferimento. La propagazione dei tag warehouse a system.billing.usage si applica solo alle visualizzazioni materializzate e alle tabelle di streaming le cui istruzioni vengono eseguite da un SQL warehouse. Vedi Attribuire i costi al warehouse SQL con tag personalizzati.
Le tabelle definite con i decoratori di pipeline presentano le seguenti limitazioni aggiuntive:
- Non puoi aggiornarli con una
REFRESHdichiarazione o aggiornare il programma conSCHEDULEoTRIGGER ON UPDATE. Vedi Aggiorna una tabella. - Aspettative, flussi aggiuntivi, flussi di raccolta dati di cambiamento (CDC), pozzi e viste temporanee non sono supportati. Vedere le API non supportate.
-
private=Truenon è supportato.