append_flow

@dp.append_flow修饰器为管道表创建追加流或回填。 该函数必须返回 Apache Spark 流式处理数据帧。 请参阅 使用 Lakeflow 管道流以增量方式加载和处理数据

追加流可以面向流式处理表或接收器。

Syntax

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 流式处理数据帧的函数。
target str 必填。 作为追加流的目标的表或接收器的名称。
name str 流名称。 如果未提供,则默认为函数名称。
once bool (可选)将流定义为一次性流,例如回填。 通过两种方式使用 once=True 更改流:
  • 返回值。 streaming-query。 在这种情况下,必须是批处理数据帧,而不是流式处理数据帧。
  • 默认情况运行一次。 如果管道通过完全刷新进行更新,则 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完全刷新不会重新导入检查点;它从一个新的空检查点开始,重新处理源代码。 要导入不同的检查点,可以使用之前未用于目标表的流名;重复使用现有流名则跳过导入过程。

局限性

  • 不支持将检查点导入已存在的表(例如原始的结构化流查询目标)。 针对流水线创建的新表,或者一个汇。
  • import_checkpoint 仅支持 append_flow