append_flow

@dp.append_flow裝飾器會為您的管線資料表建立附加流程或回填。 函式必須傳回 Apache Spark 串流 DataFrame。 請參閱 「隨湖流管線流量逐步載入與處理資料」。

追加流程可以以串流資料表或接收端為目標。

語法

from pyspark import pipelines as dp

dp.create_streaming_table("<target-table-name>") # Required only if the target table doesn't exist.

@dp.append_flow(
  target = "<target-table-name>",
  name = "<flow-name>", # optional, defaults to function name
  once = False, # optional
  spark_conf = {"<key>" : "<value", "<key" : "<value>"}, # optional
  comment = "<comment>", # optional
  import_checkpoint = "<checkpoint-path>") # optional
def <function-name>():
  return (<streaming-query>) #

參數

參數 類型 Description
函式 function 必須的。 用來從使用者定義的查詢中返回 Apache Spark 串流 DataFrame 的函數。
target str 必須的。 作為附加流程目標的資料表或接收器的名稱。
name str 流程名稱。 如果未提供,則預設為函式名稱。
once bool 或者,將流程定義為一次性流程,例如回填。 使用 once=True 將以兩種方式改變流程:
  • 傳回值。 streaming-query。 在這種情況下,必須是批次 DataFrame,而非串流 DataFrame。
  • 預設情況下,流程會執行一次。 如果管道已更新為完整刷新,則ONCE流程會再度執行以重新創建資料。
comment str 對流程的描述。
spark_conf dict 用於執行此查詢的 Spark 組態清單
import_checkpoint str 將導向現有的結構化串流檢查點匯入流程,使遷移的串流從上次提交的偏移量恢復,而非重新處理來源。 匯入檢查點目前還在 測試階段。 請參見 遷移結構化串流檢查點

範例

from pyspark import pipelines as dp

# Create a sink for an external Delta table
dp.create_sink("my_sink", "delta", {"path": "/tmp/delta_sink"})

# Add an append flow to an external Delta table
@dp.append_flow(name = "flow", target = "my_sink")
def flowFunc():
  return <streaming-query>

# Add a backfill
@dp.append_flow(name = "backfill", target = "my_sink", once = True)
def backfillFlowFunc():
    return (
      spark.read
      .format("json")
      .load("/path/to/backfill/")
    )

# Create a Kafka sink
dp.create_sink(
  "my_kafka_sink",
  "kafka",
  {
    "kafka.bootstrap.servers": "host:port",
    "topic": "my_topic"
  }
)

# Add an append flow to a Kafka sink
@dp.append_flow(name = "flow", target = "my_kafka_sink")
def myFlow():
  return read_stream("xxx").select(F.to_json(F.struct("*")).alias("value"))

遷移結構化串流檢查點

Important

匯入檢查點目前還在 測試階段

用於 import_checkpoint 將現有的結構化串流工作負載遷移到管線,無需重新處理原始碼。 設定成 checkpointLocation 你使用的結構化串流查詢,可以是雲端儲存、Unity Catalog 卷或 DBFS 路徑。 在第一次管線更新時,流程會將該檢查點複製到管線的管理儲存中。 接著流程會從最後一個提交的偏移量開始,且其狀態(如聚合、重複去重鍵和浮水印)保持完整。 後續的管線更新會使用流程的複製檢查點;原始檢查點不會被修改。

流程必須針對由 create_table建立的管理資料表。

在執行管線之前,先停止原本的結構化串流查詢。 原始的結構化串流查詢可以在匯入後重複使用,但你需要管理其檢查點狀態,並確保管線和查詢不會同時寫入同一個資料表,否則可能會產生重複資料。

將結構化串流查詢重新建立為一個流水線流程,寫入新資料表並匯入其檢查點:

from pyspark import pipelines as dp

# Create a new managed table for the pipeline
dp.create_table("target_table")

# Continue from the imported checkpoint instead of reprocessing the source.
@dp.append_flow(
  target = "target_table",
  import_checkpoint = "/Volumes/my_catalog/my_schema/checkpoints/my_stream",
)
def migrate():
  # The same source your original query read from.
  return spark.readStream.table("source_table")

檢查點僅在第一次管線更新時匯入一次;後續更新忽略 import_checkpoint完整刷新不會重新匯入檢查點;它從新的空白檢查點開始,並重新處理原始碼。 若要匯入不同的檢查點,請使用先前未用於目標資料表的流程名稱;重複使用現有流程名稱則會跳過匯入。

Limitations

  • 不支援將檢查點匯入已存在的資料表(例如原始的結構化串流查詢目標)。 針對管線建立的新資料表, 或是匯入。
  • import_checkpoint 僅支援 append_flow