Snowflake动态表更快更灵活:自定义增量实战(下)
DataHot 速览
Snowflake 2026年的Dynamic Tables迎来多项平台增强,本文是该系列第三部分,重点展示自定义增量动态表(CUSTOM_INCREMENTAL)与MERGE的实战用法。新功能允许用户自行编写MERGE或INSERT逻辑,而调度、重试、依赖跟踪和事务保障仍由Snowflake负责。文中给出了语法格式,并说明需要显式列清单、REFRESH USING中只能有一条DML,且不支持存储过程与多语句事务。作者提醒,自定义增量后数据一致性责任由用户承担,若MERGE逻辑出错,Snowflake不会自动发现。
为什么值得关注:数据从业者可从中了解Snowflake Dynamic Tables新增的自定义增量更新能力,兼顾声明式管道的便利与手动控制,适用于需要复杂MERGE逻辑的产线数据场景。
本文目录 14 节
译文
AI 逐段翻译一份针对Snowflake 2026年动态表更新的实践指南——所有示例均在真实账户上运行。

免责声明:我是Snowflake的首席技术架构师 拥有超过30年的数据战略、架构和开发经验。此处表达的观点仅代表我个人,并不一定反映我当前、前任或未来雇主的观点。
本文档概述了一套全面的平台增强功能和新引入的能力。鉴于更新内容的广泛性,内容分三个不同部分组织,以便聚焦阅读和高效知识传递。
第9部分——使用MERGE的自定义增量动态表
这是核心能力:自行编写MERGE或INSERT逻辑,而Snowflake仍负责调度、重试、依赖跟踪和事务保证。

