Important
此功能在 Beta 版中。
替换使用流程保持目标表与流源同步:它替换所有匹配指定键列的行,保持其他数据不变。
列 SEQUENCE BY 负责排序更新,确保即使更新顺序混乱,结果依然正确。 对于每个键,最高的序列获胜,较低的序列行永远不会覆盖目标中已有的更高序列。 使用相同键和相同序列的行是附加而非替换。
替代使用是如何运作的
设想有一个事件表,存储了两个区域的点击和转化事件,按 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 流可在 Databricks Runtime 18.2 及以上版本的经典计算或无服务器计算上运行。 Databricks推荐Unity Catalog。
- 源必须是流源。 REPLACE USING 拒绝接受非流式来源。
- 你必须至少指定一个键列,并且恰好指定一个
SEQUENCE BY列。
何时使用 REPLACE USING 流程
湖流管道提供三种流量覆盖现有行。 根据源数据的形式以及它如何标识要替换的行进行选择:
- 当源数据是一系列按列键控的部分快照时,请使用 REPLACE USING。 替换使用只覆盖与入数据匹配的数据,其他数据保持不动。 它不需要主密钥。
- 在以下情况下使用 AUTO CDC:如果源是具有明确 插入、更新 和 删除 操作的变更数据捕获(CDC)馈送,或者需要 缓慢变化维度(SCD)2 型 历史记录。 AUTO CDC还需要真正的主密钥。 请参阅 AUTO CDC API:使用管道简化变更数据捕获。
- 在以下情况下使用 REPLACE WHERE:当源数据是快照,并且你希望重新计算并覆盖目标表中由谓词选定的一段范围(例如过去 7 天的数据)时,可将其作为批处理操作执行。 它不需要主密钥。 请参阅 使用 REPLACE WHERE 流的批处理。
创建一个替代使用流程
使用 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
替代使用流程支持预期。
warn和fail的行为与其他流程中的表现一致:warn会保留违规行并记录违规情况,而fail会中止更新。 请参阅通过管道预期管理数据质量。
drop 期望会将违规的行视为源从未生成过该行。 被删除的行不会替换、删除或修改目标表中的匹配键:
- 丢弃发生在去重之前,因此流程会保留密钥的最新有效版本。
- 如果某个键对应的所有传入行都被丢弃,则该键现有的行将保持不变。
- 由于被丢弃的行不会设置序列号下限,因此后续的有效更新即使其序列号低于该被丢弃行的序列号,仍然会生效。
局限性
替代使用流程有以下限制:
- 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")