C データエンジニアリング — 在庫リアルタイムマート Datastream CDC(ソースタイムスタンプ+binlog位置での順序解決 × 削除イベント(_metadata_deleted)反映 × incremental merge戦略 × パーティション/クラスタリングでスキャンコスト削減 × on_schema_changeでスキーマ追従 × レプリケーション遅延のCloud Monitoring監視)(MOps 倉庫在庫管理サービス 欠品アラートLINE通知 Bad→Good)

2026-08-05 (Day 123) 水曜 C: データエンジニアリング ★★★★☆ Datastream (Cloud SQL for MySQL → BigQuery, append_only) dbt Core 1.9+ incremental/merge / Terraform 1.8+

概要

🧭

ソースDBの絶対順序で最新状態を解決

アプリ側の updated_at はクロックずれ・再送で信頼できない。Datastreamが付与する _metadata_source_timestamp と MySQL binlog位置(_metadata_log_file/_metadata_log_position)でタイブレークし、ソースDB内での実際の発生順序を一意に解決する。

🗑️

削除イベントを「消えた事実」として反映

append_only モードのCDCでは、DELETEも「行が消える」のではなく「_metadata_deleted=true の変更イベントが1行追加される」形で表現される。最新版が削除イベントならマートから除外する。

💰

incremental + パーティション/クラスタリング

追記専用の変更履歴テーブルを毎回フル再構築すると、コストが増え続ける時限爆弾になる。materialized='incremental' + partition_by/cluster_by でスキャン範囲を「差分」に固定する。

🔔

merge戦略 × on_schema_change × 監視

incremental_strategy='merge' で冪等なUPSERT、on_schema_change でソース側の無断カラム追加をサイレントに欠落させない。Datastreamのレプリケーション遅延自体もCloud Monitoringで監視する。

問題

ECサイト MOps チームでは、倉庫在庫管理サービス(Cloud SQL for MySQL 8.0)の在庫テーブル inventory.stockproduct_id, warehouse_id, quantity, reserved_quantity, updated_at)を Datastream で BigQuery にCDC(Change Data Capture)レプリケーションし、cdc_raw.inventory_changelog という変更履歴(changelog)テーブルに追記している。この変更履歴を dbt Core 1.9+ でリアルタイムに近い「現在庫数マート」(mart_inventory_realtime)に変換し、倉庫担当者向けの欠品アラートLINE通知(在庫が閾値を下回ったら自動通知)の判定に使っている。

以下の Terraform(Datastream構成)・dbtモデル・schema.yml には 7つの設計・データ品質上の問題 が潜んでいる。問題点を全て洗い出しDatastream変更履歴の正しい順序解決(source timestamp + binlog position)× 削除イベント(_metadata_deleted)の反映 × dbt incremental(merge戦略)× パーティション/クラスタリングによるスキャンコスト削減 × ソーススキーマ変更への追従(on_schema_change)× Datastreamレプリケーション遅延の監視アラート を活用した Bad→Good リファクタリングを行ってください。

制約・前提条件

  • Datastream(Cloud SQL for MySQL → BigQuery、GA機能。append_only モードで変更履歴をそのまま追記)
  • BigQuery Standard SQL、dbt Core 1.9+(incremental_strategy='merge', on_schema_change, model contract が利用可)
  • Terraform 1.8+(google プロバイダー 5.x)
  • cdc_raw.inventory_changelog のスキーマ: product_id, warehouse_id, quantity, reserved_quantity, updated_at(ソース側の更新時刻)に加え、Datastreamが自動付与するメタデータ列 _metadata_source_timestamp(ソースDBでのコミット時刻), _metadata_log_file/_metadata_log_position(MySQL binlogのファイル名・位置。同一秒内の複数変更の順序を一意に決定できる), _metadata_deleted(削除イベントなら true
  • Datastreamは変更を1行も更新・削除せず、全て新しい行として追記し続けるappend_only モード)。同一 (product_id, warehouse_id) の変更履歴が複数行たまる
  • reserved_quantity 列は最近、在庫管理サービス側チームがMOpsチームへの事前連携なしに追加した新しいカラム
  • 欠品アラートLINE通知は mart_inventory_realtime を5分ごとに参照して実行される。データが古い・不正確だと誤通知(在庫があるのに欠品通知/欠品なのに通知なし)につながる
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 Terraform コード(Datastream構成・監視アラート)+ 改善後 dbt モデル・schema.yml(incremental/merge・contract・on_schema_change)+ 実行例(input→output)+ 設計意図の説明

