> ## Documentation Index
> Fetch the complete documentation index at: https://www.integrate.io/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# ETL: FlyData(ELT)とXplenty(ETL)を組み合わせた循環型データ統合アーキテクチャ

> FlyData(ELT)でデータを収集し、Xplenty(ETL)で加工結果を再び運用DBへ書き戻す循環型データパイプラインの構成方法を解説します。

## FlyDataについて

Integrate.ioのデータ転送サービスは、実は2つのプロダクトで構成されています。ここまで中心的に扱ってきたのはETLプロダクトのXplentyですが、それとは別にELTプロダクトのFlyDataも提供されています。

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-1.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=c8f50dd40fb8bc0f5d2bd24b7e8d24a1" alt="onpremise-part03-jp image 1" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-1.webp" />
</Frame>

FlyDataは、**MySQLなどのオペレーショナルDBからAmazon Redshiftのようなデータウェアハウスへ、ほぼリアルタイムでデータを同期する**ことに特化したクラウド型データ統合サービスです。CDC（Change Data Capture）の仕組みを使い、データベース側の変更内容を短い遅延でDWHへ反映できる点が大きな特徴です。

### FlyDataの特徴

* **CDCによるリアルタイム連携**: DB側の変更をほぼ即座に検知して転送
* **セットアップが軽量**: 最小限の設定でレプリケーションを開始できる
* **DWH連携に特化**: データウェアハウスへの取り込みに最適化されたアーキテクチャ

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-2.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=048f9246498073de59ddc06267247577" alt="onpremise-part03-jp image 2" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-2.webp" />
</Frame>

| 観点      | 強み                 | 制約                           |
| ------- | ------------------ | ---------------------------- |
| データ同期   | ほぼリアルタイムの同期が可能     | 対応先が主要なDWH・一部のDB・ファイルに限定される  |
| 機能面     | DWH向けに絞ることで高い効率を発揮 | ELTのため変換処理は持たず、異常値の検知程度にとどまる |
| アーキテクチャ | エージェントベースの軽量構成     | -                            |

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-3.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=1eaa40fb4e818b679f447da4d592319d" alt="onpremise-part03-jp image 3" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-3.webp" />
</Frame>

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-4.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=ec4605c7443c07c37f026b7b043896ad" alt="onpremise-part03-jp image 4" width="1201" height="1068" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-4.webp" />
</Frame>

## 概要：なぜ「循環型」なのか

近年のデータ活用では、データを集めるだけでなく、集めたデータを分析・加工し、その結果を再び業務システムに反映する**循環的なデータフロー**が求められるようになっています。

この記事では、次の2つのプロダクトを組み合わせてその循環を実現する構成を紹介します。

| 方向           | プロダクト        | 役割                              |
| ------------ | ------------ | ------------------------------- |
| Inbound（収集）  | FlyData(ELT) | 各運用DB → DWH（Snowflakeなど）へのデータ収集 |
| Outbound（配布） | Xplenty(ETL) | DWH → 各運用DBへの加工済みデータの書き戻し       |

この2つを組み合わせることで、**収集 → 中央での分析・計算 → 結果の書き戻し**という一連の循環パイプラインを構築できます。

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-5.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=60561446bce763fd312006e3f37fbdd1" alt="onpremise-part03-jp image 5" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-5.webp" />
</Frame>

## アーキテクチャ全体像

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-6.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=bf48174f8a753c46f58a9bb429b2a5e0" alt="onpremise-part03-jp image 6" width="1202" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-6.webp" />
</Frame>

各コンポーネントの役割は以下の通りです。

| コンポーネント              | 役割                        | 使用プロダクト                                             |
| -------------------- | ------------------------- | --------------------------------------------------- |
| ソース（既存の業務システム）       | 日常の業務データを蓄積               | MySQL、PostgreSQL、Oracleなどのデータベース、SaaSやオンプレのREST API |
| ELTレイヤー              | 変更のあったデータを素早くDWHへ取り込み     | FlyData                                             |
| データウェアハウス            | データの統合・保管、SQLベースの集計・調整処理  | Snowflake、BigQuery、Redshiftなど                       |
| ETLレイヤー              | 加工データの変換と各DB向けルーティング、書き戻し | Xplenty                                             |
| デスティネーション（既存の業務システム） | 調整結果の反映、分析データの蓄積          | MySQL、PostgreSQL、Oracleなどのデータベース、SaaSやオンプレのREST API |

