1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

SnowflakeにおけるStream(CDC)を整理する

1
Posted at

Snowflakeを使っていて、「外部ストレージからの取り込み(Snowpipe)に対して、テーブルの変更差分を追跡するStream(CDC)はどのような役割を持つのか」「なぜテーブルの変更追跡(CDC)が必要なのか」が気になったので整理しました。


1. ざっくりまとめ

機能 分類 役割 扱う対象
Snowpipe Ingest 外部ステージ上のファイルをテーブルへロードする 外部ステージ上のファイル(S3/GCS等)
Streams CDC テーブルに対する変更差分(INSERT/UPDATE/DELETE)を追跡する Time Travelメタデータ
Tasks ワークフロー制御 定期実行またはストリームのデータ到着をトリガーにSQLを実行する SQLステートメント・ストアド
Dynamic Tables 宣言型データ変換 定義したクエリに基づき、ターゲットテーブルの差分リフレッシュを自動実行する ターゲットテーブルとソーステーブル

外部ステージからテーブルへデータを取り込むのが Ingest(Snowpipe)、取り込んだ後のテーブルに対する変更差分を検知・連携するのが CDC(Streams) です。


2. そもそもCDC(Change Data Capture)とは何か

CDC(Change Data Capture)とは、データベースのテーブルに対して発生した変更(INSERT、UPDATE、DELETE)を検知・キャプチャし、後続のテーブルや外部システムへ連携する技術です。

従来の2大アプローチとその限界

CDCを使わずにテーブル間のデータ同期や集計を行う場合、主に以下の2つの方法が採られてきました。しかし、どちらもデータ規模が大きくなるにつれて破綻します。

① 全件洗い替え(Full Refresh)の限界

集計対象のテーブル全体を毎回全件スキャンして洗い替える手法です。

  • 課題: 生データテーブルが数億行に達すると、1日数万件の新しいデータを反映するためだけに毎回全件を読み直すことになり、ウェアハウスのコンピュートコスト(クレジット)と処理時間が跳ね上がります。

② タイムスタンプ抽出(updated_at > ...)の限界

「前回実行した時刻」より新しいタイムスタンプを持つレコードだけを WHERE updated_at > :last_sync_time で取得する手法です。
一見合理的に見えますが、実務では以下の問題が頻発します。

  • 遅延到着データの取りこぼし: ネットワーク遅延や障害復旧で、過去のタイムスタンプを持ったデータが後から到着した場合、条件に引っかからず永続的に取りこぼします。
  • 物理削除(DELETE)の検知不可: 元テーブルで行が削除された場合、削除された行自体が存在しないため、タイムスタンプによる検索では削除イベントを検知できません。

CDCは、「データの値」を見るのではなく、「発生した変更イベント(追加・更新・削除)そのもの」をログやメタデータから確実に捉えることで、これらの問題を解決します。


3. SnowflakeにおけるCDCの実現機能

Snowflakeでは、CDCを実現するための機能がデータベースエンジンにネイティブ統合されています。

① Streams(ストリームオブジェクト)

テーブルの変更履歴を追跡するためのコア機能です。

CREATE OR REPLACE STREAM orders_stream ON TABLE raw_orders;

Streamsの内部アーキテクチャには以下の特徴があります。

  • 実データの二重保持を行わない(ゼロコピー): Streamを作成しても、追加のストレージは消費しません。Snowflakeの Time Travel メタデータ を参照し、前回処理時点からの差分(デルタ)を仮想的なテーブルとして見せています。
  • オフセットの自動前進: Streamからデータを読み取り、トランザクション内で別のテーブルへ書き込む(DMLを実行する)と、Streamの位置(オフセット)が自動的に最新時点まで進みます。明示的にオフセットを管理するコードを書く必要はありません。
  • メタデータ疑似列の付与: Streamをクエリすると、通常のカラムに加えて以下の制御列が自動付与されます。
    • METADATA$ACTION: 変更種別(INSERT または DELETE)
    • METADATA$ISUPDATE: UPDATEによる変更かどうか(TRUE / FALSE)
    • METADATA$ROW_ID: 行の一意な識別子

ストリームの種類

種類 構文オプション 追跡対象 用途
標準ストリーム 指定なし(デフォルト) INSERT / UPDATE / DELETE 通常のテーブル同期・マージ
追加専用ストリーム APPEND_ONLY = TRUE INSERT のみ 追記型ログ、ファクトデータの取り込み
外部テーブル用ストリーム 外部テーブルに対して作成 ファイルの追加・削除 データレイク上のファイル変更検知

② CHANGES 句

Streamオブジェクトを事前に作成していなくても、Time Travelが有効な期間内であれば、SQLの CHANGES 句を使ってアドホックに変更履歴を取得できます。

SELECT *
FROM raw_orders
CHANGES(INFORMATION => DEFAULT)
AT(TIMESTAMP => $last_sync_time);

③ Dynamic Tables(動的テーブル)

Streams と Tasks を使ったパイプラインを、より高水準に自動化したオブジェクトです。裏側でSnowflakeがCDCのオフセット管理と差分リフレッシュを完全に自動実行します。


4. Ingest(Snowpipe)と CDC(Streams)の境界