悪いコード (Before)

このコードには 7つの設計・データ品質上の問題 が隠れています。見つけてみてください。
bad_mart_inventory_realtime.sql — updated_atソート・削除未反映・フル再構築
{{ config(materialized='table') }}
-- 問題③: 変更履歴全体を毎回フル再構築

-- 問題①: アプリ側のupdated_atでソート
--        (クロックずれ・再送で順序が狂う)
-- 問題②: _metadata_deletedを一切フィルタしていない
--        (削除済みの在庫がそのまま残る)
SELECT
    product_id,
    warehouse_id,
    quantity,
    reserved_quantity,
    updated_at
FROM `my-project.cdc_raw.inventory_changelog`
QUALIFY ROW_NUMBER() OVER (
    PARTITION BY product_id, warehouse_id
    ORDER BY updated_at DESC   -- 問題①
) = 1
-- 問題④: 重複排除の粒度がQUALIFYの1回のみで
--        mergeによる冪等性の担保が無い
-- 問題⑤: WHERE句が無く、変更履歴テーブル
--        全体をスキャンしている
-- 問題⑥: reserved_quantity列が将来変更されても
--        検知できない(on_schema_change未設定)
bad_datastream.tf — 監視ゼロのDatastream構成
resource "google_datastream_stream" "inventory_cdc" {
  stream_id     = "inventory-cdc-stream"
  location      = "asia-northeast1"
  desired_state = "RUNNING"

  source_config {
    source_connection_profile =
      google_datastream_connection_profile.mysql_source.id
    mysql_source_config {
      include_objects {
        mysql_databases {
          database = "inventory"
          mysql_tables { table = "stock" }
        }
      }
    }
  }

  destination_config {
    destination_connection_profile =
      google_datastream_connection_profile.bq_dest.id
    bigquery_destination_config {
      single_target_dataset { dataset_id = "cdc_raw" }
      append_only {}
      data_freshness = "900s"
    }
  }
}
# 問題⑦: レプリケーション遅延・ストリーム停止を
#        検知するCloud Monitoringアラートが無い
問題点サマリー(7点)
1updated_atで順序解決 — source_timestamp+binlog位置へ
2削除イベント未フィルタ — 最新版が削除なら除外へ

ヒント(段階的開示)

