Używanie Pythona z samodzielnymi potokami

W notesie możesz za pomocą języka Python tworzyć i odświeżać samodzielne zmaterializowane widoki oraz tabele strumieniowe. Dzięki temu można zarządzać samodzielnymi potokami obok innych przepływów pracy w notatnikach opartych na języku Python.

Istnieją dwa sposoby, aby to zrobić:

  • Zdefiniuj tabelę z dekoratorami pyspark.pipelines , @dp.materialized_view a @dp.table. Używaj tego, gdy logika łatwiej jest wyrazić jako kod DataFrame. Zobacz Definiowanie tabel za pomocą dekoratorów potoków.
  • Przesyłaj te same instrukcje SQL co magazyn SQL Databricks, przekazując je do spark.sql(). Zapewnia to pełną obsługę SQL dla samodzielnych widoków materializowanych i tabel strumieniowych, w tym instrukcje REFRESH oraz harmonogramy odświeżania. Zobacz przesyłanie instrukcji SQL za pomocą spark.sql().

Kod źródłowy w języku Python dla samodzielnych potoków wymaga notatnika podłączonego do bezserwerowego środowiska obliczeniowego ogólnego przeznaczenia. Nie można użyć Python do tworzenia lub odświeżania autonomicznych potoków z usługi Databricks SQL Warehouse, ponieważ magazyn uruchamia instrukcje SQL, a nie notesy Python. Aby zamiast tego użyć magazynu SQL, zobacz Korzystanie z samodzielnych widoków zmaterializowanych i Korzystanie z samodzielnych tabel strumieniowych.

Important

Tworzenie i odświeżanie samodzielnych widoków zmaterializowanych oraz tabel strumieniowych z poziomu notesu w bezserwerowym środowisku ogólnego przeznaczenia jest dostępne w ramach wersji beta i tylko w wybranych regionach. Zobacz Laptopy.

Requirements

Aby tworzyć i odświeżać potoki autonomiczne przy użyciu Python, potrzebny jest notes dołączony do bezserwerowych ogólnych obliczeń w środowisku Databricks Runtime 18.1 lub nowszym. Aby uzyskać pełną listę wymagań, w tym dostępność regionalną i uprawnienia, zobacz Notebooks.

Definiuj tabele za pomocą dekoratorów potoków

Możesz zdefiniować samodzielny widok materializowany lub tabelę streamingową z tymi samymi dekoratorami, których używasz w pipeline Lakeflow. Każda funkcja dekorowana definiuje jedną tabelę. Gdy uruchamiasz komórkę, Azure Databricks tworzy tabelę i uruchamia pipeline bez serwera, aby ją wypełnić. Komórka wraca po zakończeniu aktualizacji.

Warning

Dekoratory potoków wymagają wersji środowiska bezserwerowego 5 lub nowszej.

Zdefiniuj widok materializowany

Użyj @dp.materialized_view w funkcji, która zwraca wsadową ramkę danych. Poniższy przykład tworzy zmaterializowany widok daily_booking_revenue z tabeli bookings w zbiorze danych próbnych 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"))
  )

Aby zdefiniować tabelę na podstawie odczytu strumieniowego, zamiast tego użyj @dp.table.

Zdefiniuj tabelę strumieniową

Użyj @dp.table na funkcji, która zwraca streamingowy DataFrame. Poniższy przykład tworzy tabelę strumieniową bookings_raw na podstawie odczytu strumieniowego tej samej tabeli 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")

Jeśli funkcja zwraca wsadowy DataFrame, @dp.table tworzy zamiast tego widok materializowany. Jedynym wyjątkiem jest replace_where, co zawsze skutkuje tabelą strumieniową. Poniższy przykład aktualizuje dzienne przychody dla zameldowań z dnia 1 lipca 2025 r. i późniejszych, bez ponownego przeliczania wcześniejszych dat:

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żde uruchomienie usuwa wiersze spełniające predykat i ponownie oblicza tylko ten zakres. Zobacz Przetwarzanie wsadowe z użyciem przepływów REPLACEWHERE.

Odśwież stół

Aby odświeżyć zdefiniowaną tabelę za pomocą dekoratora, uruchom kod, który ją definiuje, na przykład ponownie uruchamiając komórkę notatnika, cały zeszyt lub uruchamiając zeszyt jako zadanie. Każde uruchomienie tworzy tabelę, jeśli nie istnieje, i odświeża ją, jeśli istnieje.

Aby ponownie przetworzyć wszystkie dane dostępne w źródle, przekaż full_refresh=True do jednego z dwóch dekoratorów:

@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")

Nie możesz użyć instrukcji REFRESH dla tabeli zdefiniowanej za pomocą dekoratora ani planować odświeżania za pomocą SCHEDULE lub TRIGGER ON UPDATE. Aby odświeżyć harmonogram, zdefiniuj tabelę w SQL lub zaplanuj notatnik jako zadanie. Zobacz Zadania lakeflow.

Konfiguruj tabelę

Dekoratorzy akceptują te same wspólne parametry zbioru danych, które akceptują w potoku, w tym comment, table_properties, partition_cols, cluster_by, schema, oraz 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"))
  )

Listę parametrów można znaleźć w materialized_view i tabeli.

private=True nie jest obsługiwana, ponieważ prywatna tabela może być odczytana tylko przez inne zbiory danych w tym samym potoku.

Nieobsługiwane interfejsy API

