Snowflake 中的数据管道自动化
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture(变更数据捕获)

流的作用

| 流类型 | 捕获内容 | 最佳适用场景 |
|---|---|---|
| 标准 | 所有表与视图的所有 DML 变更——插入、更新、删除 | 任意行可能变化的表(如 shipments) |
| 仅追加 | 除外部表外的所有表与视图——仅跟踪插入 | 只插入一次的表(如 delivery events)——更高效 |
| 仅插入 | 外部管理的 Apache Iceberg 与外部表——仅跟踪插入 | 外部表 |
目录表 提供阶段的文件元数据(名称、大小、最后修改时间戳)
在 shipments 表上创建标准流
CREATE STREAM shipments_stream
ON TABLE logistics.shipments;
在 delivery events 表上创建仅追加流
CREATE STREAM delivery_events_stream
ON TABLE logistics.delivery_events
APPEND_ONLY = TRUE;
SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
METADATA$ACTION:INSERT 或 DELETEMETADATA$ISUPDATE:当属于更新对时为 TRUEMETADATA$ROW_ID:唯一物理行标识符METADATA$ISUPDATE = TRUE


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 中的数据管道自动化