第二行是权衡。声明式动态表保证等于其查询返回的结果。使用自定义增量方式,该保证由您承担——如果您的MERGE逻辑有误,表也会出错,而Snowflake不会察觉。
语法
CREATE [ OR REPLACE ] DYNAMIC TABLE <name> (
<col_name> <col_type> [ , … ] - REQUIRED
)
TARGET_LAG = { '<time_spec>' | DOWNSTREAM }
WAREHOUSE = <warehouse_name>
[ REFRESH_MODE = { AUTO | CUSTOM_INCREMENTAL } ]
[ INITIALIZE = ON_SCHEDULE ]
[ BACKFILL FROM <table_name> ]
[ START AT ({ STREAM => '<stream>' | TIMESTAMP => <ts>
| STATEMENT => <query_id> | OFFSET => -<seconds> }) ]
REFRESH USING ( <single_dml_statement> )硬性要求:显式列列表(模式无法从DML推断),每个REFRESH USING恰好一个 DML语句,且不支持多语句事务或存储过程。当存在REFRESH USING时,REFRESH_MODE = AUTO解析为CUSTOM_INCREMENTAL。
SELF用于引用被定义的表——作为写入目标(MERGE INTO SELF)和读取源(USING子查询中的FROM SELF AS cur,以读取当前内容)。在REFRESH USING内部不能使用表自身的对象名,并且对SELF使用CHANGES()会被拒绝。
CHANGES()取代了流语义。Snowflake将时间间隔绑定到刷新边界,因此您不能指定AT、BEFORE或END。
伴随两个元数据列:
工作示例:带删除传播的CDC增强
CREATE OR REPLACE TABLE ci_orders_cdc (
order_id NUMBER,
customer_id NUMBER,
amount NUMBER(12,2),
updated_at TIMESTAMP_NTZ
) CHANGE_TRACKING = TRUE;
CREATE OR REPLACE DYNAMIC TABLE ci_dt_orders_enriched (
order_id NUMBER,
customer_id NUMBER,
customer_nm VARCHAR,
region VARCHAR,
amount NUMBER(12,2),
updated_at TIMESTAMP_NTZ
)
TARGET_LAG = '1 minute'
WAREHOUSE = transform_wh
REFRESH_MODE = CUSTOM_INCREMENTAL
REFRESH USING (
MERGE INTO SELF AS tgt
USING (
SELECT c.order_id, c.customer_id, d.customer_nm, d.region,
c.amount, c.updated_at,
c.METADATA$ACTION AS chg_action,
c.METADATA$ISUPDATE AS chg_isupdate
FROM ci_orders_cdc CHANGES(INFORMATION => DEFAULT) AS c
LEFT OUTER JOIN ci_dim_customers AS d
ON c.customer_id = d.customer_id
QUALIFY ROW_NUMBER() OVER (
PARTITION BY c.order_id
ORDER BY c.updated_at DESC,
IFF(c.METADATA$ACTION = 'INSERT', 0, 1)
) = 1
) AS src
ON tgt.order_id = src.order_id
WHEN MATCHED AND src.chg_action = 'DELETE' THEN DELETE
WHEN MATCHED THEN UPDATE SET
tgt.customer_id = src.customer_id,
tgt.customer_nm = src.customer_nm,
tgt.region = src.region,
tgt.amount = src.amount,
tgt.updated_at = src.updated_at
WHEN NOT MATCHED AND src.chg_action = 'INSERT' THEN INSERT
(order_id, customer_id, customer_nm, region, amount, updated_at)
VALUES (src.order_id, src.customer_id, src.customer_nm,
src.region, src.amount, src.updated_at)
);在一个批次中应用插入、更新和删除:
INSERT INTO ci_orders_cdc VALUES (4,104,500.00,'2026-08-02 09:00:00');
UPDATE ci_orders_cdc SET amount = 999.99,
updated_at = '2026-08-02 10:00:00'
WHERE order_id = 2;
DELETE FROM ci_orders_cdc WHERE order_id = 3;
ALTER DYNAMIC TABLE ci_dt_orders_enriched REFRESH;{"insertedRows":1, "copiedRows":1, "deletedRows":1, "updatedRows":1}注意updatedRows——声明式动态表从未发出的统计信息,因为它们将更新表示为删除加插入。自定义增量表执行真正的UPDATE,因此会出现。这是您的MERGE分支按预期触发的可靠信号。
订单3已消失——删除已传播。
QUALIFY并非可选
当多个源行匹配一个目标行时,MERGE是不确定的,Snowflake不会保护您。使用INFORMATION => DEFAULT,更新以DELETE行和且相同键的INSERT行出现——两行匹配一个目标。
QUALIFY将它们折叠,平局决定者决定胜者:
ORDER BY c.updated_at DESC,
IFF(c.METADATA$ACTION = 'INSERT', 0, 1)updated_at DESC通常选择INSERT部分,因为新版本具有较晚的时间戳。但当列更改而未更改updated_at时,两个部分共享相同时间戳——而IFF确保INSERT仍然获胜,而不是行被错误删除。省略它会导致间歇性、静默的损坏。
工作示例:运行累加器
此模式在声明式方式下是不可能的。它读取自身先前的输出,因此从不重新扫描历史记录。
CREATE OR REPLACE DYNAMIC TABLE ci_dt_player_scores (
player_id NUMBER,
total_score NUMBER
)
TARGET_LAG = '1 minute'
WAREHOUSE = transform_wh
REFRESH_MODE = CUSTOM_INCREMENTAL
REFRESH USING (
MERGE INTO SELF AS tgt
USING (
SELECT player_id, SUM(score) AS batch_score
FROM ci_match_results CHANGES(INFORMATION => APPEND_ONLY)
GROUP BY player_id
) AS src
ON tgt.player_id = src.player_id
WHEN MATCHED THEN UPDATE SET
tgt.total_score = tgt.total_score + src.batch_score -- reuses prior state
WHEN NOT MATCHED THEN INSERT
(player_id, total_score) VALUES (src.player_id, src.batch_score)
);跨两个批次(10 + 15然后+5),玩家1最终达到30。刷新仅对两行新数据求和,并将其添加到存储的总数中。声明式等效项会每次重新聚合整个历史记录。
此处的APPEND_ONLY意味着更新和删除被忽略——通常适用于运行总计,因为撤回的匹配不应静默重写历史。对于撤回,切换到DEFAULT并在METADATA$ACTION = 'DELETE'时进行减法。
用于仅追加工作的INSERT INTO SELF
CREATE OR REPLACE DYNAMIC TABLE dt_deletions_log (
id INT, name STRING, email STRING
)
TARGET_LAG = '1 minute'
WAREHOUSE = transform_wh
INITIALIZE = ON_SCHEDULE
REFRESH USING (
INSERT INTO SELF
SELECT * EXCLUDE (METADATA$ISUPDATE, METADATA$ACTION)
FROM users CHANGES(INFORMATION => DEFAULT)
WHERE NOT METADATA$ISUPDATE AND METADATA$ACTION = 'DELETE'
);NOT METADATA$ISUPDATE使其成为对真正删除的审计,而不是被每次更新的删除部分污染日志。
初始化和迁移
没有BACKFILL FROM,初始刷新会重放每个现有源行通过您的逻辑,因为CHANGES()将它们全部视为INSERT。在大型表上,这既慢又昂贵。

