C データエンジニアリング — dbt microbatch incremental(event_time + lookback=3)× BigQuery MERGE 冪等 UPSERT × 遅延到着データ対応 × dbt source freshness × Argo Workflows DAG(dependencies + retryStrategy + onExit)(MOps 注文集計マート Bad→Good)

2026-07-08 (Day 97) 水曜 C: データエンジニアリング ★★★★☆ dbt 1.9 microbatch / BigQuery MERGE late-arriving / source freshness / Argo DAG

概要

🕒

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.tasksdependenciesstg → 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+(microbatch incremental 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 スキャンでコスト過大
期待する回答形式: 問題点の列挙(番号付き)+ 改善後 dbt incremental モデル(microbatch)+ dbt config YAML + source freshness YAML + Argo Workflows DAG マニフェスト(YAML)+ 設計意図の説明

悪いコード (Before)

このモデル・マニフェストには 7つの設計上の問題 が隠れています。見つけてみてください。
bad_mart_daily_orders.sql — full-refresh・固定1日窓・unique_key なし
-- 問題②: 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 化しても重複が出る構造)
bad_sources.yml — freshness 未設定
version: 2
sources:
  - name: events
    schema: events
    tables:
      - name: order_events
        # 問題④: freshness / loaded_at_field なし
        #         → 上流の取り込み遅延を検知できない
bad_argo_cron.yaml — 直列1ステップ・retry なし・通知なし
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 通知なし)
問題点サマリー(7点)
1固定1日窓フィルタWHERE ordered_at >= CURRENT_DATE()-1 で 3 日遅れ到着を取りこぼす。microbatch + lookback: 3 で過去 3 日を再処理
2materialized: table(full-refresh) — 毎回 400TB スキャン。materialized: incremental + パーティション剪定で日次差分に限定
3unique_key 未設定 — 重複行が蓄積。microbatch は内部 MERGE、unique_key で冪等 UPSERT
4source freshness 未設定 — 上流遅延を検知不能。freshness: {warn_after, error_after} + loaded_at_field
51ステップ直列実行 — 依存不明瞭・部分失敗で全再実行。dag.tasksdependencies で明示
6retryStrategy 未設定 — 一時エラーで即死。limit: 3 + 指数バックオフ
7失敗通知なし — サイレント失敗。onExit exit handler で Slack 通知

ヒント(段階的開示)

ヒント1 — 方向性
日次集計の incremental モデルで「昨日分だけ」を 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.tasksdependencies で明示
  • 問題⑥: retryStrategy 未設定 → limit: 3 + backoff 指数バックオフ
  • 問題⑦: 失敗通知なし → onExit exit 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日を再処理
2materialized: table(400TB フルスキャン)コストincremental + event_time パーティション剪定
3unique_key 未設定(重複行)冪等性microbatch 内部 MERGE + unique_key
4source freshness 未設定可観測性freshness warn_after/error_after + loaded_at_field
51ステップ直列(依存不明瞭)オーケストレーションdag.tasks dependencies で DAG 明示
6retryStrategy 未設定信頼性limit: 3 + backoff 指数バックオフ
7失敗通知なし(サイレント失敗)可観測性onExit exit handler で Slack 通知

処理フロー図 — Bad vs Good の変換