Samodzielna tabela to pojedynczy zbiór danych z jednym przepływem, więc API opisujące relacje między zbiorami danych nie są dostępne. Następujące wywołują błąd poza potokiem:

  • @dp.temporary_view i dp.create_streaming_table
  • @dp.append_flow oraz inne dodatkowe przepływy
  • dp.create_auto_cdc_flow i dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow oraz parametr replace_using, które definiują przepływy typu REPLACE USING. Zobacz Częściowe zastępowanie migawki za pomocą przepływów REPLACE USING.
  • dp.create_sink
  • Oczekiwania, takie jak @dp.expect i @dp.expect_or_fail

Aby z nich korzystać, należy zamiast tego tworzyć pipeline Lakeflow. Zobacz Rozwijaj kod potoku za pomocą Pythona.

Przesyłaj instrukcje SQL za pomocą spark.sql()

W notesie Python przekaż do spark.sql() te same instrukcje, których używasz w magazynie Databricks SQL. Składnia samodzielnego widoku zmaterializowanego i tabeli strumieniowej jest identyczna; różni się tylko sposób przesyłania polecenia. Podobnie jak w magazynie danych, każda instrukcja CREATE lub REFRESH uruchamia bezserwerowy potok przetwarzania, aby przetworzyć operację.

Sesja spark jest domyślnie dostępna w notesach Azure Databricks, więc importowanie nie jest wymagane.

Tworzenie zmaterializowanego widoku

Poniższy przykład tworzy zmaterializowany widok mv1 z tabeli base_table1podstawowej :

spark.sql("""
  CREATE OR REPLACE MATERIALIZED VIEW mv1
  AS SELECT
    date,
    sum(sales) AS sum_of_sales
  FROM base_table1
  GROUP BY date
""")

Aby uzyskać szczegółowe CREATE MATERIALIZED VIEW informacje, takie jak zaplanowane i wyzwalane odświeżanie, zobacz Tworzenie zmaterializowanego widoku.

Tworzenie tabeli przesyłania strumieniowego

Poniższy przykład tworzy tabelę strumieniową sales z tabeli raw_data:

spark.sql("""
  CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT product, price FROM STREAM raw_data
""")

Aby uzyskać pełne CREATE STREAMING TABLE informacje, w tym informacje o ładowaniu plików za pomocą narzędzia Auto Loader i harmonogramowaniu, zobacz Używanie samodzielnych tabel strumieniowych.

Odśwież widok zmaterializowany lub tabelę strumieniową

Użyj instrukcji REFRESH, aby zaktualizować samodzielną tabelę najnowszymi danymi z jej źródła:

spark.sql("REFRESH MATERIALIZED VIEW mv1")
spark.sql("REFRESH STREAMING TABLE sales")

W przypadku bezserwerowych ogólnych obliczeń odświeżanie jest synchroniczne. Odświeżanie asynchroniczne ( ASYNC słowo kluczowe) nie jest obsługiwane. Zobacz Ogólne obliczenia bezserwerowe.

Instrukcje parametryzacji

Aby przekazać wartości z kodu Python do zapytania zamiast kodować je na stałe, użyj nazwanych znaczników parametrów w SQL i przekaż ich wartości za pośrednictwem argumentu args w spark.sql(). Użyj znacznika, na przykład :min_sales, bezpośrednio dla wartości literałowych. Umieszczaj znacznik w IDENTIFIER() tylko wtedy, gdy parametr jest nazwą obiektu, na przykład tabeli, widoku lub schematu, ponieważ identyfikatorów nie można podstawiać jako zwykłych wartości tekstowych.

Poniższy przykład parametrizuje zarówno zmaterializowaną nazwę widoku, jak i wartość 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,
})

Aby uzyskać więcej informacji, zobacz Znaczniki parametrów i IDENTIFIER klauzula.

Uruchom inne instrukcje

Możesz uruchomić dowolną samodzielną instrukcję dotyczącą widoku materializowanego lub tabeli strumieniowej w notatniku Python, przekazując ją do spark.sql(), w tym instrukcje planowania odświeżania, modyfikowania tabeli i usuwania tabeli. Aby zrozumieć, jak używać zmaterializowanych widoków i tabel przesyłania strumieniowego, w tym składni JĘZYKA SQL, zobacz Używanie autonomicznych zmaterializowanych widoków i Używanie autonomicznych tabel przesyłania strumieniowego.

Ograniczenia

Samodzielne zmaterializowane widoki i tabele strumieniowe utworzone w bezserwerowym środowisku obliczeniowym ogólnego przeznaczenia mają dodatkowe ograniczenia, takie jak brak obsługi asynchronicznego odświeżania oraz brak przypisywania kosztów do poszczególnych tabel. Aby uzyskać pełną listę, zobacz Ogólne obliczenia bezserwerowe.

Ponieważ te potoki działają w bezserwerowym środowisku obliczeniowym ogólnego przeznaczenia, a nie w hurtowni SQL, nie dziedziczą niestandardowych tagów z nadrzędnej hurtowni. Propagacja tagów magazynu do system.billing.usage dotyczy tylko widoków zmaterializowanych i tabel strumieniowych, których instrukcje są uruchamiane z magazynu SQL. Zobacz Przypisywanie kosztów do magazynu SQL za pomocą tagów niestandardowych.

Tabele zdefiniowane za pomocą dekoratorów potoków mają następujące dodatkowe ograniczenia:

  • Nie można ich odświeżyć za pomocą instrukcji REFRESH ani planować odświeżenia za pomocą SCHEDULE lub TRIGGER ON UPDATE. Zobacz Odśwież tabelę.
  • Nie są obsługiwane oczekiwania, dodatkowe przepływy, przepływy przechwytywania danych zmian (CDC), pochłaniacze oraz widoki tymczasowe. Zobacz Nieobsługiwane API.
  • private=True nie jest obsługiwana.

Dodatkowe zasoby