アーキテクチャ設計において混同しやすいのが Snowpipe と Streams の住み分けです。

  • Snowpipe: 外部ストレージ上のファイルを検知し、Snowflakeのテーブル(Raw層)へロードする。
  • Streams: テーブルにロードされたデータのうち、追加・更新・削除された変更差分のみを検知し、後続の変換処理(Silver/Gold層・データマート)へ渡す。

Snowpipeはデータの「取り込み(Ingest)」を担い、Streamsはその後の「変更差分の追跡・連携(CDC)」を担うため、両者はパイプラインの前段と後段として組み合わせて利用されます。


5. 具体的なユースケースと実装パターン

ユースケース1: ELTパイプラインの差分マージ(Raw ➔ Mart)

最も一般的なユースケースです。生データテーブル(raw_orders)の変更を、集計マート(mart_orders)へ差分のみ MERGE します。

実装手順(Streams + Tasks)

1. Streamの作成

CREATE OR REPLACE STREAM raw_orders_stream ON TABLE raw_orders;

2. Taskの作成(データがある時だけ動かす)

SYSTEM$STREAM_HAS_DATA を条件に指定することで、変更データが1件も到着していない時間帯はウェアハウスを一切起動せず、クレジットを消費しません。

CREATE OR REPLACE TASK sync_orders_task
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('raw_orders_stream')
AS
MERGE INTO mart_orders t
USING (
    -- Streamから差分を取得
    SELECT * FROM raw_orders_stream
) s
ON t.order_id = s.order_id

-- ① 削除されたレコードの反映
WHEN MATCHED AND s.METADATA$ACTION = 'DELETE' AND s.METADATA$ISUPDATE = FALSE THEN
  DELETE

-- ② 更新されたレコードの反映
WHEN MATCHED AND s.METADATA$ACTION = 'INSERT' AND s.METADATA$ISUPDATE = TRUE THEN
  UPDATE SET
    t.customer_id = s.customer_id,
    t.amount      = s.amount,
    t.status      = s.status,
    t.updated_at  = s.updated_at

-- ③ 新規追加されたレコードの挿入
WHEN NOT MATCHED AND s.METADATA$ACTION = 'INSERT' THEN
  INSERT (order_id, customer_id, amount, status, updated_at)
  VALUES (s.order_id, s.customer_id, s.amount, s.status, s.updated_at);

3. タスクの開始

ALTER TASK sync_orders_task RESUME;

このTaskが正常に実行完了すると、raw_orders_stream のオフセットが自動的に前進し、次回はそれ以降の変更のみが対象になります。


ユースケース2: 宣言型パイプライン(Dynamic Tables)

Streams と Tasks による MERGE 文の記述・保守を省き、SQLクエリと目標鮮度(TARGET_LAG)の指定だけで差分更新を行いたい場合のパターンです。

CREATE OR REPLACE DYNAMIC TABLE mart_orders_dt
  TARGET_LAG = '10 MINUTE'
  WAREHOUSE = compute_wh
AS
SELECT
    order_id,
    customer_id,
    amount,
    status,
    updated_at
FROM raw_orders;

Snowflakeが内部でソーステーブルの変更履歴を自動追跡し、TARGET_LAG に応じて増分リフレッシュ(Incremental Refresh)を実行します。


6. 更新手法の判断基準

すべてのテーブルにStreamを導入する必要はありません。データの特性や要件に応じて選択します。

手法 適する条件 メリット 注意点
全件洗い替え
(INSERT OVERWRITE)
- データ量が小さい(数十万行程度)
- 複雑な集計があり差分更新が難しい
- 最もシンプルで壊れにくい
- 差分ロジックの管理が不要
- データ量が増えるとコストと実行時間が急増する
Streams + Tasks - 大規模テーブルの差分更新
- DELETEの厳密な同期が必要
- 複数テーブルへの分岐更新や外部通知
- きめ細やかな制御が可能
- SYSTEM$STREAM_HAS_DATA による確実なコスト抑制
- DML文(MERGE)の自前実装と保守が必要
Dynamic Tables - 中〜大規模テーブルの加工・集計
- 宣言的なSQLで簡潔に保ちたい
- Stream/Taskの管理が不要
- データリネージの可視化が容易
- クエリ構造によっては増分更新できず全件リフレッシュになる場合がある

7. まとめ

  • CDCは差分処理の必須要件: 全件洗い替えのコスト増加や、タイムスタンプ抽出による遅延取りこぼし・削除未検知を防ぐために不可欠な技術。
  • SnowpipeとStreamsの役割分離: 外部からテーブルへロードする「Ingest」がSnowpipe、取り込んだ後のテーブル変更を追跡する「CDC」がStreams。
  • ゼロコピーとオフセット管理: SnowflakeのStreamはTime Travelメタデータを活用するため追加ストレージがかからず、DML実行によってオフセットが安全に自動更新される。
  • 実装の選択肢: 手動制御とコスト最適化を極めるなら Streams + Tasks、パイプライン運用のシンプルさを取るなら Dynamic Tables。

テーブル規模と必要な制御レベルに合わせて、適切な差分更新パターンを選定することが重要です。


参考

1
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
1
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?