概要
microbatch で遅延到着データを自然に取り込む
dbt 1.9 GA の microbatch incremental strategy は event_time を基準に batch_size(day/hour)ごとに時間窓を区切り、各バッチを独立処理する。lookback: 3 を設定すると最新バッチに加え過去 3 バッチも毎回再処理されるため、最大 3 日遅れで到着する注文レコード(late-arriving)を取りこぼさない。手動の is_incremental() や WHERE 句は不要で、dbt がバッチ範囲フィルタを自動注入する。
MERGE 冪等 UPSERT で再実行安全
microbatch は内部で BigQuery MERGE 相当の UPSERT を発行する。unique_key を集計グレイン(order_date + channel)に設定することで、同じバッチを再実行しても既存行が上書きされ重複が生じない。「昨日のバッチが失敗 → 今日リトライ」しても正しい結果になる。
パーティション剪定でスキャンコスト削減
materialized: table(full-refresh)は毎回全期間 400TB をスキャンするが、microbatch は直近 4 日分(当日 + lookback 3)のパーティションのみ読む。event_time がソースのパーティションキー(ordered_at)と一致すればパーティション剪定が効き、オンデマンドスキャンを 99% 以上削減できる。
Argo DAG で順序制御・部分リトライ・失敗通知
dag.tasks の dependencies で stg → mart → freshness を宣言し、失敗タスクのみ retryStrategy(指数バックオフ)でリトライ。onExit exit handler で成功/失敗を Slack 通知し、朝 05:00 バッチのサイレント失敗を防ぐ。dbt source freshness で上流の取り込み遅延も監視する。
問題
ECサイト MOps チームでは、注文イベント(events.order_events、日次到着)を BigQuery × dbt Core 1.9+ で日次集計マートに変換し、Argo Workflows CronWorkflow の DAG でオーケストレーションしている。以下の dbt モデル・スキーマ・Argo マニフェストには 7つの設計上の問題 が潜んでいる。問題点を全て洗い出し、dbt microbatch incremental strategy(dbt 1.9)× BigQuery MERGE 冪等 UPSERT × 遅延到着データ対応 × dbt source freshness × Argo Workflows DAG(dependencies + retryStrategy + exit handler) を活用した Bad→Good リファクタリングを行え。
制約・前提条件
- dbt Core 1.9+(
microbatchincremental strategy が GA、event_time/batch_size/lookback設定) - BigQuery Standard SQL(GA 版機能のみ)
- テーブル
events.order_events(partition:ordered_at TIMESTAMP, cluster:order_id)— ネットワーク遅延で最大 3 日遅れで到着するレコードあり(late-arriving) - Argo Workflows CronWorkflow で毎朝 JST 05:00 実行、
stg_order_events → mart_daily_orders → freshness_checkの順序依存あり - BigQuery オンデマンド課金($5/TB)。full-refresh は 1 回 400TB スキャンでコスト過大
悪いコード (Before)
-- 問題②: materialized table で毎回 400TB フルスキャン
{{ config(
materialized='table'
) }}
SELECT
DATE(ordered_at) AS order_date,
channel,
COUNT(DISTINCT order_id) AS order_count,
SUM(amount) AS gmv_jpy
FROM {{ source('events', 'order_events') }}
-- 問題①: 固定1日窓 → 3日遅れ到着レコードを取りこぼす
WHERE ordered_at >= CURRENT_DATE() - 1
GROUP BY order_date, channel
-- 問題③: incremental でないため unique_key なし
-- (後で incremental 化しても重複が出る構造)
version: 2
sources:
- name: events
schema: events
tables:
- name: order_events
# 問題④: freshness / loaded_at_field なし
# → 上流の取り込み遅延を検知できない
apiVersion: argoproj.io/v1alpha1
kind: CronWorkflow
metadata:
name: daily-orders-mart
namespace: mops
spec:
schedule: "0 20 * * *" # JST 05:00
workflowSpec:
entrypoint: run-all
templates:
- name: run-all
container:
image: .../mops/dbt-runner:latest
command: ["dbt"]
# 問題⑤: stg / mart / freshness を1コンテナ直列
# → 依存不明瞭・部分失敗で全再実行
args: ["build"]
# 問題⑥: retryStrategy なし
# → rateLimitExceeded で即ワークフロー失敗
# 問題⑦: onExit exit handler なし
# → 失敗してもサイレント(Slack 通知なし)
WHERE ordered_at >= CURRENT_DATE()-1 で 3 日遅れ到着を取りこぼす。microbatch + lookback: 3 で過去 3 日を再処理materialized: incremental + パーティション剪定で日次差分に限定microbatch は内部 MERGE、unique_key で冪等 UPSERTfreshness: {warn_after, error_after} + loaded_at_fielddag.tasks の dependencies で明示limit: 3 + 指数バックオフonExit exit handler で Slack 通知ヒント(段階的開示)
ヒント1 — 方向性
WHERE ordered_at >= CURRENT_DATE()-1 で処理すると、3 日遅れで到着したレコードを取りこぼす。dbt 1.9 の microbatch incremental strategy は event_time + lookback を設定すると複数のタイムバッチを自動的に再処理し、遅延到着データを自然に取り込める。また、DELETE+INSERT ではなく BigQuery MERGE による UPSERT で冪等性(同じバッチを再実行しても結果が変わらない)を担保する。Argo 側は 1 コンテナ直列を DAG に分割し、失敗タスクのみリトライ・失敗通知を付ける。
ヒント2 — アプローチ
- 問題①:
WHERE ordered_at >= CURRENT_DATE()-1(固定 1 日窓)→microbatch+lookback: 3で遅延到着を再処理 - 問題②:
materialized: table(毎回 400TB)→incremental+event_timeパーティション剪定 - 問題③:
unique_key未設定(重複行)→microbatchの内部 MERGE、unique_keyで冪等 - 問題④: source freshness 未設定 →
freshness: {warn_after, error_after}+loaded_at_field - 問題⑤: 全ステップ 1 コンテナ直列 →
dag.tasksのdependenciesで明示 - 問題⑥:
retryStrategy未設定 →limit: 3+backoff指数バックオフ - 問題⑦: 失敗通知なし →
onExitexit handler で Slack 通知({{workflow.status}})
ヒント3 — コードの骨格
-- dbt microbatch incremental モデルの骨格
{{ config(
materialized='incremental',
incremental_strategy='microbatch',
event_time='ordered_at',
batch_size='day',
lookback=3, -- 修正①: 過去3日を毎回再処理(遅延到着対応)
begin='2026-01-01',
unique_key=['order_date','channel'], -- 修正③: 冪等 UPSERT
partition_by={'field':'order_date','data_type':'date','granularity':'day'},
cluster_by=['channel']
) }}
-- microbatch では dbt が ordered_at のバッチ範囲フィルタを自動注入
-- → 手動 WHERE / is_incremental() ブロック不要
# Argo Workflows DAG の骨格
dag:
tasks:
- name: stg-order-events
template: dbt-run
- name: mart-daily-orders
dependencies: [stg-order-events] # 修正⑤
- name: freshness-check
dependencies: [mart-daily-orders]
# retryStrategy: {limit: 3, backoff: {duration: 30s, factor: 2}} # 修正⑥
# onExit: notify-slack # 修正⑦
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | 固定1日窓フィルタ(遅延到着取りこぼし) | 正確性 | microbatch + lookback: 3 で過去3日を再処理 |
| 2 | materialized: table(400TB フルスキャン) | コスト | incremental + event_time パーティション剪定 |
| 3 | unique_key 未設定(重複行) | 冪等性 | microbatch 内部 MERGE + unique_key |
| 4 | source freshness 未設定 | 可観測性 | freshness warn_after/error_after + loaded_at_field |
| 5 | 1ステップ直列(依存不明瞭) | オーケストレーション | dag.tasks dependencies で DAG 明示 |
| 6 | retryStrategy 未設定 | 信頼性 | limit: 3 + backoff 指数バックオフ |
| 7 | 失敗通知なし(サイレント失敗) | 可観測性 | onExit exit handler で Slack 通知 |
処理フロー図 — Bad vs Good の変換
模範解答
{{ config(materialized='table') }}
SELECT
DATE(ordered_at) AS order_date,
channel,
COUNT(DISTINCT order_id) AS order_count,
SUM(amount) AS gmv_jpy
FROM {{ source('events','order_events') }}
-- 固定1日窓 → 3日遅れを取りこぼす
WHERE ordered_at >= CURRENT_DATE() - 1
GROUP BY order_date, channel
-- models/marts/mart_daily_orders.sql
-- dbt Core 1.9+ | 遅延到着データ対応 microbatch
{{ config(
materialized='incremental',
incremental_strategy='microbatch', -- 修正②③: 差分 + 内部 MERGE
event_time='ordered_at', -- バッチ分割の基準列
batch_size='day',
lookback=3, -- 修正①: 過去3日を再処理
begin='2026-01-01',
unique_key=['order_date','channel'], -- 修正③: 冪等 UPSERT のキー
partition_by={
'field': 'order_date',
'data_type': 'date',
'granularity': 'day'
},
cluster_by=['channel']
) }}
-- microbatch は dbt が ordered_at のバッチ範囲を自動注入する:
-- WHERE ordered_at >= '{start}' AND ordered_at < '{end}'
-- → 手動 WHERE / is_incremental() ブロックは不要
SELECT
DATE(ordered_at) AS order_date, -- パーティションキー
channel, -- web / app / store
COUNT(DISTINCT order_id) AS order_count, -- 重複排除した注文数
SUM(amount_jpy) AS gmv_jpy, -- 流通総額
AVG(amount_jpy) AS avg_order_value_jpy
FROM {{ ref('stg_order_events') }}
GROUP BY order_date, channel
-- ── staging(型変換・NULL ガード、microbatch 連携)──
-- models/staging/stg_order_events.sql
-- {{ config(materialized='incremental',
-- incremental_strategy='microbatch', event_time='ordered_at',
-- batch_size='day', lookback=3, begin='2026-01-01',
-- unique_key='order_id') }}
-- SELECT order_id,
-- SAFE_CAST(amount AS NUMERIC) AS amount_jpy, -- 不正値は NULL
-- LOWER(TRIM(channel)) AS channel,
-- ordered_at
-- FROM {{ source('events','order_events') }}
-- WHERE order_id IS NOT NULL -- 破損レコードの除外のみ明示
# models/staging/_sources.yml
# 修正④: source freshness で上流の取り込み遅延を検知
version: 2
sources:
- name: events
database: my-gcp-project
schema: events
# ソース全体のデフォルト鮮度基準
freshness:
warn_after: {count: 6, period: hour} # 6時間遅延で警告
error_after: {count: 25, period: hour} # 25時間でエラー(日次跨ぎ許容)
loaded_at_field: ordered_at # 鮮度判定に使う列
tables:
- name: order_events
description: "注文確定イベント(最大3日 late-arriving あり)"
# 遅延到着があるため error_after を緩めに上書き
freshness:
warn_after: {count: 12, period: hour}
error_after: {count: 72, period: hour} # 3日 late-arriving を許容
columns:
- name: order_id
tests: [not_null, unique]
- name: channel
tests:
- accepted_values:
values: ["web", "app", "store"]
error_after を緩め(72h)にしないと正常な late-arriving を誤ってエラー扱いする。warn_after(12h)で軽い警告を出し、異常な停滞のみ早期検知する。
# argo-cron-daily-orders.yaml
apiVersion: argoproj.io/v1alpha1
kind: CronWorkflow
metadata:
name: daily-orders-mart
namespace: mops
spec:
schedule: "0 20 * * *" # UTC 20:00 = JST 05:00
timezone: "Asia/Tokyo"
concurrencyPolicy: Forbid # 前回未完了なら新規起動しない
startingDeadlineSeconds: 300
workflowSpec:
entrypoint: main
serviceAccountName: dbt-runner # WIF で BigQuery 認証
onExit: notify-slack # 修正⑦: 終了時に必ず実行
arguments:
parameters:
- name: dbt-target
value: "prod"
templates:
# ── DAG 定義(修正⑤: 依存を明示)──
- name: main
dag:
tasks:
- name: stg-order-events
template: dbt-run
arguments:
parameters: [{name: select, value: "stg_order_events"}]
- name: mart-daily-orders
template: dbt-run
dependencies: [stg-order-events] # staging 完了後
arguments:
parameters: [{name: select, value: "mart_daily_orders"}]
- name: freshness-check
template: dbt-source-freshness
dependencies: [mart-daily-orders] # マート完了後
# ── dbt run(修正⑥: retryStrategy)──
- name: dbt-run
inputs:
parameters: [{name: select}]
retryStrategy:
limit: 3
retryPolicy: "OnError" # システムエラーのみリトライ
backoff:
duration: "30s"
factor: 2 # 30s → 60s → 120s
maxDuration: "5m"
container:
image: .../mops/dbt-runner@sha256:abc123
command: ["dbt"]
args: ["run", "--select", "{{inputs.parameters.select}}",
"--target", "{{workflow.parameters.dbt-target}}"]
- name: dbt-source-freshness
retryStrategy:
limit: 2
backoff: {duration: "20s", factor: 2, maxDuration: "2m"}
container:
image: .../mops/dbt-runner@sha256:abc123
command: ["dbt"]
args: ["source", "freshness",
"--target", "{{workflow.parameters.dbt-target}}"]
# ── 修正⑦: Exit handler(Slack 通知)──
- name: notify-slack
container:
image: .../mops/slack-notifier@sha256:def456
command: ["python", "/app/notify.py"]
env:
- name: WORKFLOW_STATUS
value: "{{workflow.status}}" # Succeeded / Failed / Error
- name: WORKFLOW_NAME
value: "{{workflow.name}}"
- name: FAILURES
value: "{{workflow.failures}}"
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① 固定1日窓 | microbatch + lookback: 3 | 3日遅れ到着を毎回再処理 → 正確な集計 |
| ② full-refresh | incremental + event_time パーティション剪定 | 400TB → 〜4TB(99%+ スキャン削減) |
| ③ unique_key なし | 内部 MERGE + unique_key [order_date, channel] | 再実行しても重複しない冪等 UPSERT |
| ④ freshness なし | warn_after/error_after + loaded_at_field | 上流遅延の早期検知(late-arriving 許容) |
| ⑤ 1ステップ直列 | dag.tasks dependencies | 順序制御 + 失敗タスクのみ再実行 |
| ⑥ retry なし | retryStrategy limit: 3 + 指数バックオフ | rateLimitExceeded 等の一時エラーに耐性 |
| ⑦ 通知なし | onExit exit handler で Slack 通知 | サイレント失敗を防止 → SLA 保護 |
ポイント解説
event_time で指定した列を基準に batch_size(day/hour/month)ごとに時間窓を区切り、各バッチを独立処理する。dbt が自動で WHERE {event_time} >= {batch_start} AND {event_time} < {batch_end} を注入するため、モデル SQL に手動の is_incremental() ブロックや WHERE 句は不要。lookback: 3 で最新バッチ + 過去 3 バッチ(3 日分)を毎回再処理し、遅延到着データを自然に取り込む。microbatch は内部で BigQuery MERGE 相当の UPSERT を発行する。unique_key: [order_date, channel] を集計グレインに設定することで、同じバッチを再実行しても既存行が上書きされ重複が生じない。「昨日のバッチが失敗 → 今日リトライ」しても正しい結果になる。materialized: table(full-refresh)は毎回全期間(400TB)をスキャンするが、microbatch は直近 4 日分(当日 + lookback 3)のパーティションのみ読む。event_time 列がソースのパーティションキー(ordered_at)と一致すればパーティション剪定が効き、オンデマンドスキャンを 99% 以上削減できる。error_after を緩め(72h)にしないと正常な late-arriving を誤ってエラー扱いする。一方 warn_after(12h)で軽い警告を出し、上流パイプラインの異常な停滞を早期検知する。dbt source freshness は loaded_at_field の最大値と現在時刻の差分で鮮度を判定する。dag.tasks の dependencies で stg → mart → freshness の順序を宣言すると、Argo が依存グラフを解決して並列化・順序制御する。あるタスクが失敗しても retryStrategy はそのタスク単位でリトライされ、全ワークフローの再実行を避けられる。retryPolicy: "OnError" はシステムエラー(コンテナクラッシュ等)のみ、"OnFailure" はアプリ終了コード非ゼロもリトライする点に注意。onExit: notify-slack はワークフローの成功・失敗に関わらず終了時に必ず実行される。{{workflow.status}}(Succeeded/Failed/Error)と {{workflow.failures}}(失敗タスクの JSON)を環境変数で渡すことで、通知スクリプトが失敗内容を Slack に投稿できる。朝 05:00 のバッチ失敗を即検知し SLA 違反を防ぐ。実務への応用
ECサイト MOps チームの注文集計マート × キャンペーン ROAS 分析では本パターンが直接効く。決済確定が遅延した注文(3 日遅れ)も lookback: 3 で過去バッチが再処理されるため、月次 GMV レポートが後から正しい値に自動修正される。手動のバックフィルスクリプトが不要になる。
microbatch により full-refresh(400TB/回)を差分処理(〜4TB/回)に変え、オンデマンド課金を月あたり大幅削減できる。dbt run --select mart_daily_orders --event-time-start 2026-06-01 --event-time-end 2026-07-01 で任意期間の手動バックフィルも可能。
Argo DAG の部分リトライにより、mart-daily-orders タスクだけが rateLimitExceeded で失敗した場合でも stg-order-events は再実行されず、失敗タスクのみ指数バックオフでリトライされる。BigQuery スロット逼迫時の一時エラーに強くなる。
dbt source freshness の結果を DataDog / Slack に送り、上流の注文イベント取り込みが 12h 停滞したら P3 アラート、72h でエラー通知。exit handler と組み合わせて運用の可観測性を高める。
今日のまとめ
次のステップ
- 発展問題:
microbatchモデルの並列バッチ実行(Argo Workflows のwithItemsで複数バッチを並列化)を設計し、過去 90 日のバックフィルを 1 時間以内に完了させる DAG を実装せよ。BigQuery スロット予約(Editions)との兼ね合いも考慮すること。 - 参考: dbt Incremental microbatch 公式ドキュメント(1.9+)/ BigQuery MERGE 文リファレンス / Argo Workflows DAG・retryStrategy・exit handler / dbt source freshness