データインジェスチョンは、業務システムやイベントログ、外部ファイルなど多様なソースから、分析システムへデータを取り込む処理を指す。ソースの状態をどの粒度で読み取り、ターゲットへどう書き込むかの組み合わせによって、インジェストの方式は複数のパターンに分かれる。
インジェスチョンは「ソースから何を取るか (抽出)」と「ターゲットへどう書くか (着地)」の 2 軸で整理できる。全件か差分か、追記か上書きか。この 2 択の組み合わせが、履歴の残り方・物理削除への追従・実行コストを左右する。
## 抽出と着地の 2 軸
インジェストパターンは次の 8 つ。
1. Full Refresh
2. Full Snapshot Append
3. Incremental Append
4. Upsert
5. Partition Overwrite
6. Rolling Window Append
7. Rolling Window Reload
8. CDC (Change Data Capture)
これらを抽出の粒度 (縦) と着地の書き込み方 (横) で整理すると、次のマトリクスになる。
| | 追記 (履歴を残す) | 上書き (最新のみ) |
| --------------------- | ---------------------------- | ------------------------------------ |
| **全件** | Full Snapshot Append | Full Refresh |
| **差分** | Incremental Append / CDC | Upsert / Partition Overwrite |
| **直近 N 日ウィンドウ** | Rolling Window Append | Rolling Window Reload |
## 共通の比較データ
8 つのパターンを同じ条件で比べるため、共通の注文データを用いる。1 日目 (1/6) と 2 日目 (1/7) の 2 時点でソースの変化を固定し、各パターンがこの変化をどうインジェストするかを見る。
- 下記 2 表は ETL インジェスト「前」のソースの真値。`load_ts` (インジェスト時に ETL で付与) はここには含めず、各インジェストパターンの結果表で初めて登場する。
- 各 Day の差分 (追加/更新/削除/変更なし) は「前回スナップショットとの比較」を表す。Day 1 は初回インジェストのため、基準は空集合。
### Day 1: ソース (1/6 時点)
| order_id (PK) | order_date | status | updated_at |
| ------------- | ---------- | ------- | ---------- |
| 101 | 2020-01-01 | paid | 2020-01-01 |
| 102 | 2020-01-02 | pending | 2020-01-02 |
| 103 | 2020-01-02 | paid | 2020-01-02 |
| 104 | 2020-01-04 | paid | 2020-01-04 |
| 105 | 2020-01-05 | pending | 2020-01-05 |
初回インジェストのため、全 5 行が追加。更新・削除・変更なしはない。
### Day 2: ソース (1/7 時点)
| order_id (PK) | order_date | status | updated_at |
| ------------- | ---------- | ------- | ---------- |
| 101 | 2020-01-01 | shipped | 2020-01-06 |
| 102 | 2020-01-02 | shipped | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 |
| 105 | 2020-01-05 | pending | 2020-01-05 |
| 106 | 2020-01-06 | pending | 2020-01-06 |
- 追加: `106`
- 更新: `101`, `102`
- 削除: `103`
- 変更なし: `104`, `105`
## 各パターンの詳細
各パターンとも、Day 1 は全 5 行を `load_ts` = 1/6 でインジェストした状態から始まる。以降の各節は Day 2 のインジェスト結果だけを示す (Day 1 で入った行は diff の変更なし行として残る)。
行データは共通だが、`load_ts` を PK にするか `change_type` 列を足すかはパターンごとに異なる。
```diff
Day 1 のインジェスト結果 (全パターン共通)
| order_id (PK) | order_date | status | updated_at | load_ts |
| ------------- | ---------- | ------- | ---------- | ---------- |
+ | 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
+ | 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
+ | 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
+ | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
+ | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
```
### 1. Full Refresh
毎回ソース全体を取得し、テーブルを全置換する。履歴は残らない。
```diff
Day 2 (全置換: 旧 5 行を捨て、ソースの最新 5 行を入れ直す)
| order_id (PK) | order_date | status | updated_at | load_ts |
| ------------- | ---------- | ------- | ---------- | ---------- |
- | 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
- | 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
- | 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
- | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
- | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 |
+ | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-07 |
+ | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-07 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
- **`load_ts`**: 全行が最新バッチ時刻に揃う。行ごとの新旧は見分けられない。
- **最新状態の取得**: テーブルが常に最新断面そのもの。そのまま読める。
- **削除の追従**: ソースで消えた行 (103) は全置換で自動的に消える。物理削除を確実に反映できる。
- **使いどころ**: 小規模・履歴不要・差分が取れないソース。
- **トレードオフ**: 実装の単純さと確実な削除反映を得る代わりに、毎回全件を I/O するコストを払い、過去の状態は一切残せない。
### 2. Full Snapshot Append
毎回ソース全体を取得し、消さずに追記する。各回の全断面が積層していく。
```diff
Day 2 (最新の全断面をまるごと追記)
| order_id (PK) | order_date | status | updated_at | load_ts (PK) |
| ------------- | ---------- | ------- | ---------- | ------------ |
| 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
| 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 |
+ | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-07 |
+ | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-07 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
- **`load_ts`**: 断面の識別キー (PK)。`WHERE load_ts = X` で過去の任意の断面を復元できる。
- **最新状態の取得**: `WHERE load_ts = MAX(load_ts)` で最新断面だけを取り出す。
- **削除の追従**: 物理削除 (103) は行としては消えないが、連続する 2 断面の差分で検知できる。
- **使いどころ**: 過去の各断面 (インジェスト粒度) をそのまま復元したい。
- **トレードオフ**: 過去のどの断面でも復元できる代わりに、変化のない行も毎回コピーし、ストレージが線形に増える (`load_ts` パーティション+失効で抑制)。
### 3. Incremental Append
前回以降に変化した行だけを、単調増加するカーソル列 (`updated_at`・バージョン番号・シーケンス ID など) で抽出して追記する。
```diff
Day 2 (updated_at > 2020-01-05 の差分 3 行だけを追記)
| order_id (PK) | order_date | status | updated_at | load_ts (PK) |
| ------------- | ---------- | ------- | ---------- | ------------ |
| 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
| 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
- **`load_ts`**: 変化した行だけが新しい `load_ts` で積まれる。更新ありは同一 `order_id` に複数バージョンが残り、イベントは 1 行 1 イベントで重複しない。
- **最新状態の取得**: 更新ありは `QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY load_ts DESC) = 1` で最新版を取る。イベントは重複がないため不要。
- **削除の追従**: 更新ありは物理削除 (103) を取りこぼす。イベントは削除自体が起きない。
- **使いどころ**: 大規模で全件取得が重く、信頼できるカーソルを持つソース。
- **トレードオフ**: 全件を読まず差分だけで済む代わりに、カーソルの正確さに依存し、更新ありでは物理削除を取りこぼす。
### 4. Upsert
差分を取得し、主キーで既存行に突き合わせる。一致すれば UPDATE、なければ INSERT する (MERGE)。Incremental Append と同じ差分を入力にするが、積まずに上書きするため各主キーは常に最新の 1 行だけを保つ。
```diff
Day 2 (差分を order_id で MERGE: 一致は UPDATE、新規は INSERT)
| order_id (PK) | order_date | status | updated_at | load_ts |
| ------------- | ---------- | ------- | ---------- | ---------- |
- | 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
- | 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
- **`load_ts`**: 行ごとに最後の更新時刻で上書きされる。履歴は残らず、各 `order_id` に最新 1 バージョンだけが存在する。
- **最新状態の取得**: テーブルが常に最新断面そのもの。絞り込みなしでそのまま読める。
- **削除の追従**: 物理削除 (103) を取りこぼす。消えた行は差分に来ないため、古いまま残り続ける。
- **使いどころ**: 最新状態だけ必要で、主キーで更新でき、差分が取れる大規模ソース。
- **トレードオフ**: 全件 I/O せず最新断面を保てる代わりに、履歴が残らず、物理削除を取りこぼす。
- **実装の補助**: ソースが更新時刻を持たない場合は、[[Python のデータフレームにハッシュ差分列を追加する|ハッシュ差分列]] で変更検知を代替できる。
### 5. Partition Overwrite
差分が属する `order_date` パーティションを特定し、そのパーティションをソースから全件取り直して洗い替える。行単位ではなくパーティション単位で置換するため、同じパーティションに同居する他の行も巻き込んで最新化される。
この方式は、dbt の増分モデルでは [[dbt の insert_overwrite はパーティション単位で洗い替えて増分更新する|insert_overwrite ストラテジー]] として実装できる。
```diff
Day 2 (変更があった order_date パーティション [1/1, 1/2, 1/6] を全件取り直して洗い替え)
| order_id (PK) | order_date | status | updated_at | load_ts |
| ------------- | ---------- | ------- | ---------- | ---------- |
- | 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
- | 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 |
- | 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
Upsert との結果差がここに出る。1/2 パーティションを取り直すので、同居する 103 の削除が反映され **103 が消える**。Upsert は行単位の突合なので 103 が残る。これが両者の本質的な違い。
- **`load_ts`**: 洗い替えたパーティション (1/1, 1/2, 1/6) の行だけ 1/7、触らないパーティション (1/4, 1/5) は 1/6 のまま。
- **最新状態の取得**: テーブルが常に最新断面そのもの。そのまま読める。
- **削除の追従**: 同一パーティション内の物理削除 (103) を拾える。1/2 パーティションを洗い替える際に巻き込んで消える。
- **使いどころ**: 遅延到着するファクト (ログ・課金・広告費など) の直近 N 日再処理、冪等なバックフィル、MERGE が高くつく大規模ファクト。
- **トレードオフ**: パーティション単位で削除も追従でき再処理も冪等な代わりに、パーティションキーが不変であることが前提で、取り直されないパーティションの削除やキー外の更新は取りこぼす。
### 6. Rolling Window Append
ソースの制約で直近 N 日分しか取得できないとき、その窓を毎回まるごと取り直して消さずに追記する。Full Snapshot Append を時間窓に限定した版。
下の例は lookback 3 日 (`order_date` 1/4〜1/6) の窓を再取得した場合で、初回 (1/6) はバックフィルで全件インジェスト済みとする。
```diff
Day 2 (lookback 3 日: order_date 1/4〜1/6 の窓を再取得して追記)
| order_id (PK) | order_date | status | updated_at | load_ts (PK) |
| ------------- | ---------- | ------- | ---------- | ------------ |
| 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
| 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-07 |
+ | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-07 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
- **`load_ts`**: 窓内 (1/4, 1/5, 1/6) の行を毎回追記する。104, 105 は値が同じでも別 `load_ts` で重複する (窓内の履歴を残すため)。
- **最新状態の取得**: `QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY load_ts DESC) = 1` で 1 行に畳む。
- **削除の追従**: 物理削除は窓内・窓外とも取りこぼす。追記なので消えた行 (103) は古いバージョンが残り続ける。
- **使いどころ**: 更新 API がなく過去 N 日だけ取得でき、窓内の履歴を残したいソース。
- **トレードオフ**: 差分カーソルがなくても直近の変化を取り込める代わりに、窓より古い更新 (101, 102) と削除 (103) を取りこぼし、窓を広げると取得・ストレージコストが増える。
### 7. Rolling Window Reload
Rolling Window Append と同じく直近 N 日窓を毎回取り直すが、窓を delete + insert で洗い替える (追記しない)。
この方式は Partition Overwrite を時間窓に固定した版であり、常に 1 行に畳まれる。窓・lookback の条件は Rolling Window Append と同じ。
```diff
Day 2 (lookback 3 日: order_date 1/4〜1/6 を delete + insert で洗い替え)
| order_id (PK) | order_date | status | updated_at | load_ts |
| ------------- | ---------- | ------- | ---------- | ---------- |
| 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 |
| 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 |
- | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 |
- | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 |
+ | 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-07 |
+ | 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-07 |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 |
```
Rolling Window Append との差は、窓を洗い替えるので旧 104, 105 が消えて重複しない点。窓外の 101, 102, 103 は同じく手付かず。
- **`load_ts`**: 窓内 (1/4, 1/5, 1/6) の行だけ 1/7、窓外は 1/6 のまま。最新 1 行だけを保つので単独 PK。
- **最新状態の取得**: テーブルが常に最新断面そのもの。そのまま読める。
- **削除の追従**: 窓内の削除は洗い替えで反映できる (Partition Overwrite と同じ)。窓外の削除 (103) は取りこぼす。
- **使いどころ**: 過去 N 日だけ取得でき、最新状態を 1 行で持ちたいソース。冪等なので再実行・遅延データに強い。
- **トレードオフ**: 窓内なら削除も追従でき冪等な代わりに、窓外の更新 (101, 102) と削除 (103) を取りこぼす。
### 8. CDC (Change Data Capture)
ソース DB のトランザクションログ (WAL / binlog / redo) を読み、すべての insert / update / delete を変更イベントとして捕捉する。クエリで現在値を取りに行く他の方式と違い、ログを読むため物理削除も `updated_at` を上げない更新も漏らさない。
CDC は厳密には「抽出」の方式である。出力は `change_type` 付きの追記専用な変更ログで、これを着地させれば Incremental Append の上位互換になる。可変データ (マスタ・可変トランザクション) が対象で、不変イベントには不要。
```diff
Day 2 (トランザクションログの変更を追記)
| order_id (PK) | order_date | status | updated_at | load_ts (PK) | change_type |
| ------------- | ---------- | ------- | ---------- | ------------ | ----------- |
| 101 | 2020-01-01 | paid | 2020-01-01 | 2020-01-06 | snapshot |
| 102 | 2020-01-02 | pending | 2020-01-02 | 2020-01-06 | snapshot |
| 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-06 | snapshot |
| 104 | 2020-01-04 | paid | 2020-01-04 | 2020-01-06 | snapshot |
| 105 | 2020-01-05 | pending | 2020-01-05 | 2020-01-06 | snapshot |
+ | 101 | 2020-01-01 | shipped | 2020-01-06 | 2020-01-07 | update |
+ | 102 | 2020-01-02 | shipped | 2020-01-06 | 2020-01-07 | update |
+ | 103 | 2020-01-02 | paid | 2020-01-02 | 2020-01-07 | delete |
+ | 106 | 2020-01-06 | pending | 2020-01-06 | 2020-01-07 | insert |
```
物理削除が `delete` イベントとして残るのが要点。他方式が取りこぼした 103 の削除を `change_type = delete` として捕捉し、その行はソースから消える直前の値 (before-image) を持つ。据え置きの 104, 105 はイベントを生まない。
- **`change_type`**: 各行の操作種別 (`insert` / `update` / `delete` / `snapshot`)。列名・値はツールで異なる (Debezium `op` = `c`/`u`/`d`/`r`、AWS DMS `Op` = `I`/`U`/`D`、Google Datastream・Delta Lake は `change_type`)。
- **最新状態の取得**: キーごとに最新の非 `delete` 行を取る。`delete` 行はそのキーが消えたことを表す。
- **削除の追従**: トランザクションログを読むので、物理削除も `updated_at` を上げない更新も漏らさない。クエリ系 (Incremental Append・Upsert・Partition Overwrite) が取りこぼす削除への回答になる。
- **使いどころ**: 可変データで削除や中間状態まで正確に追いたい。不変イベントには不要。
- **トレードオフ**: 削除も中間状態も漏れなく捕捉できる代わりに、DB のトランザクションログへのアクセス権が要り、コネクタ運用・順序保証・スキーマ変更対応の複雑さを負う。
> [!note]
> CDC の着地は、近年はマイクロバッチ (数分間隔) が主流。マネージドコネクタ (Google Datastream・Fivetran・Airbyte) が短い間隔で変更ログをまとめて流す形で、真のストリーミング (Debezium + Kafka など) は運用負荷が高く採用が絞られる。
>
> いずれも日次バッチより `load_ts` は変更時刻に近づく。ここでは他方式と比較するため日次バッチの枠で描いている。