概要
WITH RECURSIVE + path サイクルガード
1次/2次/3次紹介を3本の自己JOINでハードコードすると段数が固定され、循環紹介(データ不整合)も検知できない。BigQueryの WITH RECURSIVE は基底ケースと再帰ケースを1つの定義にまとめ、path 配列で辿ったノードを保持して NOT IN UNNEST(path) により循環を安全に停止する。
ARRAY<STRUCT>は集約してからJOIN
UNNESTした注文明細行を紹介レベルと直接JOINすると、同じ注文が「マッチする紹介レベルの数」だけ再スキャン・再結合される(行爆発)。先に注文粒度へ GROUP BY 集約してからJOINすることで結合対象行数を最小化する。
dbt Mesh クロスプロジェクト参照
別dbtプロジェクトが管理するテーブルをハードコードされた完全修飾名で参照すると、dbtの依存グラフに載らず上流のスキーマ変更をサイレントに踏み抜く。{{ ref('shared_dw', 'dim_member') }} でlineageと環境切り替えを両立する。
exposures + tests + contract で防波堤
下流Looker Studioダッシュボードの依存を exposures で明示し、QUALIFY ROW_NUMBER() で紹介関係の重複を排除。bonus_amount のNULL/マイナス値・一意性は dbt tests + model contract でbuild時にガードする。
問題
ECサイト MOps チームでは「友達紹介キャンペーン」を運用しており、直接紹介(1次)だけでなく、紹介された会員がさらに別の会員を紹介した場合の2次・3次紹介まで遡って紹介ボーナスを支払う 多段階リファラルボーナス集計マート を BigQuery × dbt Core 1.9+ で構築している。以下の dbt モデルには 7つの設計・パフォーマンス・データガバナンス上の問題 が潜んでいる。
問題点を全て洗い出し、BigQuery WITH RECURSIVE 再帰CTE(深さ制限+サイクルガード)× ARRAY<STRUCT> 明細の事前集約による行爆発回避 × dbt Mesh クロスプロジェクト参照 × 紹介関係の重複排除(QUALIFY)× dbt exposures × incremental + partition_by/cluster_by × dbt tests / model contract を活用した Bad→Good リファクタリングを行ってください。
制約・前提条件
- BigQuery Standard SQL(GA機能のみ。
WITH RECURSIVE,QUALIFY,ARRAY_CONCAT,UNNEST利用可) - dbt Core 1.9+(Mesh のクロスプロジェクト
ref()、exposures、model contractenforcedが利用可) - テーブル
orders.order_events(order_id,member_id,order_date DATE,order_items ARRAY<STRUCT<product_id STRING, qty INT64, unit_price NUMERIC>>。partition:order_date) - テーブル
enrollment.referral_edges(member_id,referred_by_member_id,referred_at)。データ不整合で同一member_idに複数の紹介者登録が紛れ込むことがある dim_member(会員ランク等)は 別の dbt プロジェクトshared_dwが管理・公開している共有モデル- 紹介ボーナス率: 1次5% / 2次2% / 3次1%。3次より先は対象外
- このマートは Looker Studio の「紹介ボーナスダッシュボード」から参照される。BigQuery オンデマンド課金($5/TB)、注文は日々積み上がる
悪いコード (Before)
{{ config(materialized='table') }}
-- 問題①: 1次/2次/3次を3本の自己JOINでハードコード(段数固定・循環紹介を検知不可)
WITH level1 AS (
SELECT e1.member_id, e1.referred_by_member_id AS referrer_id, 1 AS level
FROM `my-project.enrollment.referral_edges` e1
-- 問題④: 重複紹介登録(同一member_idに複数referrer)を排除していない
),
level2 AS (
SELECT e2.member_id, e1.referrer_id, 2 AS level
FROM `my-project.enrollment.referral_edges` e2
JOIN level1 e1 ON e2.referred_by_member_id = e1.member_id
),
level3 AS (
SELECT e3.member_id, e2.referrer_id, 3 AS level
FROM `my-project.enrollment.referral_edges` e3
JOIN level2 e2 ON e3.referred_by_member_id = e2.member_id
),
all_levels AS (
SELECT * FROM level1
UNION ALL SELECT * FROM level2
UNION ALL SELECT * FROM level3
),
-- 問題②: UNNEST明細をレベルとJOINする前に集約していない → 行爆発
order_lines AS (
SELECT o.order_id, o.member_id, o.order_date, item.qty, item.unit_price
FROM `my-project.orders.order_events` o, UNNEST(o.order_items) AS item
),
-- 問題③: 別プロジェクトのdim_memberをハードコードFQNで参照(dbt ref()を使わない)
member_dim AS (
SELECT member_id, member_rank
FROM `project-a.shared_dw.dim_member`
)
SELECT
al.referrer_id,
al.level,
ol.order_id,
SUM(ol.qty * ol.unit_price) AS order_amount,
SUM(ol.qty * ol.unit_price)
* CASE al.level WHEN 1 THEN 0.05 WHEN 2 THEN 0.02 ELSE 0.01 END AS bonus_amount
FROM all_levels al
JOIN order_lines ol ON al.member_id = ol.member_id -- 明細行×レベル数だけ結合が膨張
JOIN member_dim md ON ol.member_id = md.member_id
GROUP BY al.referrer_id, al.level, ol.order_id
version: 2
models:
- name: bad_referral_bonus
# 問題⑤: Looker Studioダッシュボードがこのマートに依存しているが
# exposures が宣言されておらず影響範囲が誰にも見えない
# 問題⑦: bonus_amount のNULL/マイナス値・一意性の防御なし
columns:
- name: referrer_id
- name: bonus_amount
WITH RECURSIVEへヒント(段階的開示)
ヒント1 — 方向性
WITH RECURSIVE を使い、path 配列で「これまで辿ってきたノード」を保持して NOT IN UNNEST(path) でサイクルを弾く。(2) ネスト配列の集約順序 — ARRAY<STRUCT> の注文明細を UNNEST した直後に紹介レベルとJOINすると、同じ注文が「マッチする紹介レベルの数」だけ再スキャンされる。先に注文粒度へ集約してからJOINすれば結合対象行数を最小化できる。(3) 組織的なデータガバナンス — 別プロジェクトの共有モデルをハードコードFQNで参照する、下流ダッシュボードの依存が誰にも見えない、紹介関係の重複を検知しない、といった「動くが壊れやすい」状態を dbt Mesh・exposures・tests で防波堤化する。
ヒント2 — アプローチ
- 問題①: 1次/2次/3次を3本の自己JOINでハードコード(段数固定・循環紹介を検知不可)→
WITH RECURSIVE+level <= 3の深さ制限 +path配列によるサイクルガード - 問題②:
UNNESTした注文明細行を紹介レベルと直接JOINしてから集約(行爆発)→ 明細を先に注文粒度へGROUP BY集約したorder_totalsを作り、それを紹介チェーンにJOIN - 問題③: 別dbtプロジェクトの
dim_memberをハードコードされた完全修飾テーブル名で参照 → dbt Mesh のクロスプロジェクト{{ ref('shared_dw', 'dim_member') }} - 問題④:
referral_edgesに「1会員につき紹介者は1人」の保証がなく、重複登録で紹介ボーナスが多重計上されうる →QUALIFY ROW_NUMBER() OVER (PARTITION BY member_id ORDER BY referred_at DESC) = 1で最新1件のみ採用 - 問題⑤: Looker Studioダッシュボードがこのマートに依存していることがdbtプロジェクト上どこにも宣言されていない →
exposures.ymlでdepends_onを明示 - 問題⑥:
materialized: table(フルスキャン、partition/cluster なし)で注文増加に比例してビルドコストが増える →incremental+partition_by(order_date)+cluster_by(referrer_id) - 問題⑦:
bonus_amountのNULL/マイナス値防御なし、(referrer_id, order_id, level)の一意性保証なし → dbt tests(not_null/accepted_range/unique_combination_of_columns)+ model contract
ヒント3 — コードの骨格
-- 再帰CTEの骨格(BigQuery WITH RECURSIVE)
WITH RECURSIVE
referral_chain AS (
-- 基底: 1次紹介
SELECT member_id, referrer_id, 1 AS level,
[member_id, referrer_id] AS path -- 修正①: サイクル検出用の経路
FROM referral_edges_dedup
UNION ALL
-- 再帰: 紹介者のさらに上位を辿る(深さ制限 + サイクルガード)
SELECT rc.member_id, re.referrer_id, rc.level + 1,
ARRAY_CONCAT(rc.path, [re.referrer_id])
FROM referral_chain rc
JOIN referral_edges_dedup re ON rc.referrer_id = re.member_id
WHERE rc.level < 3 -- 修正①: 3次まで
AND re.referrer_id NOT IN UNNEST(rc.path) -- 修正①: 循環紹介ガード
)
-- 修正②: order_totals は「先に注文粒度へ集約してから」referral_chain とJOINする
-- 修正③: member_dim は {{ ref('shared_dw', 'dim_member') }} で参照する
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | 1次/2次/3次を3本の自己JOINでハードコード(段数固定・循環紹介を検知不可) | 拡張性/正確性 | WITH RECURSIVE + 深さ制限 + pathサイクルガード |
| 2 | UNNEST明細を集約前に紹介レベルとJOIN(行爆発) | パフォーマンス | 明細を先に注文粒度へ集約してからJOIN |
| 3 | 別プロジェクトのdim_memberをハードコードFQNで参照 | ガバナンス/lineage | dbt Mesh クロスプロジェクトref() |
| 4 | 紹介関係の重複登録を検知せず多重計上されうる | データ品質 | QUALIFY ROW_NUMBER()で最新1件のみ採用 |
| 5 | 下流Looker Studioダッシュボードの依存が宣言されていない | ガバナンス | dbt exposures |
| 6 | materialized: table(フルスキャン、partition/cluster なし) | コスト | incremental + partition_by/cluster_by |
| 7 | bonus_amountのNULL/マイナス値・一意性の防御なし | データ品質 | dbt tests + model contract |
再帰チェーンの仕組み — WITH RECURSIVE + path サイクルガード / 事前集約
模範解答
{{ config(materialized='table') }}
WITH level1 AS ( -- 問題①
SELECT e1.member_id, e1.referred_by_member_id AS referrer_id, 1 AS level
FROM `my-project.enrollment.referral_edges` e1
),
level2 AS (
SELECT e2.member_id, e1.referrer_id, 2 AS level
FROM `my-project.enrollment.referral_edges` e2
JOIN level1 e1 ON e2.referred_by_member_id = e1.member_id
),
level3 AS (
SELECT e3.member_id, e2.referrer_id, 3 AS level
FROM `my-project.enrollment.referral_edges` e3
JOIN level2 e2 ON e3.referred_by_member_id = e2.member_id
),
all_levels AS (
SELECT * FROM level1
UNION ALL SELECT * FROM level2
UNION ALL SELECT * FROM level3
),
order_lines AS ( -- 問題②
SELECT o.order_id, o.member_id, item.qty, item.unit_price
FROM `my-project.orders.order_events` o, UNNEST(o.order_items) AS item
),
member_dim AS ( -- 問題③
SELECT member_id, member_rank
FROM `project-a.shared_dw.dim_member`
)
SELECT
al.referrer_id, al.level, ol.order_id,
SUM(ol.qty * ol.unit_price) AS order_amount,
SUM(ol.qty * ol.unit_price)
* CASE al.level WHEN 1 THEN 0.05 WHEN 2 THEN 0.02 ELSE 0.01 END AS bonus_amount
FROM all_levels al
JOIN order_lines ol ON al.member_id = ol.member_id
JOIN member_dim md ON ol.member_id = md.member_id
GROUP BY al.referrer_id, al.level, ol.order_id
-- models/marts/mart_referral_bonus.sql
-- dbt Core 1.9+ | BigQuery Standard SQL
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite', -- 修正⑥
partition_by={
'field': 'order_date', 'data_type': 'date', 'granularity': 'day'
},
cluster_by=['referrer_id'], -- 修正⑥
on_schema_change='fail'
) }}
WITH RECURSIVE
-- 修正④: 1会員につき有効な紹介者は最新1件のみ採用
referral_edges_dedup AS (
SELECT member_id, referred_by_member_id AS referrer_id, referred_at
FROM {{ ref('stg_referral_edges') }}
WHERE referred_by_member_id IS NOT NULL
QUALIFY ROW_NUMBER() OVER (
PARTITION BY member_id ORDER BY referred_at DESC
) = 1
),
-- 修正①: ハードコード3段JOIN → 再帰CTEで多段紹介チェーンをたどる
referral_chain AS (
SELECT
member_id, referrer_id, 1 AS level,
[member_id, referrer_id] AS path -- サイクル検出用の経路
FROM referral_edges_dedup
UNION ALL
SELECT
rc.member_id, re.referrer_id, rc.level + 1,
ARRAY_CONCAT(rc.path, [re.referrer_id])
FROM referral_chain rc
JOIN referral_edges_dedup re ON rc.referrer_id = re.member_id
WHERE rc.level < 3 -- 3次紹介まで
AND re.referrer_id NOT IN UNNEST(rc.path) -- 循環紹介ガード
),
-- 修正②: 明細行は「先に注文粒度へ集約」してからJOIN(行爆発を回避)
order_totals AS (
SELECT
o.order_id, o.member_id, o.order_date,
SUM(item.qty * item.unit_price) AS order_amount
FROM {{ ref('stg_order_events') }} o, UNNEST(o.order_items) AS item
{% if is_incremental() %}
WHERE o.order_date >= DATE_SUB(CURRENT_DATE(), INTERVAL 3 DAY)
{% endif %}
GROUP BY o.order_id, o.member_id, o.order_date
),
-- 修正③: 別dbtプロジェクト(shared_dw)をdbt Meshのクロスプロジェクトref()で参照
member_dim AS (
SELECT member_id, member_rank
FROM {{ ref('shared_dw', 'dim_member') }}
)
SELECT
rc.referrer_id, rc.level, ot.order_id, ot.order_date, ot.order_amount,
ot.order_amount * CASE rc.level
WHEN 1 THEN 0.05 WHEN 2 THEN 0.02 ELSE 0.01
END AS bonus_amount
FROM referral_chain rc
JOIN order_totals ot ON rc.member_id = ot.member_id
JOIN member_dim md ON ot.member_id = md.member_id
# models/marts/_marts.yml
version: 2
models:
- name: mart_referral_bonus
description: "多段階(最大3次)紹介ボーナス集計マート。Looker Studio「紹介ボーナスダッシュボード」の参照元。"
config:
contract:
enforced: true # 修正⑦: スキーマ不一致・制約違反でbuildをfailさせる
columns:
- name: referrer_id
data_type: string
constraints:
- type: not_null
- name: level
data_type: int64
constraints:
- type: not_null
tests:
- accepted_values:
values: [1, 2, 3]
- name: order_id
data_type: string
constraints:
- type: not_null
- name: order_date
data_type: date
constraints:
- type: not_null
- name: order_amount
data_type: numeric
- name: bonus_amount
data_type: numeric
tests:
- dbt_utils.accepted_range: # 修正⑦: マイナス値の混入をbuild時に検知
min_value: 0
tests:
- dbt_utils.unique_combination_of_columns: # 修正⑦: 二重支払いの元となる重複行を検知
combination_of_columns:
- referrer_id
- order_id
- level
# 修正⑤: 下流ダッシュボードの依存を明示(スキーマ変更時の影響範囲を可視化)
exposures:
- name: referral_bonus_dashboard
label: "紹介ボーナスダッシュボード"
type: dashboard
maturity: medium
url: https://lookerstudio.google.com/reporting/referral-bonus-dashboard
description: "MOpsチームが友達紹介キャンペーンの支払い状況を確認するダッシュボード"
depends_on:
- ref('mart_referral_bonus')
owner:
name: MOps Team
email: mops-team@example.com
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① ハードコード3段JOIN | WITH RECURSIVE + 深さ制限 + pathガード | 段数変更が容易・循環紹介でも安全に停止 |
| ② UNNEST明細を集約前にJOIN | order_totalsで先に集約 | 結合対象行数を明細点数から独立させる |
| ③ 別プロジェクトをハードコードFQN | dbt Mesh ref('shared_dw', ...) | lineage維持・環境切り替え自動化 |
| ④ 紹介関係の重複未排除 | QUALIFY ROW_NUMBER()で最新1件 | ボーナス多重計上を防止 |
| ⑤ exposuresなし | exposures.ymlでdepends_on明示 | 下流ダッシュボードへの影響を可視化 |
| ⑥ materialized:table | incremental + partition/cluster | ビルドコストを差分処理に抑制 |
| ⑦ tests/contractなし | dbt tests + model contract | NULL/マイナス/重複をbuild時に検知 |
実行例(input → output)
input(referral_edges_dedup、簡略化):
member_id referrer_id
B A
C B
D C ← A(1次) ← B(2次) ← C(3次) の3段紹介チェーン
X Y
Y X ← データ不整合: X と Y が互いを紹介し合う循環データ
再帰CTE展開結果(referral_chain、D起点):
member_id referrer_id level path
D C 1 [D, C]
D B 2 [D, C, B]
D A 3 [D, C, B, A]
X/Y の循環データは、再帰項に到達した時点で re.referrer_id NOT IN UNNEST(rc.path) に引っかかり無限ループにならず1段で停止する。
output(mart_referral_bonus、Dの注文1万円分):
referrer_id level order_id order_amount bonus_amount
C 1 O-100 10000 500 (5%)
B 2 O-100 10000 200 (2%)
A 3 O-100 10000 100 (1%)
ポイント解説
WITH RECURSIVEへの置き換え — 「1次」「2次」「3次」を別々のCTEとしてコピペすると、段数を4次・5次に増やすたびにモデル修正が必要になる。BigQueryの再帰CTEは「基底ケース」(1次紹介)と「再帰ケース」(さらに上位を辿る)を UNION ALL で1つの再帰定義にまとめ、level の深さ制限だけで任意段数のロジックを表現できる。再帰項では集約関数・ORDER BY・LIMIT・外部結合の外側での自己参照が禁止されている点に注意。path 配列によるサイクルガード — 紹介関係はマスタデータであり、システム的な不整合(誤登録・手動修正ミス)で循環(A→B→A)が発生しうる。再帰CTEに深さ制限だけを設けても循環自体は誤った計算結果を生みうる。各行が「これまで辿ってきたノードの配列」を path として引き継ぎ、次の紹介者が既に path に含まれていたら停止する(NOT IN UNNEST(path))ことで、データ不整合があっても安全に停止する。ARRAY<STRUCT> の注文明細を UNNEST した直後の行を、紹介レベルという別粒度のテーブルと直接JOINすると、1つの注文が「マッチする紹介レベルの数」だけ重複してスキャン・結合される。明細行はまず注文単位に GROUP BY して order_amount に集約し、それを軽量な粒度(注文1件=1行)でJOINすることで結合対象行数と計算量を最小化できる。QUALIFY ROW_NUMBER() による重複紹介の防御 — 紹介関係は「1会員につき紹介者は1人」が業務ルールだが、データ入力ミスで同一 member_id に複数の紹介者登録が紛れ込むことがある。これを無防備に結合するとボーナスが多重計上される。QUALIFY ROW_NUMBER() OVER (PARTITION BY member_id ORDER BY referred_at DESC) = 1 で「最新の紹介登録のみ」を機械的に1件へ絞り込む。ref()でlineageを繋ぐ — ハードコードされた完全修飾名は、dbtの依存グラフ(DAG)に載らない。上流の shared_dw プロジェクトがスキーマを変更しても検知できず実行時に初めて壊れる。{{ ref('shared_dw', 'dim_member') }} はプロジェクト名を明示した参照で、dev/staging/prod の環境ごとに解決先を自動で切り替えつつ、依存関係をlineageグラフに正しく反映する。exposures にダッシュボードを宣言し depends_on でこのマートへの依存を明示することで、lineageグラフにBIツールまで含めた影響範囲が表示され、スキーマ変更前に「このカラムを消すとどのダッシュボードが壊れるか」を事前に把握できる。incremental + パーティション/クラスタリング + model contract + tests — materialized: table は注文が積み上がるほどビルドコストが線形に増加する。insert_overwrite で対象日のパーティションだけを洗い替える増分処理に変え、cluster_by(referrer_id) で紹介者単位の集計クエリを高速化する。加えて contract.enforced: true と bonus_amount への accepted_range、複合一意性テストで下流の会計システムに二重支払いや不正な値がサイレントに流れることをbuild時点で止める。実務への応用
ECサイト MOps チームの友達紹介キャンペーン運用に直結する。多段階紹介ボーナスのような「グラフ構造を辿る集計」は、段数が増えるたびにJOINを手書きで増やすアプローチが破綻しやすく、再帰CTE化しておくことでキャンペーン仕様変更(3次→4次への拡張等)に level の上限値を変えるだけで対応できる。循環データへの防御は、紹介キャンペーンに限らず「組織階層」「カテゴリツリー」「返信スレッド」など木構造・グラフ構造を扱うあらゆるマートで再利用できるパターン。
dbt Mesh のクロスプロジェクト参照と exposures は、チームが複数(MOps / 会員基盤 / BI)にまたがる組織でのデータガバナンスの要。shared_dw のような共通ディメンションを持つ組織では、ハードコードFQN参照が「誰にも気づかれないサイレント破壊」の温床になりやすく、早期にMesh化しておくことが技術的負債の予防になる。
今日のまとめ
次のステップ
- 発展問題: 紹介者が退会済み(
dim_memberに存在しない、またはis_active = false)の場合、その紹介者への以降のボーナス計上を止めつつ、さらに上位の紹介者への計上は継続する「孤立ノードのスキップ」ロジックを再帰CTEに組み込め。また、shared_dw側のdim_memberを dbt Mesh の model versions(v1/v2)で管理し、破壊的変更時に契約側へ段階移行させる方法を設計せよ。 - 参考: BigQuery
WITH RECURSIVE(再帰CTE)公式ドキュメント /ARRAY_CONCAT/UNNESTとネスト配列集約のベストプラクティス / dbt Mesh(cross-project ref, model versions)/ dbt exposures /dbt_utils.unique_combination_of_columns/QUALIFY