Snowflakeにおけるデータパイプラインの自動化
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: 変更データキャプチャ

ストリームの機能

| ストリームタイプ | キャプチャ対象 | 適したユースケース |
|---|---|---|
| Standard | すべてのテーブル・ビュー、すべてのDML変更(INSERT・UPDATE・DELETE) | 任意の行が変更されうるテーブル(例:shipments) |
| Append-only | 外部テーブル以外のすべてのテーブル・ビュー、行のINSERTのみ | 一度だけ挿入されるテーブル(例:配送イベント)- 高効率 |
| Insert-only | 外部管理のApache Icebergテーブルおよび外部テーブル、行のINSERTのみ | 外部テーブル |
ディレクトリテーブルはステージのファイルメタデータ(名前・サイズ・最終更新日時)を参照可能にする
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におけるデータパイプラインの自動化