在獨立管線中使用 Python

你可以用 Python 從筆記本建立並刷新獨立的實體化檢視表和串流資料表。 這讓你可以將獨立管線與其他以 Python 為基礎的筆記本工作流程一併管理。

作法有二:

  • 使用 pyspark.pipelines 裝飾器、@dp.materialized_view 和 @dp.table 來定義表格。 當邏輯更容易用 DataFrame 程式碼表達時,就用這個方法。 請參閱 使用管線裝飾器定義資料表。
  • 將 Databricks SQL 倉儲執行的相同 SQL 陳述傳遞給 spark.sql(),即可提交這些陳述。 這會讓你完整地獲得獨立的實體化檢視和串流資料表的 SQL 表面,包括 REFRESH 語句和刷新排程。 請參閱 使用 spark.sql() 提交 SQL 陳述式。

Python Source 用於獨立管線需要一本筆記本連接到無伺服器的一般運算系統。 你不能用 Python 從 Databricks SQL 倉庫建立或刷新獨立管線,因為倉庫是執行 SQL 語句,不是 Python 筆記本。 若想改用 SQL 倉庫,請參見 「使用獨立實體化檢視 」及 「使用獨立串流資料表」。

Important

在無伺服器通用運算平台上,從筆記本建立並刷新獨立的實體化視圖與串流資料表,目前仍處於 測試階段 ,且在特定區域提供。 詳見 筆記本。

要求

要用 Python 建立並刷新獨立管線,你需要一個筆記本連接到 Databricks 18.1 以上的無伺服器通用運算系統。 完整需求清單,包括區域可用性與權限,請參閱 筆記本。

定義帶有管線裝飾器的表格

你可以用 Lakeflow 管線中使用的裝飾工具,定義獨立的實體化視圖或串流表。 每個裝飾過的函式定義一個表格。 當你執行儲存格時,Azure Databricks 會建立資料表,並運行無伺服器的管線來填充它。 更新完成後,儲存格會再次顯示。

Warning

管線裝飾器需要 無伺服器環境版本 5 或以上。

定義具體化的觀點

在會傳回批次 DataFrame 的函式上使用 @dp.materialized_view。 以下範例是從 Wanderbricks 範例資料集中的bookings表格中建立物質化的視圖daily_booking_revenue:

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

若要從串流讀取定義資料表,請使用 @dp.table 代替。

定義串流表

在會傳回串流 DataFrame 的函式上使用 @dp.table。 以下範例是從同一bookings資料表的串流讀取建立串流表bookings_raw:

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

如果函式回傳的是批次資料框,則 @dp.table 會建立一個物質化的視圖。 唯一的例外是 replace_where,它總是會產生串流表。 以下範例會讓 2025 年 7 月 1 日當天及之後的簽到每日收入保持最新狀態,而不必重新計算更早的日期:

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

每次執行都會刪除與謂詞相符的列,並只重新計算該範圍。 請參見 使用 REPLACE WHERE 流程的批次處理。

重新整理表格

要重新整理你用裝飾程式定義的資料表,請重新執行定義該資料表的程式碼,例如重跑筆記本儲存格、整個筆記本,或是將筆記本當作工作執行。 每次執行如果不存在資料表,會建立該資料表;如果存在,則重新整理。

若要重新處理來源中所有可用的資料,請在任一裝飾器上傳遞 full_refresh=True:

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

你不能對使用裝飾器定義的資料表使用 REFRESH 陳述式,也不能使用 SCHEDULE 或 TRIGGER ON UPDATE 來排程重新整理。 若要刷新排程,請在 SQL 中定義表格,或將筆記本排程為工作。 請參閱 Lakeflow 職位。

設定表格

裝飾器接受與其在管線內使用時相同的通用資料集參數,包括 comment、table_properties、partition_cols、cluster_by、schema 和 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"))
  )

關於參數列表,請參見 materialized_view 與 表格。

private=True 不支援,因為私有資料表只能被同一管線中的其他資料集讀取。