ヒント1 — 方向性
問題は3層に分類できる。(1) 変更履歴の正しい解決 — Datastreamは変更を追記するだけなので、「今この商品・倉庫の在庫は何個か」を求めるには、変更履歴の中から正しい順序で最新の1件を選び出す必要がある。単純に updated_at(アプリ側がセットする時刻。クロックずれ・再送で不正確なことがある)でソートしたり、削除イベント(_metadata_deleted=true)を無視すると、古い値や削除済みの在庫が「現在の値」として採用されてしまう。(2) スキャンコストと冪等性 — 変更履歴は際限なく増え続けるテーブルであり、dbtモデルがこれを毎回フル再構築すると、コストが増大し続ける。パーティション・クラスタリングで対象範囲を絞り、merge incremental戦略で「差分だけを冪等に反映」する設計が必要。(3) スキーマ変更への追従と監視 — ソースDB側で列が追加された場合にモデルが追従できるか、そしてDatastreamのレプリケーション自体が止まっている・遅延している場合にそれを検知できるかという「気づく仕組み」の欠如がガバナンス上の問題になる。
ヒント2 — アプローチ
  • 問題①: 変更履歴から最新状態を選ぶ際、アプリ側の updated_at でソートしている → Datastreamの _metadata_source_timestamp(ソースDBのコミット時刻)を主キーにし、同一タイムスタンプ内の順序は _metadata_log_file/_metadata_log_position(MySQL binlogの物理位置)でタイブレークする
  • 問題②: 削除イベント(_metadata_deleted=true)をフィルタしておらず、削除済みの在庫行が「在庫あり」として残り続ける → 最新版が削除イベントの場合はマートから除外する
  • 問題③: dbtモデルが変更履歴を毎回フル再構築(materialized='table')しており、履歴が増えるほどコストが増大する → incremental + partition_by_metadata_source_timestampの日付)+ cluster_byproduct_id, warehouse_id)で対象範囲を絞る
  • 問題④: incrementalモデルが単純追記で、同一キーの重複行がマートに混在している → incremental_strategy='merge' + unique_key で冪等なUPSERTにする
  • 問題⑤: 増分実行のたびに変更履歴テーブル全体をスキャンしており、パーティションを絞れていない → WHERE句で「前回実行以降のパーティションのみ」に限定する
  • 問題⑥: ソース側で reserved_quantity 列が無断追加された際、dbtモデル・contractがそれを検知せずサイレントに欠落する → on_schema_change='append_new_columns' を設定し、想定外の列追加をビルドログで可視化する
  • 問題⑦: Datastreamのレプリケーション遅延・ストリーム停止を検知するアラートが無く、欠品アラートLINE通知が古いデータのまま動き続けるリスクに気づけない → Datastreamのレイテンシメトリクスに対する Cloud Monitoring アラートポリシーを追加する
ヒント3 — コードの骨格
-- 最新版の解決(骨格)
QUALIFY ROW_NUMBER() OVER (
  PARTITION BY product_id, warehouse_id
  ORDER BY _metadata_source_timestamp DESC,
           _metadata_log_file DESC,
           _metadata_log_position DESC
) = 1
# dbt incremental設定の骨格
{{ config(
    materialized='incremental',
    incremental_strategy='merge',
    unique_key=['product_id', 'warehouse_id'],
    partition_by={'field': 'change_date', 'data_type': 'date'},
    cluster_by=['product_id', 'warehouse_id'],
    on_schema_change='append_new_columns'
) }}

問題点分析(7点)

#問題点分類改善方法
1最新状態の解決にアプリ側updated_atを使い順序が不安定データ品質_metadata_source_timestamp + binlog位置でタイブレーク
2削除イベント(_metadata_deleted)を無視データ品質最新版が削除イベントならマートから除外
3変更履歴を毎回フル再構築(materialized='table')コストincremental + partition_by/cluster_by
4単純追記で同一キーの重複行が混在データ品質incremental_strategy='merge' + unique_key
5増分実行のたびに変更履歴全体をスキャンコストパーティションフィルタで対象範囲を限定
6ソース側の無断カラム追加にモデルが追従しない保守性/ガバナンスon_schema_change='append_new_columns'
7Datastreamのレプリケーション遅延・停止の監視なし監視Cloud Monitoring アラートポリシー

CDC順序解決の仕組み — Bad vs Good の比較

Bad(変更前)— updated_atソート・削除未反映 cdc_raw.inventory_changelog(P100 / W1) ① quantity=50 updated_at=10:00:00 (先発) ② quantity=30 updated_at=10:00:00 (再送・同一時刻) ③ DELETE updated_at=10:00:01 (実は①②より先にDBでcommit済み) 問題①: updated_atでは②③の到着順序が保証されない 問題②: _metadata_deletedを見ていないので③(削除)が無視される → 誤って quantity=30 が「現在庫」として採用されうる bad_mart_inventory_realtime(materialized='table') 問題③: 毎回変更履歴の全期間をフルスキャン・フル再構築 問題④⑤: 重複排除・パーティション絞り込みが不十分 欠品アラートLINE通知(5分ごと) 問題⑥⑦: スキーマ追従なし・レプリケーション遅延の監視なし × 削除済みの在庫が「在庫あり」として通知判定に使われる × 変更履歴が増えるほどBigQueryコストが増加し続ける × reserved_quantity追加のような変更にサイレントに追従できない × Datastreamが遅延・停止していても誰も気づけない Good(変更後)— source_timestamp+binlog順序解決 同じ3件を _metadata_source_timestamp + binlog位置で解決 ① log_position=4820 quantity=50 ② log_position=4860 DELETE (_metadata_deleted=true) ③ log_position=4915 quantity=30 修正①: binlog位置がソースDB内の絶対順序を一意に表す 修正②: 最新版(③)を採用しつつ_metadata_deletedも正しく判定 → 最終的に quantity=30 が正しい現在庫として確定 stg_inventory_changelog(incremental + merge) 修正③: partition_by(change_date) + cluster_by(product_id, warehouse_id) 修正④⑤: merge戦略 + パーティションフィルタで差分のみスキャン mart_inventory_realtime → 欠品アラートLINE通知 修正⑥⑦: on_schema_change + Datastream遅延の Cloud Monitoring 監視 ✓ ソースDB内の絶対順序で現在の在庫状態を正しく解決 ✓ merge戦略でスキャン範囲・重複行を最小化しコストを抑制 ✓ 無断カラム追加もon_schema_changeでビルドログに可視化 ✓ レプリケーション遅延をビジネス影響が出る前に検知 修正

