Lakeflow管道中的处理保障

在任何真实的管道中,重试和重新运行都是不可避免的,因此本页将说明 Lakeflow 管道为你提供的处理保证,以及如何使你编写的部分能够安全地重新运行。

Overview

有两个相关属性决定了重新运行管道是否安全:

  • 幂等性 意味着无论你对同一输入运行多少次,流水线都会产生相同的结果。 失败后重跑、重复填充同一日期范围,或手动重新触发作业都不会生成重复行或损坏状态。
  • 处理保证 描述了每条记录对结果影响的次数。 至少一次处理可以保证所有记录都能被处理,但失败和重试可能会让部分记录处理多次,这可能导致重复。 恰好一次处理语义保证每条记录对结果产生的影响,都如同它被精确地处理了一次,即使在重试情况下,也不会出现重复或遗漏。

Lakeflow 管道默认对其管理的部分具有幂等性,并且在它们自己的管理表中只提供一次处理。 重要的是要明白这些保证在哪里不再自动,这样你才能在管道边缘添加合适的保障措施。

工作原理

Lakeflow 管道为其管理的数据流提供恰好一次处理和幂等性,并提供工具来帮助你确保所编写的逻辑也具有幂等性。

托管表的精确一次处理

在托管表中,默认是精确一次的处理。 流式表结合了结构化流检查点和Delta Lake的事务写入:每个微批次都会提交其源偏移量和输出,因此失败后的重试批处理要么完全成功,要么被完全回滚重试,绝不会部分应用两次。 这同样适用于 Auto Loader 文件引入、Kafka、Kinesis 和 Azure 事件中心 读取,以及 AUTO CDC 更新,无需编写任何代码。

如果一个至少一次的源发送同一条记录多次,流水线会将它们作为唯一记录处理,并全部写入你的表。 删除这些重复的记录是你的责任。 参见至少 一次去重复的来源

读取操作的幂等性也来自这些相同的检查点。 自动加载器和流式表检查点保证每个源文件或偏移量被处理一次以进行状态追踪,因此失败后重新处理流水线更新时,会从检查点恢复,而不是重新处理或跳过数据。 你可以通过在 spark.readStream 之上使用流式表,而不是手写的批处理循环,来做到这一点。 请参阅流式处理表

使用 AUTO CDC 代替手写的 MERGE 语句

AUTO CDC INTO 就其 keyssequence_by 而言,本质上是幂等的。 将相同的变更记录应用两次,或乱序应用这些记录,都会产生相同的最终状态,因为流水线使用序列列来判断传入的行是否确实比已存储的行更新:

CREATE FLOW customers_cdc_flow AS AUTO CDC INTO customers_silver
FROM stream(customers_cdc_bronze)
KEYS (customer_id)
SEQUENCE BY sequence_num
STORED AS SCD TYPE 1;

如果你在 AUTO CDC 之外编写自己的 upsert 逻辑(这种情况虽然少见,但在复杂的合并条件下有时是必要的),应基于稳定的业务键,并确保该逻辑可安全地重复执行两次,例如使用以 order_id 为键的 MERGE ... WHEN MATCHED,而不是盲目地执行 INSERT。 更多信息请参见 《自动CDC API:通过管道简化变更数据采集》。

确保你自己的转换是幂等的

为使逻辑在重新运行写入操作时保持幂等性,请遵循以下两项准则:

  • 避免在具象化视图中使用非确定性变换。 由于实体化视图可以完全或增量地重新计算,避免使用输出依赖于运行时间而非输入值的函数。 例如,不要使用 current_timestamp() 来计算写入后就应保持固定的业务值;应从源事件中获取时间戳,或将其作为参数传入,这样重新计算时才能产生完全相同的输出。
  • 设计全面刷新以确保安全。 完整刷新会删除现有表并从头重新计算,只有在每个上游源仍能提供全部历史数据时,这样做才是安全的。 如果上游源端仅保留一个滚动变化窗口,那么对下游 AUTO CDC 表执行全量刷新时,可能会在无提示的情况下丢失历史数据,因此在设计源端和主题的保留策略时应考虑这一点。

