概要
ソース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.stock(product_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分ごとに参照して実行される。データが古い・不正確だと誤通知(在庫があるのに欠品通知/欠品なのに通知なし)につながる
悪いコード (Before)
{{ 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未設定)
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アラートが無い
ヒント(段階的開示)
ヒント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_by(product_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' |
| 7 | Datastreamのレプリケーション遅延・停止の監視なし | 監視 | Cloud Monitoring アラートポリシー |
CDC順序解決の仕組み — Bad vs Good の比較
模範解答
# 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]
}
{{ 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なし
-- 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 -- 修正②: 削除イベントを除外
{{ 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=4915(quantity=30、_metadata_deleted=false)となり、正しく「現在庫30個」として mart_inventory_realtime に反映される。仮にDELETEが最後に発生していた場合は _metadata_deleted=FALSE フィルタにより当該商品はマートから正しく除外される。
ポイント解説
updated_at は、アプリサーバー間のクロックずれ・リトライによる重複送信・タイムゾーン設定ミスなどで信頼できないことが多い。Datastreamが自動付与する _metadata_source_timestamp と、MySQLなら _metadata_log_file/_metadata_log_position(binlogの物理位置)を組み合わせることで、「ソースDBの中で実際にどちらが先に起きたか」を一意に決定できる。DELETE はDatastreamの世界では「行が消える」のではなく「_metadata_deleted=true の変更イベントが1行追加される」形で表現される。この設計を理解していないと、削除フィルタを書き忘れて「消えたはずのデータが下流で生き続ける」というバグを埋め込みやすい。incremental + パーティションフィルタによって、スキャン対象を「前回実行以降の差分」に固定化するのが定石。unique_key を指定した merge 戦略は「同じキーが来たら上書き」という冪等なUPSERTになるため、同じバッチを2回流しても結果が変わらないという実務上非常に重要な性質を得られる。total_latencies メトリクスへのアラートは、「データが古いまま止まっている」という状態を、ビジネス影響(誤った欠品アラート)が出る前に検知するための仕組み。reserved_quantity(引当済み数量)の意味づけや計算式は業務要件であり、将来ルールが増える可能性がある。CDCの順序解決・重複排除という「データ基盤としての正しさ」を担保するstaging層と、業務ロジックを表現するmart層を分離しておくことで、業務ルール変更時の影響範囲をmart層に限定できる。実務への応用
ECサイト MOps チームが扱う在庫データは、倉庫管理サービス・注文管理サービス・返品サービスなど複数のマイクロサービスにまたがって更新される。CDCでリアルタイムに近い連携をする設計は、バッチ集計に比べて鮮度は高いが、その分「変更履歴から正しい状態をどう組み立てるか」という設計の巧拙が直接ビジネス影響(誤った欠品通知・過剰発注・機会損失)に跳ね返る。特に本問のような「アラート・自動発注のトリガーとして使われるデータ」は、鮮度だけでなく正確性(重複排除・削除反映・順序保証)の3点セットが揃って初めて信頼できるデータ基盤になる。
on_schema_change やmodel contractのようなdbt側の防御に加えて、実務では「ソースDB側の変更管理プロセスにMOpsチームをレビュワーとして加える」といった組織的な連携も並行して整備するのが望ましい。
今日のまとめ
次のステップ
- 発展問題: 現在の設計は「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/ dbton_schema_change/ dbt source freshness / BigQuery パーティション分割テーブル・クラスタ化テーブル / Change Data Capture(CDC)設計パターン