模範解答

# infra/datastream_governance.tf

resource "google_datastream_stream" "inventory_cdc" {
  stream_id     = "inventory-cdc-stream"
  location      = "asia-northeast1"
  desired_state = "RUNNING"

  source_config {
    source_connection_profile = google_datastream_connection_profile.mysql_source.id
    mysql_source_config {
      include_objects {
        mysql_databases {
          database = "inventory"
          mysql_tables { table = "stock" }
        }
      }
    }
  }

  destination_config {
    destination_connection_profile = google_datastream_connection_profile.bq_dest.id
    bigquery_destination_config {
      single_target_dataset { dataset_id = "cdc_raw" }
      append_only {}
      data_freshness = "300s"   # 欠品アラート通知の5分周期に合わせSLA明示
    }
  }
}

# 修正⑦: Datastreamのレプリケーション遅延を監視し、閾値超過でアラート
resource "google_monitoring_alert_policy" "datastream_latency" {
  project      = var.project_id
  display_name = "inventory-cdc-stream レプリケーション遅延超過"
  combiner     = "OR"
  conditions {
    display_name = "total_latencies が 600秒(欠品アラートSLAの2倍)を超過"
    condition_threshold {
      filter          = <<-EOT
        resource.type="datastream.googleapis.com/Stream"
        resource.labels.stream_id="inventory-cdc-stream"
        metric.type="datastream.googleapis.com/stream/total_latencies"
      EOT
      comparison      = "COMPARISON_GT"
      threshold_value = 600
      duration        = "120s"
      aggregations {
        alignment_period   = "300s"
        per_series_aligner = "ALIGN_PERCENTILE_99"
      }
    }
  }
  notification_channels = [var.oncall_notification_channel_id]
}

# 修正⑦: ストリーム自体が異常停止した場合の検知
resource "google_monitoring_alert_policy" "datastream_stream_down" {
  project      = var.project_id
  display_name = "inventory-cdc-stream ストリーム異常停止"
  combiner     = "OR"
  conditions {
    display_name = "イベント流入が10分間停止"
    condition_threshold {
      filter          = <<-EOT
        resource.type="datastream.googleapis.com/Stream"
        resource.labels.stream_id="inventory-cdc-stream"
        metric.type="datastream.googleapis.com/stream/event_count"
      EOT
      comparison      = "COMPARISON_LT"
      threshold_value = 1
      duration        = "600s"
    }
  }
  notification_channels = [var.oncall_notification_channel_id]
}
Before — updated_atソート・削除未反映・フル再構築
{{ config(materialized='table') }}

SELECT
    product_id, warehouse_id, quantity,
    reserved_quantity, updated_at