## 各ステップの役割

* **ステップ1：FlyDataによる収集（ELT）**
  * CDCにより各運用DBの変更をリアルタイム〜準リアルタイムでSnowflakeに連携
  * 変換は行わず、まず生データをそのままDWHへ保存する（ELTの考え方の核）
  * 異なる種類の複数DBを1つのDWHに統合できる

* **ステップ2：Snowflake内での分析・計算**
  * 蓄積した生データをもとに、View・ストアドプロシージャ・dbtなどで**調整済みデータ**を作成
  * 例：在庫の調整数量、精算金額、集計済みのユーザースコアなど

* **ステップ3：Xplentyによる書き戻し（ETL）**
  * XplentyのビジュアルパイプラインでSnowflake上の加工データを読み込み
  * 必要な変換（フィールドマッピング、型変換、フィルタ）を適用
  * Merge（Upsert）、Append、Truncate & Insertなどの方式で各運用DBへ書き戻す

## ユースケース：ECプラットフォームの決済・在庫データ同期

### 背景

* 日本・米国・韓国それぞれに個別の運用DBを持つECサービス
* 各リージョンの決済データを中央で集計し、**グローバルな決済調整額**を算出したい
* 算出した調整額を各リージョンのDBへ反映する必要がある

### ステップ1：FlyDataで各リージョンDB → Snowflakeへ収集

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-7.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=5a8672dd9d4ab4622136189b5b04abc2" alt="onpremise-part03-jp image 7" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-7.webp" />
</Frame>

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-8.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=89576b3ac5b762d584adb4712e3a9a71" alt="onpremise-part03-jp image 8" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-8.webp" />
</Frame>

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-9.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=bdf501c96c70cda72e46aaf0bd89cc3f" alt="onpremise-part03-jp image 9" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-9.webp" />
</Frame>

* 日本のMySQL → CDC → Snowflake: `raw.jp_payments`
* 米国のPostgreSQL → CDC → Snowflake: `raw.us_payments`
* 韓国のMySQL → CDC → Snowflake: `raw.kr_payments`

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-10.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=4a15cc8b74c585c5a59ebde09c1aa09d" alt="onpremise-part03-jp image 10" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-10.webp" />
</Frame>

これらはSnowflake側でスキーマとして統合され、FlyData上のパイプライン設定によって継続的に更新されます。

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-11.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=d1d3d3266b3d5d66b2ee6d8894738f37" alt="onpremise-part03-jp image 11" width="1400" height="637" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-11.webp" />
</Frame>

### ステップ2：Snowflakeでグローバル調整額を算出

```sql theme={null}
-- 例：全体平均を基準にリージョンごとの調整額を算出
CREATE OR REPLACE VIEW analytics.adjusted_payments AS
SELECT
   payment_id,
   'JP' AS region,
   jp.amount - (global_avg.avg_amount * 0.3) AS adjusted_amount
...
...;
-- US、KRについても同様のロジックを適用
```

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-12.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=149d775c15d355a18fa6b2d173fe2d89" alt="onpremise-part03-jp image 12" width="477" height="278" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-12.webp" />
</Frame>

### ステップ3：Xplentyで各リージョンDBへ書き戻し

Xplentyパイプラインの構成例：

```
[Snowflake Source]
 analytics.adjusted_payments
       ↓
[Filterコンポーネント]
 region = 'JP' で絞り込み
       ↓
[Selectコンポーネント]
 payment_id, adjusted_qty のフィールドマッピング・型変換
       ↓
[Database Destination: JP MySQL]
 テーブル: payment_adjustment
 モード: Merge（Upsert）by payment_id
```

<Frame>
  <img src="https://mintcdn.com/integrateio/outYH6eYsuVJOWlW/images/japanese-knowledge-base/onpremise-part03-jp/image-13.webp?fit=max&auto=format&n=outYH6eYsuVJOWlW&q=85&s=266e250aaaaed0e3b5e5c6e177916d04" alt="onpremise-part03-jp image 13" width="1201" height="829" data-path="images/japanese-knowledge-base/onpremise-part03-jp/image-13.webp" />
