流与变更数据捕获

Snowflake 中的数据管道自动化

Emily Melhuish

Technical Curriculum Developer, Snowflake

只处理已变更的数据

CDC: Change Data Capture(变更数据捕获) 两种数据处理方式的对比

Snowflake 中的数据管道自动化

流(Streams)

流的作用

  • 跟踪源表上的每次 INSERT、UPDATE、DELETE

Screenshot 2026-05-11 at 10.50.31 am.png

  • 维护持续变更日志——不复制数据
  • 被消费后,偏移量前移;下一次读取从最新开始
1 * Snowflake 学习资源
Snowflake 中的数据管道自动化

流的类型

 

流类型 捕获内容 最佳适用场景
标准 所有表与视图的所有 DML 变更——插入、更新、删除 任意行可能变化的表(如 shipments)
仅追加 除外部表外的所有表与视图——仅跟踪插入 只插入一次的表(如 delivery events)——更高效
仅插入 外部管理的 Apache Iceberg 与外部表——仅跟踪插入 外部表

目录表 提供阶段的文件元数据(名称、大小、最后修改时间戳)

Snowflake 中的数据管道自动化

创建流

在 shipments 表上创建标准流

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

在 delivery events 表上创建仅追加流

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Snowflake 中的数据管道自动化

流的元数据列

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTIONINSERTDELETE
  • METADATA$ISUPDATE:当属于更新对时为 TRUE
  • METADATA$ROW_ID:唯一物理行标识符
  • 更新显示为一对 DELETE + INSERT,两者均标记为 METADATA$ISUPDATE = TRUE

显示 METADATA$ACTION、METADATA$ISUPDATE、METADATA$ROW_ID 样例数据的流元数据列

Snowflake 中的数据管道自动化

流的偏移量(Offset)

时间线图——创建流(偏移量从此开始)→ 源表发生变更(流累积记录)→ 在事务中消费流(偏移量推进到当前)

Snowflake 中的数据管道自动化

管道中的流:概览

  • 流与任务配合使用——任务是按计划运行 SQL 的 Snowflake 对象
  • 任务仅读取变更行;1000 万行示例:2 秒 vs 2 分钟

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake 学习资源
Snowflake 中的数据管道自动化

管道中的流:示例查询

  • 流与任务配合使用——任务是按计划运行 SQL 的 Snowflake 对象
  • 任务仅读取变更行;1000 万行示例:2 秒 vs 2 分钟
CREATE TASK logistics.sync_shipments
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('logistics.staging_shipments_stream')
AS
  INSERT INTO logistics.shipments
  SELECT shipment_id, region, carrier, delivery_days
  FROM logistics.staging_shipments_stream
  WHERE METADATA$ACTION = 'INSERT';
Snowflake 中的数据管道自动化

Vamos praticar!

Snowflake 中的数据管道自动化

Preparing Video For Download...