FROM `my-project.cdc_raw.inventory_changelog`
QUALIFY ROW_NUMBER() OVER (
    PARTITION BY product_id, warehouse_id
    ORDER BY updated_at DESC
) = 1
-- 削除フィルタなし・パーティション絞り込みなし
-- on_schema_changeなし
After — 順序解決 + merge + パーティション/クラスタリング
-- models/staging/stg_inventory_changelog.sql
-- 修正③⑤: incremental化しパーティション+クラスタリングで
--         スキャン範囲を「新しく追記された変更履歴」に限定
{{ config(
    materialized='incremental',
    incremental_strategy='merge',
    unique_key=['product_id', 'warehouse_id'],
    partition_by={'field': 'change_date', 'data_type': 'date',
                  'granularity': 'day'},
    cluster_by=['product_id', 'warehouse_id'],
    on_schema_change='append_new_columns'   -- 修正⑥
) }}

WITH changelog AS (
    SELECT
        product_id, warehouse_id, quantity, reserved_quantity,
        _metadata_source_timestamp, _metadata_log_file,
        _metadata_log_position, _metadata_deleted,
        DATE(_metadata_source_timestamp) AS change_date
    FROM {{ source('cdc_raw', 'inventory_changelog') }}
    {% if is_incremental() %}
    -- 修正⑤: 前回実行以降のパーティションのみスキャン
    WHERE DATE(_metadata_source_timestamp) >= (
        SELECT DATE_SUB(MAX(change_date), INTERVAL 1 DAY) FROM {{ this }}
    )
    {% endif %}
),
latest_per_key AS (
    -- 修正①: ソースDBのコミット時刻+binlog位置で
    --        実際の発生順序を一意に解決する
    SELECT *,
        ROW_NUMBER() OVER (
            PARTITION BY product_id, warehouse_id
            ORDER BY _metadata_source_timestamp DESC,
                     _metadata_log_file DESC,
                     _metadata_log_position DESC
        ) AS rn
    FROM changelog
)
SELECT
    product_id, warehouse_id, quantity, reserved_quantity,
    _metadata_source_timestamp AS last_changed_at, change_date
FROM latest_per_key
WHERE rn = 1
  AND _metadata_deleted = FALSE   -- 修正②: 削除イベントを除外
mart_inventory_realtime.sql — 修正④: merge戦略で冪等なUPSERT
{{ config(
    materialized='incremental',
    incremental_strategy='merge',
    unique_key=['product_id', 'warehouse_id'],
    partition_by={'field': 'change_date', 'data_type': 'date', 'granularity': 'day'},
    cluster_by=['product_id', 'warehouse_id']
) }}

SELECT
    product_id,
    warehouse_id,
    quantity - reserved_quantity AS available_quantity,  -- 欠品判定に使う実質在庫数
    quantity,
    reserved_quantity,
    last_changed_at,
    change_date
FROM {{ ref('stg_inventory_changelog') }}
{% if is_incremental() %}
WHERE change_date >= (SELECT DATE_SUB(MAX(change_date), INTERVAL 1 DAY) FROM {{ this }})
{% endif %}
# models/staging/_sources.yml
version: 2

sources:
  - name: cdc_raw
    database: my-project
    schema: cdc_raw
    tables:
      - name: inventory_changelog
        description: "Datastream (Cloud SQL for MySQL → BigQuery) による在庫変更履歴。append_onlyモードで全変更を追記。"
        loaded_at_field: _metadata_source_timestamp
        freshness:
          warn_after: { count: 10, period: minute }
          error_after: { count: 20, period: minute }   # 欠品アラート5分周期に対する遅延許容の上限

# models/marts/_marts.yml
models:
  - name: mart_inventory_realtime
    description: "商品×倉庫の現在庫数マート。欠品アラートLINE通知の判定に使用。"
    config:
      contract:
        enforced: true
    columns:
      - name: product_id
        data_type: string
        constraints: [{ type: not_null }]
      - name: warehouse_id
        data_type: string
        constraints: [{ type: not_null }]
      - name: available_quantity
        data_type: int64
      - name: last_changed_at
        data_type: timestamp
    tests:
      - dbt_utils.unique_combination_of_columns:   # 修正④の検証: 重複キーが残っていないか
          combination_of_columns: [product_id, warehouse_id]