Bad — full-refresh・固定窓・直列・通知なし events.order_events(3日 late-arriving あり) 問題④: source freshness なし → 上流遅延を検知できない partition: ordered_at / cluster: order_id mart_daily_orders(materialized: table) 問題②: 毎回 400TB フルスキャン($5/TB → $2,000/回) 問題①: WHERE ordered_at >= CURRENT_DATE()-1(固定窓) 問題③: unique_key なし → 重複行が蓄積 3日遅れ注文は集計から漏れる Argo CronWorkflow(run-all 1コンテナ) 問題⑤: dbt build を1ステップ直列実行 → 部分失敗しても全 stg/mart/freshness を再実行 問題⑥: retryStrategy なし → rateLimitExceeded で即死 問題⑦: onExit なし → 失敗してもサイレント × 3日遅れ注文が集計に反映されない × 毎回 400TB → オンデマンド課金過大 × 一時エラーで即失敗・気づかない × 部分失敗で全再実行(時間・コスト無駄) Good — microbatch・MERGE・DAG・通知 events.order_events + source freshness 修正④: warn_after 12h / error_after 72h(late-arriving 許容) loaded_at_field: ordered_at → 鮮度を日次監視 accepted_values / not_null / unique テスト mart_daily_orders(microbatch incremental) 修正①: event_time=ordered_at / batch_size=day / lookback=3 修正②: 直近4日パーティションのみスキャン(〜4TB/回) 修正③: 内部 MERGE + unique_key [order_date, channel] 3日遅れ注文も lookback で再処理 → 正確な GMV 手動 WHERE / is_incremental() 不要(自動注入) Argo DAG(dependencies + retry + onExit) 修正⑤: stg-order-events → mart-daily-orders → freshness-check 修正⑥: retryStrategy limit:3 / backoff duration 30s factor 2 修正⑦: onExit notify-slack({{workflow.status}} / failures) 失敗タスクのみリトライ → 全再実行を回避 concurrencyPolicy: Forbid / WIF で BigQuery 認証 retryPolicy OnError vs OnFailure の使い分け ✓ 3日遅れ注文も lookback で正確に集計 ✓ スキャン 400TB → 〜4TB(99%+ 削減) ✓ 冪等 UPSERT → 再実行安全(重複なし) ✓ 部分リトライ + Slack 通知で運用可観測性 修正

模範解答

Before — table・固定窓・unique_key なし
{{ 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
After — microbatch + パーティション剪定 + 冪等 UPSERT
-- 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: 33日遅れ到着を毎回再処理 → 正確な集計
② full-refreshincremental + 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 保護

ポイント解説

1microbatch incremental strategy(dbt 1.9 GA)の仕組み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 日分)を毎回再処理し、遅延到着データを自然に取り込む。
2lookback + unique_key で冪等性を担保microbatch は内部で BigQuery MERGE 相当の UPSERT を発行する。unique_key: [order_date, channel] を集計グレインに設定することで、同じバッチを再実行しても既存行が上書きされ重複が生じない。「昨日のバッチが失敗 → 今日リトライ」しても正しい結果になる。
3パーティション剪定によるスキャンコスト削減materialized: table(full-refresh)は毎回全期間(400TB)をスキャンするが、microbatch は直近 4 日分(当日 + lookback 3)のパーティションのみ読む。event_time 列がソースのパーティションキー(ordered_at)と一致すればパーティション剪定が効き、オンデマンドスキャンを 99% 以上削減できる。
4source freshness の warn_after / error_after 設計 — 遅延到着があるテーブルでは error_after を緩め(72h)にしないと正常な late-arriving を誤ってエラー扱いする。一方 warn_after(12h)で軽い警告を出し、上流パイプラインの異常な停滞を早期検知する。dbt source freshnessloaded_at_field の最大値と現在時刻の差分で鮮度を判定する。
5Argo DAG 依存と retryStrategy の役割分担dag.tasksdependenciesstg → mart → freshness の順序を宣言すると、Argo が依存グラフを解決して並列化・順序制御する。あるタスクが失敗しても retryStrategy はそのタスク単位でリトライされ、全ワークフローの再実行を避けられる。retryPolicy: "OnError" はシステムエラー(コンテナクラッシュ等)のみ、"OnFailure" はアプリ終了コード非ゼロもリトライする点に注意。
6onExit exit handler によるサイレント失敗の防止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 と組み合わせて運用の可観測性を高める。

今日のまとめ

遅延到着データを含む日次集計の7点チェックリスト: ① microbatch + lookback: 3(late-arriving 再処理)② incremental + event_time パーティション剪定(400TB → 〜4TB)③ unique_key で冪等 MERGE UPSERT ④ source freshness warn_after/error_after(error は緩めに)⑤ Argo dag.tasks dependencies(順序制御)⑥ retryStrategy limit + 指数バックオフ(部分リトライ)⑦ onExit exit handler で Slack 通知。 full-refresh を差分処理に変えてコストを削減し、DAG の依存・リトライ・通知でバッチ運用の信頼性と可観測性を高めることが実務のポイント。

次のステップ

  • 発展問題: 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

自己評価(あとで記入)