在 Snowflake 進行資料管線自動化
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture(變更資料擷取)

串流能做什麼

| 串流類型 | 捕捉內容 | 最適用於 |
|---|---|---|
| Standard | 所有資料表與檢視,且記錄所有 DML 變更——追蹤 insert、update、delete | 任一列都可能變更的資料表(如 shipments) |
| Append-only | 除外部表以外的所有資料表與檢視——僅追蹤列插入 | 僅插入一次的資料表(如 delivery events),效率更高 |
| Insert-only | 外部管理的 Apache Iceberg 與外部表——僅追蹤列插入 | 外部表 |
Directory tables 會呈現某個 stage 的檔案中繼資料(名稱、大小、最後修改時間戳)
在 shipments 資料表上建立 Standard 串流
CREATE STREAM shipments_stream
ON TABLE logistics.shipments;
在 delivery events 資料表上建立 Append-only 串流
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 進行資料管線自動化