START AT接受TIMESTAMP、STATEMENT => <query_id>、STREAM => <stream_name>或OFFSET => -<seconds>。
STREAM选项是从流和任务进行干净切换:将新动态表指向现有流的偏移量,以便不遗漏或重复计数任何更改。两个子句都仅在创建时生效——不能通过ALTER设置。
限制:
FROZEN WHERE和INSERT ONLY INPUTS不能与REFRESH USING结合使用。 冻结区域的成本控制和自定义MERGE逻辑是互斥的——选择其一。
不支持dbt或DCM集成。 只有CREATE OR ALTER可以修改REFRESH USING定义或其属性。
不推导主键。 声明式表从GROUP BY和QUALIFY分区推断键;这些不推断。添加一个,以便下游消费者保持高效:
ALTER TABLE ci_dt_orders_enriched ADD PRIMARY KEY (order_id) RELY;保留期必须超过刷新间隔。 如果在下次刷新之前,通过CHANGES()读取的基本表上的保留期到期,该刷新将失败。考虑计划中的暂停。
上游模式更改会导致下次刷新失败。 如果上游被修改,使用CREATE OR ALTER和更新的REFRESH USING恢复;刷新从上次成功位置继续。如果被CREATE OR REPLACE,更改跟踪被破坏,下游必须重新创建。
维度在快照时读取,而非增量。 CHANGES()之外的对象在刷新时按当前时间读取,因此编辑维度不会重新丰富现有行——只有后续更改集才会采用新值。如果需要,那是声明式连接。
每次刷新是一个自动提交事务。 任何失败都会回滚整个刷新。
类型限制: 在MERGE … ON中不支持结构化OBJECT/ARRAY/MAP、INTERVAL、地理空间列。不支持UDTF;CHANGES()内的非SQL标量UDF必须是IMMUTABLE。
RELY改变CHANGES()语义
如果删除后插入相同行应该对您的合并逻辑不可见——通常是INSERT OVERWRITE的结果——添加RELY,以便Snowflake在它们到达CHANGES()之前压缩这些对。APPEND_ONLY从不使用主键进行身份识别。
“无新数据”并非失败
在启用 TARGET_LAG 的情况下,调度器可能会在您手动执行 ALTER … REFRESH 之前消费变更集,并且 CHANGES() 只会被消费一次。在测试期间,连续两次手动刷新报告“无新数据”,而累加器实际上已被调度器正确更新。在假设刷新失败之前,请检查内容。
何时使用
将自定义增量用于 CDC,用于删除传播、流静态连接、状态复用(累加器、Top-K 排行榜)、审计跟踪,或将单语句 MERGE/INSERT逻辑从流和任务中迁移出来。
当 SELECT 能表达结果时,保持声明式。您可以获得延迟视图等价性、自动键派生、FROZEN WHERE、dbt 集成以及 Snowflake 的正确性保证。自定义增量会以所有这些换取表达力——仅在需要时才采用这种权衡。
迁移清单
按依赖顺序:
- 迁移到 Gen2 仓库。收益最大,无需重写。
- 向基表添加 PRIMARY KEY … RELY,然后对下游 Dynamic Table 执行 CREATE OR REPLACE,因为此操作不具有追溯性。
- 在历史不可变的位置添加冻结区域。使用 DATEADD,而不是 CURRENT_DATE() - n。
- 对任何具有现有历史的迁移使用 BACKFILL FROM。
- 将去重逻辑转换为 QUALIFY ROW_NUMBER() = 1,如果您希望进行模式演进,则使用 SELECT * EXCLUDE。
- 拆分单体表——先处理连接,再处理聚合。
- 设置 INITIALIZATION_WAREHOUSE以实现双仓库策略。
- 显式固定 REFRESH_MODE。使用 AUTO 发现,使用 INCREMENTAL 或 ADAPTIVE 部署。
- 使用 SHOW DYNAMIC TABLES 和刷新历史统计信息进行验证。
- 考虑自定义增量仅当 SELECT 确实无法表达转换时——带有删除传播的 CDC、流静态连接或累加器。请记住它与 FROZEN WHERE 不兼容,因此请决定哪个对该表更重要。
结语
本版本的主线是 Dynamic Table 在不放弃声明式模型的情况下获得了更强的表达能力。您描述所需的结果和新鲜度;Snowflake 负责增量维护、调度和依赖关系。
发生变化的是,边缘不再是障碍。不可变历史可以被冻结,而不是无限重算。现有数据可以被采纳,而不是重建。刷新策略可以针对每次运行自适应。当声明式模型确实需要逃生舱口时——在十年冻结历史中执行 GDPR 删除——现在有了一个狭窄且定义明确的舱口。而对于那些无法用 SELECT 表示的转换,自定义增量让您提供 MERGE,而 Snowflake 则保留调度、重试和依赖跟踪。
最后一个才是真正的转变。Dynamic Table 过去迫使您做出选择:声明式的便利 或命令式的控制。您现在可以在一个管道中同时使用两者——对于适合的转换使用声明式表,对于不适合的 CDC(带删除传播)或运行累加器使用自定义增量表。请记住权衡:自定义增量将正确性保证和控制权都交给了您。
如果您因为编排复杂性或表达能力缺口而一直推迟使用 Dynamic Table,那么这些阻碍大多已消除。从 Gen2 仓库和冻结区域开始;它们以最小的改动带来最大的价值。
所有示例均已在 2026 年 8 月的 AWS us-east-1 上的 Snowflake 中验证。功能可用性因账户和区域而异——在依赖预览功能之前,请先在您自己的环境中验证。
Dynamic Table 变得更快且更灵活——第 3 部分(完) 最初发布在 Snowflake Builders Blog: Data Engineers, App Developers, AI, & Data Science 上,Medium 上的人们正在通过亮点和回应来继续讨论这个故事。
这篇内容对你有用吗?
反馈只用于改善内容筛选,不等同于收藏