使用 REPLACE USING 流进行部分快照替换

Important

此功能在 Beta 版中。

替换使用流程保持目标表与流源同步:它替换所有匹配指定键列的行,保持其他数据不变。

SEQUENCE BY 负责排序更新,确保即使更新顺序混乱,结果依然正确。 对于每个键,最高的序列获胜,较低的序列行永远不会覆盖目标中已有的更高序列。 使用相同键和相同序列的行是附加而非替换。

REPLACE USING 的工作原理

设想有一个事件表,存储了两个区域的点击和转化事件,按 seq 排序:

region_id device_type 事件类型 seq
1 iOS click 1
1 Android 转换 1
2 iOS click 1
2 桌面 click 1

REPLACE USING (region_id) SEQUENCE BY seq 流接收针对区域 1 和区域 3 的这些更新。 第二区没有更新:

region_id device_type event_type seq
1 iOS click 2
1 Android 转换 2
1 桌面 click 2
3 iOS click 1
3 桌面 click 2

目标变为:

region_id device_type 事件类型 seq 结果
1 iOS click 2 被替换,因为序列2大于序列1
1 Android 转换 2 被替换,因为序列2大于序列1
1 桌面 click 2 被替换,因为序列2大于序列1
2 iOS click 1 未更改,因为此次更新中不存在该密钥
2 桌面 click 1 未更改,因为此更新中不存在该键
3 桌面 click 2 已添加。 区域 3 的序列号为 1 的行未被添加,因为仅应用键的最高序列号。

要求

REPLACE USING 流具有以下要求:

  • REPLACE USING 流在 Databricks Runtime 18.2 及更高版本上运行,支持经典计算或无服务器计算。 Databricks推荐Unity Catalog。
  • 源必须是流源。 REPLACE USING 会拒绝非流式处理源。
  • 必须指定至少一个键列和恰好一个 SEQUENCE BY 列。

何时使用 REPLACE USING 流程

Lakeflow 管道提供了三种覆盖现有行的流。 根据源的形式及其标识要替换行的方式进行选择:

  • 在以下情况下使用 REPLACE USING:源是一系列以列为键的部分快照。 REPLACE USING 仅覆盖在传入数据中具有匹配项的数据,其余数据保持不变。 它不需要主密钥。
  • 在以下情况下使用 AUTO CDC:源是包含明确插入更新删除操作的变更数据捕获 (CDC) 源,或者你需要渐变维度 (SCD) 类型 2 历史记录。 AUTO CDC还需要真正的主密钥。 请参阅 AUTO CDC API:使用管道简化变更数据捕获
  • 在以下情况下使用 REPLACE WHERE:当源数据是快照,并且你希望重新计算并覆盖目标表中由谓词选定的一段范围(例如过去 7 天的数据)时,可将其作为批处理操作执行。 它不需要主密钥。 请参阅使用 REPLACE WHERE 流的批处理

创建 REPLACE USING 流

使用 SQL 或 Python 定义 REPLACE USING 流程。

SQL

FLOW REPLACE USING 子句与 CREATE STREAMING TABLE 内联使用:

CREATE STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

或者,使用长格式 CREATE FLOW 语法:

CREATE STREAMING TABLE payments_current;

CREATE FLOW payments_flow AS
INSERT INTO payments_current BY NAME
REPLACE USING (payment_id) SEQUENCE BY payment_date
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

注释

BY NAME 在 SQL 中是必需的。 它根据列名而非位置进行匹配。

Python

将表和流一起声明为 @dp.table

from pyspark import pipelines as dp