不支援的 API

獨立資料表是單一資料集,只有單一流程,因此無法提供描述資料集間關係的 API。 以下會產生流水線外的錯誤:

  • @dp.temporary_view 與 dp.create_streaming_table
  • @dp.append_flow 以及其他額外的流量
  • dp.create_auto_cdc_flow 與 dp.create_auto_cdc_from_snapshot_flow
  • @dp.replace_flow 和 replace_using 參數,用於定義 REPLACE USING 流程。 請參見使用 REPLACE USING 流程進行部分快照取代。
  • dp.create_sink
  • 期望值,例如 @dp.expect 和 @dp.expect_or_fail

若要使用這些功能,請改為撰寫 Lakeflow 管線。 請參閱 使用 Python 開發管線程式代碼。

使用 spark.sql() 提交 SQL 陳述式

在 Python 筆記本中,將你從 Databricks SQL 倉庫執行的相同語句傳給 spark.sql()。 獨立的實體化檢視與串流資料表語法相同;只有提交聲明的方式不同。 與倉庫類似,每個 CREATE 或 REFRESH 語句都會運行無伺服器的管線來處理操作。

spark該會話預設可在 Azure Databricks 筆記本中使用,因此不需要匯入。

建立實體化視圖

下列範例會從基底資料表 mv1 中建立具體化檢視 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
""")

如需完整 CREATE MATERIALIZED VIEW 細節,例如排程與觸發刷新,請參閱 建立實體化視圖。

建立串流表

以下範例從 sales 資料表建立串流資料表 raw_data:

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

完整 CREATE STREAMING TABLE 細節,包括用自動載入器載入檔案及排程,請參見 使用獨立串流表。

重新整理實體化的視圖或串流表

使用 REFRESH 陳述式,以來源端的最新資料更新獨立資料表:

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

在無伺服器一般運算中,刷新是同步的。 不支援非同步重新整理(ASYNC 關鍵字)。 參見 無伺服器一般計算。

參數化陳述

要將 Python 程式碼的值傳入語句,而非硬編碼,可以在 SQL 中使用命名參數標記,並透過 args 的spark.sql()參數提供其值。 直接使用如 :min_sales 這樣的標記來表示字面值。 僅在參數是物件名稱時,才以 IDENTIFIER() 包住標記,例如資料表、檢視表或綱要,因為識別碼不能以一般字串值替代。

以下範例同時參數化了實體化的視圖名稱與濾波器值:

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,
})

欲了解更多資訊,請參閱 參數標記 與 IDENTIFIER 子句。

執行其他陳述式

你可以在 Python 筆記本中執行任何獨立的實體化檢視或串流資料表陳述式,只要將其傳遞給 spark.sql(),包括用於排定重新整理、修改資料表或卸除資料表的陳述式。 要了解如何使用實體化檢視與串流資料表,包括 SQL 語法,請參閱 使用獨立實體化檢視 與 使用獨立串流資料表。

Limitations

在無伺服器一般運算上建立的獨立實體化視圖與串流資料表還有額外限制,例如不支援非同步刷新,且無法按資料表歸因成本。 完整列表請參見 無伺服器一般計算。

由於這些管線是在無伺服器一般運算上執行,而非在 SQL 倉儲上執行,因此不會繼承其所屬 SQL 倉儲的自訂標籤。 將倉庫標籤傳播至 system.billing.usage 僅適用於其陳述式從 SQL 倉庫執行的實體化檢視和串流資料表。 請參考 SQL 倉庫的屬性成本與自訂標籤。

使用管道裝飾器定義的表格另有下列限制:

  • 你無法使用 REFRESH 陳述式重新整理它們,也無法使用 SCHEDULE 或 TRIGGER ON UPDATE 排程重新整理。 參見 「重新整理表格」。
  • 不支援預期、額外流程、變更資料擷取(CDC)流程、匯項及臨時檢視。 請參見 不支援的 API。
  • 不支援 private=True。

其他資源