在边缘端实现恰好一次

“恰好一次”语义不再是自动保证的地方,恰恰是在超出流水线直接控制范围的边界处,例如写入外部系统时。 当你将数据写入外部系统时,要确保写入操作本身是幂等的,例如在接收端按键执行 upsert(更新插入),否则重试的微批次可能会将同一批数据写入两次。 以下汇入器从执行者写入批次的每个分区,并使用幂等密钥,确保重试批次不会重复写入:

from pyspark import pipelines as dp

@dp.foreach_batch_sink(name="orders_to_external_api")
def write_orders_to_api(batch_df, batch_id):
    def write_partition(rows):
        # Open one client per partition.
        for row in rows:
            # Use an idempotency key (order_id) so a retried batch doesn't double-write.
            upsert_to_external_system(key=row.order_id, payload=row.asDict())

    batch_df.select("order_id", "amount").foreachPartition(write_partition)

关于写入外部系统的更多信息,请参见 “湖流管道中的汇”。

重复数据删除至少一次源

当某个源可能多次产生同一条记录时,请在下游进行去重。 将水印与 dropDuplicatesWithinWatermark 结合使用,后者能够感知水印,并且无需使用无界状态来检测重复项。 基于可唯一标识事件的列进行去重。 当没有单一列单独唯一时,身份可以跨越多个列。 在以下例子中,点击序列号仅在其会话中唯一,因此两列共同标识事件:

from pyspark import pipelines as dp

@dp.table(name="clicks_deduped")
def clicks_deduped():
    return (
        spark.readStream.table("clicks_bronze")
        .withWatermark("click_ts", "5 minutes")
        .dropDuplicatesWithinWatermark(["session_id", "click_seq_num"])
    )

应根据源端关于唯一性的约定来选择这些列,而不是根据样本数据中表面上看起来不同的内容来判断。 对于本来就可能重复的列,如果将其视为唯一标识,就会丢弃真实发生的事件。 用户连续点击同一则广告两次是常见例子:用户删除重码后,广告无声地删除了第二次点击。

AUTO CDC 基于键的 upsert 语义也会自然合并重复项,因此,将至少一次投递的数据路由到以稳定业务键为键的 AUTO CDC 流中,也是收敛到精确一次状态的另一种方式。

局限性

精确一次处理适用于受管理的Delta到Delta流量。 将以下边视为至少一次,并在此添加显式的去重或幂等写逻辑:

  • foreach_batch_sink 以及自定义的外部写入操作。 Spark 保证每个批次至少会被尝试一次,但如果某个批次在部分写入后被重试,可能会导致某些行在外部系统中可见两次。 使外部写入具备幂等性,例如按自然键执行插入或更新操作,或写入一个批次 ID,以便接收方据此去重。
  • Kafka 作为接收端。 Kafka 主题不像 Delta 那样支持事务性的精确一次写入,因此,重试的微批次写入 Kafka 时可能会产生重复消息。 如果下游消费者对重复数据敏感,可以在消费端进行去重,例如根据事件 ID 去重。
  • 自定义的 Python 数据源作为来源。 读取是否恰好一次取决于你的源实现是否正确报告并从偏移量恢复。 如果不跟踪偏移量,就应将其视为至少一次处理语义,并在下游基于事件 ID 使用 AUTO CDC 去重,或依赖 dropDuplicates 的基于键的 upsert 语义。

通常来说,如果你的整个管道都是 Delta 到 Delta(流式表和物化视图通过托管流读取和写入 Delta 表),那么你已经具备恰好一次语义。 一旦你添加了 foreach_batch_sink、非 Delta 接收器或未经验证的自定义源,就要将该特定边缘视为至少一次语义,并在该处添加幂等写入或去重逻辑。

其他资源