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 流支持预期。
warn 和 fail 的行为与在其他流中一致:warn 会保留违规行并记录违规情况,而 fail 会停止更新。 请参阅通过管道预期管理数据质量。
drop 预期会将违规的行视为源从未生成过该行。 被删除的行不会替换、删除或修改目标表中的匹配键:
- 删除操作发生在重复数据删除之前,因此流会保留键的最新有效版本。
- 如果某个键的所有传入行都被删除,则该键的现有行保持不变。
- 由于被丢弃的行不会设置序列号下限,因此后续的有效更新即使其序列号低于该被丢弃行的序列号,仍然会生效。
局限性
REPLACE USING 流具有以下限制:
- REPLACE USING 支持每个目标表一个流。 不支持在同一目标上将 REPLACE USING 与另一种流类型结合使用。
- 目标表必须在管道内创建。
- 源必须是流源。
- 你必须指定至少一个关键列和一个
SEQUENCE BY列。 键列不能重复,每个键列的类型必须是可排序的。 原子类型,如整数、字符串和日期,可以是键,而MAP和VARIANT不能。 - 对于独立流式表,请参阅使用 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")