返回
RSS Snowflake Engineering (Medium) AI 逐段翻译 精选 发布 2026-08-27 22:01 收录于 09-03

Snowflake动态表更快更灵活:自定义增量实战(下)

DataHot 速览

Snowflake 2026年的Dynamic Tables迎来多项平台增强,本文是该系列第三部分,重点展示自定义增量动态表(CUSTOM_INCREMENTAL)与MERGE的实战用法。新功能允许用户自行编写MERGE或INSERT逻辑,而调度、重试、依赖跟踪和事务保障仍由Snowflake负责。文中给出了语法格式,并说明需要显式列清单、REFRESH USING中只能有一条DML,且不支持存储过程与多语句事务。作者提醒,自定义增量后数据一致性责任由用户承担,若MERGE逻辑出错,Snowflake不会自动发现。

为什么值得关注:数据从业者可从中了解Snowflake Dynamic Tables新增的自定义增量更新能力,兼顾声明式管道的便利与手动控制,适用于需要复杂MERGE逻辑的产线数据场景。

本文目录 14 节
  1. 一份针对Snowflake 2026年动态表更新的实践指南——所有示例均在真实账户上运行。
  2. 第9部分——使用MERGE的自定义增量动态表
  3. 语法
  4. 工作示例:带删除传播的CDC增强
  5. QUALIFY并非可选
  6. 工作示例:运行累加器
  7. 用于仅追加工作的INSERT INTO SELF
  8. 初始化和迁移
  9. 限制:
  10. RELY改变CHANGES()语义
  11. “无新数据”并非失败
  12. 何时使用
  13. 迁移清单
  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 的正确性保证。自定义增量会以所有这些换取表达力——仅在需要时才采用这种权衡。

迁移清单

按依赖顺序:

  1. 迁移到 Gen2 仓库。收益最大,无需重写。
  2. 向基表添加 PRIMARY KEY … RELY,然后对下游 Dynamic Table 执行 CREATE OR REPLACE,因为此操作不具有追溯性。
  3. 在历史不可变的位置添加冻结区域。使用 DATEADD,而不是 CURRENT_DATE() - n。
  4. 对任何具有现有历史的迁移使用 BACKFILL FROM
  5. 将去重逻辑转换为 QUALIFY ROW_NUMBER() = 1,如果您希望进行模式演进,则使用 SELECT * EXCLUDE。
  6. 拆分单体表——先处理连接,再处理聚合。
  7. 设置 INITIALIZATION_WAREHOUSE以实现双仓库策略。
  8. 显式固定 REFRESH_MODE。使用 AUTO 发现,使用 INCREMENTAL 或 ADAPTIVE 部署。
  9. 使用 SHOW DYNAMIC TABLES 和刷新历史统计信息进行验证
  10. 考虑自定义增量仅当 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 上的人们正在通过亮点和回应来继续讨论这个故事。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存