問題修正内容効果
① updated_atで順序解決_metadata_source_timestamp + binlog位置ソースDB内の絶対順序で最新状態を確定
② 削除イベント未フィルタ_metadata_deleted=FALSEでフィルタ削除済み在庫がマートに残らない
③ 毎回フル再構築incremental + partition_by/cluster_byスキャン範囲を差分に固定しコスト抑制
④ 重複行が混在incremental_strategy='merge' + unique_key冪等なUPSERTでリトライ時も結果不変
⑤ 全期間スキャンパーティションフィルタ差分のみスキャンしコストをさらに抑制
⑥ 無断カラム追加に非追従on_schema_change='append_new_columns'スキーマ変更をビルドログで可視化
⑦ 遅延・停止の監視なしCloud Monitoring アラートポリシービジネス影響前にレプリケーション異常を検知

実行例(input → output)

input(cdc_raw.inventory_changelog、商品P100・倉庫W1の変更履歴、簡略化):

product_id  warehouse_id  quantity  updated_at(app側)         _metadata_source_timestamp  _metadata_log_file  _metadata_log_position  _metadata_deleted
P100        W1            50        10:00:00                   10:00:00.100                binlog.000123       4820                     false
P100        W1            30        10:00:00 (再送で同一値!)   10:00:00.310                binlog.000123       4915                     false
P100        W1            NULL      10:00:01                   10:00:00.900                binlog.000123       4860                     true   (削除イベント。DB上では上記2件より先にcommitされていたが到着が遅れた)

output(Before・問題のある実装):

updated_at でソートすると quantity=30(アプリ側時刻の粒度・再送で到着順が不安定)と削除イベントの前後関係が崩れ、削除フィルタも無いため quantity=30 が「現在庫」として誤って採用され続けるリスクがある。

output(After・改善後):

_metadata_log_position(binlog位置)でソートすると、ソースDB内の絶対順序は 4820 → 4860(削除) → 4915 と確定する。最新版は log_position=4915quantity=30_metadata_deleted=false)となり、正しく「現在庫30個」として mart_inventory_realtime に反映される。仮にDELETEが最後に発生していた場合は _metadata_deleted=FALSE フィルタにより当該商品はマートから正しく除外される。

ポイント解説

1CDCの「変更履歴」から「現在の状態」を求めるには、ソースDB内の絶対順序が必要 — アプリケーションがセットする updated_at は、アプリサーバー間のクロックずれ・リトライによる重複送信・タイムゾーン設定ミスなどで信頼できないことが多い。Datastreamが自動付与する _metadata_source_timestamp と、MySQLなら _metadata_log_file/_metadata_log_position(binlogの物理位置)を組み合わせることで、「ソースDBの中で実際にどちらが先に起きたか」を一意に決定できる。
2append_onlyモードの変更履歴は「削除された」という事実も1行のデータとして表現される — リレーショナルDBの DELETE はDatastreamの世界では「行が消える」のではなく「_metadata_deleted=true の変更イベントが1行追加される」形で表現される。この設計を理解していないと、削除フィルタを書き忘れて「消えたはずのデータが下流で生き続ける」というバグを埋め込みやすい。
3materialized='table'の全件再構築は、変更履歴テーブルのように無限に増え続けるソースでは破滅的にスケールしない — 変更履歴は性質上、追記専用で日々増加し続ける。これを毎回フルスキャン・フル再構築すると、テーブルサイズに比例してBigQueryコストが増加し続ける「時限爆弾」になる。incremental + パーティションフィルタによって、スキャン対象を「前回実行以降の差分」に固定化するのが定石。
4incremental_strategy='merge'は「冪等性」を保証するための選択 — 単純な追記のみの増分戦略では、dbt実行が失敗してリトライされた際や同じ変更履歴の行を2回処理してしまった際に、同一キーの重複行がマートに混入するリスクがある。unique_key を指定した merge 戦略は「同じキーが来たら上書き」という冪等なUPSERTになるため、同じバッチを2回流しても結果が変わらないという実務上非常に重要な性質を得られる。
5on_schema_change='append_new_columns'は「サイレントなデータ欠落」を防ぐガードレール — ソース側チームが事前連携なしに列を追加するのは実務でよく起きる。何の設定もしないと、限定列挙モデルでは新しい列が単に無視され続け、気づかれないまま「実は取得できていたはずのデータ」を長期間失う。明示的に設定し、model contractでスキーマを固定することで、変更を「エラーとして気づける」形に変える。
6監視は「パイプラインが動いているつもり」を防ぐ最後の防波堤 — Datastreamのレプリケーションが遅延・停止しても、BigQuery側のテーブル自体は(過去のデータのまま)存在し続けるため、dbtやダッシュボード側からは一見正常に見えてしまう。total_latencies メトリクスへのアラートは、「データが古いまま止まっている」という状態を、ビジネス影響(誤った欠品アラート)が出る前に検知するための仕組み。
7available_quantity = quantity - reserved_quantityという業務ロジックはstaging層ではなくmart層に置くreserved_quantity(引当済み数量)の意味づけや計算式は業務要件であり、将来ルールが増える可能性がある。CDCの順序解決・重複排除という「データ基盤としての正しさ」を担保するstaging層と、業務ロジックを表現するmart層を分離しておくことで、業務ルール変更時の影響範囲をmart層に限定できる。