@dp.table(name="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_current():
  return spark.readStream.table("samples.wanderbricks.payments")

或者,使用 @dp.replace_flow 指定现有流式表作为目标:

from pyspark import pipelines as dp

dp.create_streaming_table("payments_current")

@dp.replace_flow(target="payments_current", replace_using=["payment_id"], sequence_by="payment_date")
def payments_flow():
  return spark.readStream.table("samples.wanderbricks.payments")

replace_using 是关键列的列表。 sequence_by 是列名或 Column 表达式,且每当 replace_using 被设置时都是必需的。

测序与错序数据

SEQUENCE BY 列使结果与更新到达的顺序无关。 仅当某一行的序列号大于该键当前已存储的序列号时,才会将该行应用于该键,因此,晚到或重放且早于当前值的行会被忽略。 更新中未包含的键将保持不变。

遵循以下做法,使替换行为可预测:

练习 原因
使用严格按密钥版本增加的序列,比如时间戳、版本号或日志偏移量。 具有相同键和相同序列的两行都会被保留,从而导致该键出现重复行。
使用非空序列。 空序列可能导致行为未定义。

Expectations

REPLACE USING 流支持预期。 warnfail 的行为与在其他流中一致:warn 会保留违规行并记录违规情况,而 fail 会停止更新。 请参阅通过管道预期管理数据质量

drop 预期会将违规的行视为源从未生成过该行。 被删除的行不会替换、删除或修改目标表中的匹配键:

  • 删除操作发生在重复数据删除之前,因此流会保留键的最新有效版本。
  • 如果某个键的所有传入行都被删除,则该键的现有行保持不变。
  • 由于被丢弃的行不会设置序列号下限,因此后续的有效更新即使其序列号低于该被丢弃行的序列号,仍然会生效。

局限性

REPLACE USING 流具有以下限制:

  • REPLACE USING 支持每个目标表一个流。 不支持在同一目标上将 REPLACE USING 与另一种流类型结合使用。
  • 目标表必须在管道内创建。
  • 源必须是流源。
  • 你必须指定至少一个关键列和一个 SEQUENCE BY 列。 键列不能重复,每个键列的类型必须是可排序的。 原子类型,如整数、字符串和日期,可以是键,而 MAPVARIANT 不能。
  • 对于独立流式表,请参阅使用 REPLACE USING 流应用部分快照替换以了解语法差异。

Examples

以下示例取自 samples.wanderbricks.booking_updates,这是一个在每个支持 Unity Catalog 的工作区中可用的预订状态变更示例表。 每次发生变更时,每个预订都会显示一次,因此 booking_id 会随着新的 booking_update_id 重复出现。 参见 Wanderbricks 数据集

示例1:保持每个密钥的最新记录

只保留每个预订的当前状态。 该流以 booking_id 为键,并按 booking_update_id 排序,因此预订的最新更新会替换其早期版本。 当源是包含明确插入、更新和删除操作的更改源时,请改用 AUTO CDC。

SQL

CREATE OR REFRESH STREAMING TABLE bookings_current
FLOW REPLACE USING (booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_current",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
def bookings_current():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

这个例子是按 booking_update_id 排序,而不是按 updated_at 时间戳排序,因为同一预订的多次更新可能具有相同的时间戳。 序列号相同的行将被追加而非替换,这将导致这些预订保留多行记录。

示例 2:使用多个列作为键

当记录由多列组合标识时,将所有列列在 REPLACE USING中列出。 这里每个预订都由 (property_id, booking_id) 标识,因此该流程会保留每个物业下每个预订的当前状态。 如果键列可以为 null,REPLACE USING 会将 null 与 null 匹配,而不是跳过该行。

SQL

CREATE OR REFRESH STREAMING TABLE bookings_by_property
FLOW REPLACE USING (property_id, booking_id) SEQUENCE BY booking_update_id BY NAME
SELECT property_id, booking_id, status, total_amount, booking_update_id
FROM STREAM(samples.wanderbricks.booking_updates);

Python

from pyspark import pipelines as dp

@dp.table(
  name="bookings_by_property",
  replace_using=["property_id", "booking_id"],
  sequence_by="booking_update_id"
)
def bookings_by_property():
  return spark.readStream.table("samples.wanderbricks.booking_updates")

示例 3:使用预期条件丢弃无效记录

添加预期以将错误行排除在目标范围之外。 丢弃的行被视为源从未生成过:它不会替换或删除匹配的密钥,流程会回落到该密钥的最新有效行。 该流会删除没有正 total_amount 的更新。

from pyspark import pipelines as dp

@dp.table(
  name="bookings_validated",
  replace_using=["booking_id"],
  sequence_by="booking_update_id"
)
@dp.expect_or_drop("positive_amount", "total_amount > 0")
def bookings_validated():
  return spark.readStream.table("samples.wanderbricks.booking_updates")