Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Pomocí Python můžete vytvářet a aktualizovat samostatná materializovaná zobrazení a streamované tabulky z poznámkového bloku. To vám umožní spravovat samostatné pipeline vedle ostatních pracovních postupů notebooků založených na Pythonu.
Můžete to udělat dvěma způsoby:
- Definujte tabulku pomocí dekorátorů
pyspark.pipelines,@dp.materialized_viewa@dp.table. Použijte to, když je logika snazší vyjádřit jako kód v DataFrame. Viz Definice tabulek pomocí dekorátorů pipeline. - Odešlete stejné příkazy SQL, které spouští sklad Databricks SQL, jejich předáním do
spark.sql(). Získáte tak úplné samostatné rozhraní SQL pro materializovaný pohled i streamovací tabulku, včetně příkazůREFRESHa plánů aktualizací. Podívejte se na Odeslání příkazů SQL pomocíspark.sql().
Zdrojový kód v Pythonu pro samostatné pipeline vyžaduje notebook připojený k obecnému bezserverovému výpočetnímu prostředí. Nemůžete použít Python k vytvoření nebo aktualizaci samostatných kanálů ze služby Databricks SQL Warehouse, protože sklad spouští příkazy SQL, ne Python poznámkové bloky. Pokud chcete místo toho použít SQL Warehouse, přečtěte si téma Použití samostatných materializovaných zobrazení a použití samostatných streamovacích tabulek.
Důležité
Vytváření a aktualizace samostatných materializovaných zobrazení a streamovaných tabulek z poznámkového bloku v bezserverovém obecném výpočetním prostředí je v beta verzi a k dispozici ve vybraných oblastech. Viz Poznámkové bloky.
Requirements
K vytváření a obnovování samostatných pipeline v Pythonu potřebujete poznámkový blok připojený k bezserverovému obecnému výpočetnímu prostředí v Databricks Runtime 18.1 nebo novějším. Úplný seznam požadavků, včetně regionální dostupnosti a oprávnění, najdete v poznámkových blocích.
Definujte tabulky pomocí dekorátorů kanálů
Můžete definovat samostatné materializované zobrazení nebo streamovací tabulku pomocí stejných dekorátorů, které používáte v kanálu Lakeflow. Každá zdobená funkce definuje jednu tabulku. Když buňku spustíte, Azure Databricks vytvoří tabulku a spustí bezserverové kanálové zpracování k jejímu naplnění daty. Buňka se po dokončení aktualizace znovu zobrazí.
Warning
Dekorátory kanálů vyžadují verzi prostředí serverless 5 nebo vyšší.
Definujte materializovaný pohled
Použijte @dp.materialized_view na funkci, která vrací dávkový DataFrame. Následující příklad vytváří materializovaný pohled daily_booking_revenue z tabulky bookings v ukázkové datové sadě 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"))
)
Pro definici tabulky ze streamovaného čtení použijte @dp.table místo toho.
Definujte streamovací tabulku
Použijte @dp.table na funkci, která vrací streamovaný DataFrame. Následující příklad vytvoří streamovací tabulku bookings_raw na základě streamovaného čtení ze stejné tabulky 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")
Pokud funkce vrátí dávkový DataFrame, @dp.table vytvoří místo toho materializovaný pohled. Jedinou výjimkou je replace_where, jehož výsledkem je vždy streamovací tabulka. Následující příklad udržuje denní příjmy z check-inů k 1. červenci 2025 nebo po něm aktuální, aniž by přepočítával dřívější data:
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"))
)
Každý běh smaže řádky odpovídající predikátu a znovu spočítá pouze tento rozsah. Viz Dávkové zpracování pomocí toků REPLACE WHERE.
Obnovte stůl
Chcete-li aktualizovat tabulku, kterou jste definovali pomocí dekorátoru, znovu spusťte kód, který ji definuje, například opakovaným spuštěním buňky notebooku, spuštěním celého notebooku nebo spuštěním notebooku jako úlohy. Každý běh vytvoří tabulku, pokud neexistuje, a obnoví ji, pokud existuje.
Pro opakované zpracování všech dat dostupných ve zdroji zadejte full_refresh=True u kteréhokoli z dekorátorů:
@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")
Nemůžete použít příkaz REFRESH pro tabulku definovanou pomocí dekorátoru ani plánovat aktualizace pomocí SCHEDULE nebo TRIGGER ON UPDATE. Pro obnovení rozvrhu definujte tabulku v SQL nebo plánujte zápisník jako úlohu. Podívejte se na Úlohy Lakeflow.
Konfigurujte tabulku
Dekorátoři přijímají stejné běžné parametry datové sady, jaké přijímají uvnitř pipeline, včetně comment, table_properties, partition_cols, cluster_by, schema, a :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"))
)
Pro seznam parametrů viz materialized_view a tabulku.
private=True není podporováno, protože soukromou tabulku mohou číst pouze jiné datové sady ve stejném pipeline.
Nepodporovaná rozhraní API
Samostatná tabulka je jedna datová sada s jedním tokem, takže API, která popisují vztahy mezi datovými sadami, nejsou dostupná. Následující vyvolají chybu mimo pipeline:
-
@dp.temporary_viewadp.create_streaming_table -
@dp.append_flowa další dodatečné toky -
dp.create_auto_cdc_flowadp.create_auto_cdc_from_snapshot_flow -
@dp.replace_flowa parametrreplace_using, které definují postupy REPLACE USING. Viz částečné nahrazení snapshotu pomocí toků REPLACE USING. dp.create_sink- Očekávání, například
@dp.expecta@dp.expect_or_fail
Chcete-li je použít, vytvořte místo toho kanál Lakeflow. Viz Vývoj kódu kanálu pomocíPythonu .
Odesílejte SQL příkazy pomocí spark.sql()
V poznámkovém bloku Python předejte stejné příkazy, které byste spustili ze služby Databricks SQL Warehouse do spark.sql(). Syntaxe pro samostatný materializovaný pohled a streamovací tabulku je stejná; liší se pouze způsob podání příkazu. Stejně jako u skladu spouští každý CREATE příkaz REFRESH kanál bez serveru, který operaci zpracuje.
Relace spark je ve výchozím nastavení dostupná v poznámkových blocích Azure Databricks, takže není potřeba žádný import.
Vytvoření materializovaného zobrazení
Následující příklad vytvoří materializované zobrazení mv1 ze základní tabulky 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
""")
Úplné CREATE MATERIALIZED VIEW podrobnosti, jako jsou plánované a aktivované aktualizace, najdete v tématu Vytvoření materializovaného zobrazení.
Vytvořte streamovací tabulku
Následující příklad vytvoří streamovací tabulku sales z raw_data tabulky:
spark.sql("""
CREATE OR REFRESH STREAMING TABLE sales
AS SELECT product, price FROM STREAM raw_data
""")
Podrobné CREATE STREAMING TABLE informace, včetně načítání souborů pomocí Auto Loaderu a plánování, viz Použití samostatných streamovacích tabulek.
Aktualizovat materializované zobrazení nebo streamovací tabulku
REFRESH Pomocí příkazu aktualizujte samostatnou tabulku nejnovějšími daty ze svého zdroje:
spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")
Na bezserverových obecných výpočetních prostředcích jsou aktualizace synchronní. Asynchronní aktualizace ( ASYNC klíčové slovo) se nepodporují. Viz Obecné výpočetní prostředky bez serveru.
Parametrizujte příkazy
Pokud chcete předat hodnoty z vašeho Python kódu do příkazu místo jejich pevně zakódování, použijte v SQL pojmenované značky parametrů a zadejte jejich hodnoty prostřednictvím argumentu argsspark.sql(). Použijte přímo značku, například :min_sales, pro literálové hodnoty. Zabalte značku IDENTIFIER() pouze v případě, že je parametr názvem objektu, jako je tabulka, zobrazení nebo schéma, protože identifikátory nelze nahradit hodnotami prostého řetězce.
Následující příklad parametrizuje materializovaný název zobrazení i hodnotu filtru:
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,
})
Další informace najdete v tématu Značky parametrů a IDENTIFIER klauzule.
Spusťte další příkazy
Z poznámkového bloku Pythonu můžete spustit libovolný samostatný příkaz pro materializované zobrazení nebo streamovací tabulku jeho předáním do spark.sql(), včetně příkazů k naplánování obnovení, úpravě tabulky nebo odstranění tabulky. Informace o použití materializovaných zobrazení a streamovaných tabulek, včetně syntaxe SQL, najdete v tématu Použití samostatných materializovaných zobrazení a použití samostatných streamovacích tabulek.
Limitations
Samostatná materializovaná zobrazení a streamované tabulky vytvořené na bezserverových obecných výpočetních prostředcích mají další omezení, jako je například žádná podpora asynchronních aktualizací a nepřisouzení nákladů na tabulku. Úplný seznam najdete v tématu Obecné výpočetní prostředí bez serveru.
Protože tyto kanály běží na bezserverových obecných výpočetních prostředcích, nikoli v datovém skladu SQL, nedědí vlastní štítky z nadřazeného datového skladu. Šíření tagů skladu do system.billing.usage platí pouze pro materializovaná zobrazení a streamovací tabulky, jejichž příkazy se spouštějí z SQL skladu. Viz Přiřazení nákladů k datovému skladu SQL pomocí vlastních značek.
Tabulky definované pomocí dekorátorů potrubí mají následující další omezení:
- Nemůžete je obnovit pomocí příkazu
REFRESHani naplánovat obnovení pomocíSCHEDULEneboTRIGGER ON UPDATE. Viz Obnovit tabulku. - Nejsou podporovány očekávání, další toky, toky pro zachycování změn dat (CDC), záchyty a dočasné pohledy. Viz Nepodporovaná API.
-
private=Truese nepodporuje.