実務への応用

ECサイト MOps チームが扱う在庫データは、倉庫管理サービス・注文管理サービス・返品サービスなど複数のマイクロサービスにまたがって更新される。CDCでリアルタイムに近い連携をする設計は、バッチ集計に比べて鮮度は高いが、その分「変更履歴から正しい状態をどう組み立てるか」という設計の巧拙が直接ビジネス影響(誤った欠品通知・過剰発注・機会損失)に跳ね返る。特に本問のような「アラート・自動発注のトリガーとして使われるデータ」は、鮮度だけでなく正確性(重複排除・削除反映・順序保証)の3点セットが揃って初めて信頼できるデータ基盤になる。

組織的な連携も並行して整備する: ソースDB側チームとの「無断スキーマ変更」問題は、マイクロサービスが増えるほど頻発する。on_schema_change やmodel contractのようなdbt側の防御に加えて、実務では「ソースDB側の変更管理プロセスにMOpsチームをレビュワーとして加える」といった組織的な連携も並行して整備するのが望ましい。

今日のまとめ

Datastream CDCによるリアルタイム在庫マートの7点チェックリスト: ① アプリ側updated_atではなく_metadata_source_timestamp+binlog位置で正しい順序を解決 ② _metadata_deletedを必ずフィルタし削除イベントを反映 ③ 変更履歴の全件再構築をやめincremental+パーティション/クラスタリングでコスト削減 ④ incremental_strategy='merge'で冪等なUPSERT ⑤ 増分実行のスキャン範囲をパーティションフィルタで限定 ⑥ on_schema_changeでソース側の無断スキーマ変更をサイレントな欠落にしない ⑦ Datastreamのレプリケーション遅延・停止をCloud Monitoringで監視。 CDCパイプラインの本質は「変更履歴から正しい現在の状態をどう一意に組み立てるか」にあり、順序解決・削除反映・冪等性の3点を欠くと、下流のアラート・自動化が静かに誤動作し続けるリスクを抱える。

次のステップ

  • 発展問題: 現在の設計は「5分ごとにバッチでmartを再計算する」擬似リアルタイムだが、真にリアルタイムな欠品アラートが必要な場合、Datastream → Pub/Sub → Dataflow(Streaming)で直接ストリーム処理し、BigQueryへの書き込みを介さずにアラート判定する設計に拡張するとどう変わるか、コスト・複雑性とのトレードオフを整理せよ。また、_metadata_log_file/_metadata_log_position はMySQL固有のメタデータだが、ソースがPostgreSQLの場合(LSN: Log Sequence Number)ではどう読み替えるべきかを調べよ。
  • 参考: Datastream(Cloud SQL for MySQL → BigQuery, append_only モード)公式ドキュメント / dbt incremental_strategy: merge / dbt on_schema_change / dbt source freshness / BigQuery パーティション分割テーブル・クラスタ化テーブル / Change Data Capture(CDC)設計パターン

自己評価(あとで記入)