</Frame>

* JP MySQL、US PostgreSQL、KR MySQLそれぞれに個別のパイプラインを用意するか、単一パイプライン内で分岐させる
* SnowflakeでのView更新完了後にXplentyパイプラインが自動実行されるよう、依存関係をスケジューラ側で設定する

## メリット・デメリット

### メリット

| メリット                      | 内容                                                       |
| ------------------------- | -------------------------------------------------------- |
| 役割分担が明確                   | 収集（FlyData）と配布（Xplenty）の責務が分かれ、それぞれの強みを活かせる              |
| Single Source of Truthの確立 | Snowflakeを中心に置くことで全リージョン・全サービスのデータが一箇所に集約され、一貫した基準で計算できる |
| ノーコードでのパイプライン管理           | Xplentyのドラッグ&ドロップUIにより、非エンジニアでも複雑なロジックを理解・修正しやすい         |
| 拡張しやすい                    | リージョンやサービスDBが増えても、FlyDataの接続とXplentyパイプラインを追加するだけで対応できる  |
| ELT区間のリアルタイム性             | FlyDataのCDCはバッチ方式より遅延が小さく、Snowflakeを常に最新状態に保てる           |
| 運用DBへの負荷軽減                | 重い集計処理はすべてSnowflake側で行われるため、運用DBへの負荷が抑えられる               |

### デメリット

| デメリット          | 内容                                                               |
| -------------- | ---------------------------------------------------------------- |
| 書き戻しの遅延        | パイプライン全体の遅延が積み重なるため、決済や在庫引当のようなリアルタイム性が必須な処理には不向き                |
| データループのリスク     | 書き戻したデータが再びFlyDataのCDCで検知され、Snowflakeへ再流入する可能性がある。これを防ぐフィルタが別途必要 |
| 運用が複雑になる       | FlyData・Snowflake・Xplentyの3つを同時に運用するため、監視・障害対応・コスト管理の負担が増える      |
| コストが2プロダクト分かかる | FlyDataとXplentyそれぞれのライセンス費用に加え、Snowflakeのコンピューティング費用も発生する        |
| スキーマ不整合のリスク    | ソースDB → Snowflake → 対象DBの間でカラム名・型・NULLの扱いが異なると、マッピングエラーが起きやすい    |

## 運用時の注意点

1. **データループの防止（最重要）**
   * 書き戻すデータには`source_system`や`is_adjusted`、`updated_by`のような出所を示すメタカラムを付与する
   * FlyDataのCDCフィルタ、またはSnowflake側のViewロジックで、書き戻したデータを再収集対象から除外する

2. **冪等性の担保**
   * Xplentyの書き戻しパイプラインが再実行されても重複挿入が起きないよう、Merge（Upsert）モードと明確な主キーを設定する

3. **実行順序の管理**
   * Snowflake側の計算（View更新、dbt runなど）が完了してからXplentyパイプラインが動くよう、スケジュールや依存関係トリガーを設定する
   * Xplentyの Dependent Execution や Webhook Trigger の活用が有効

4. **失敗時の部分ロールバック対策**
   * 書き戻しの途中で一部のDBだけ更新され残りが失敗すると、データ不整合が生じる
   * Xplentyの Single Transaction Mode や Pre/Post-action SQL を使ってロールバック戦略を用意する

5. **監視とアラート**
   * FlyData・Xplenty双方に障害通知を設定し、どの区間で失敗してもすぐ気づけるようにする
   * Slack、PagerDuty、Emailなどのフックが有効

6. **機密データの取り扱い**
   * 個人情報や金融情報を含むデータを書き戻す場合は、XplentyのSelectやFilterコンポーネントで不要な機密項目を除外・マスキングする

## まとめ

FlyData（ELT）とXplenty（ETL）を組み合わせた循環型のデータパイプラインは、**中央集約型のデータ管理**と**業務システムへのデータ反映**を同時に実現できる構成です。

複数リージョンや複数サービスの運用DBを抱えつつ、中央で計算した結果を各システムに反映したいEC・金融・マルチリージョンSaaSのようなケースで特に有効に機能します。

一方で、リアルタイム性や厳密な一貫性が求められるトランザクション処理には向かないため、要件に応じてアーキテクチャを使い分けることが重要です。
