オープンソースETL・データパイプライン比較:dbt vs Apache Airflow vs Meltano でデータ変換をセルフホストする
オープンソースラボ編集部 ・ 2026年6月13日
オープンソースETL・データパイプライン比較:dbt vs Apache Airflow vs Meltano でデータ変換をセルフホストする
FivetranやStitchなどの有料ETLサービスに依存せず、データパイプライン・データ変換・スケジューリングをオープンソースで構築しましょう。dbt・Apache Airflow・Meltanoはデータウェアハウスへの取り込みから変換・スケジューリングまでの全工程をコードで管理できます。
ETL・データパイプラインが必要な場面
- データ統合: 複数のSaaS(Salesforce・HubSpot・Stripe)のデータをBigQuery・Redshiftに集約
- データ変換: 生データをアナリストが使いやすいデータモデルに変換
- スケジューリング: 夜間バッチ処理・リアルタイムデータ更新を自動化
- コスト削減: Fivetranは$0.01/MAR(月間アクティブ行)〜、大量データで高額になる
- データリネージ: データがどこから来てどう変換されたかの完全な追跡
主要ツールの概要
dbt(data build tool)
データウェアハウス内でのSQL変換に特化したツールです。SQLファイルをモデルとして管理し、依存関係を自動解決して変換を実行します。テスト・ドキュメント・バージョン管理(Git)を組み合わせることで、データエンジニアリングにソフトウェア開発のベストプラクティスを適用できます。
# dbt-coreをインストール(BigQuery用)
pip install dbt-bigquery
# または他のウェアハウス
# pip install dbt-snowflake dbt-redshift dbt-postgres
# 新しいdbtプロジェクトを作成
dbt init my_analytics
cd my_analytics
# profiles.yml(~/.dbt/profiles.yml)の設定例(BigQuery)
# my_analytics:
# target: dev
# outputs:
# dev:
# type: bigquery
# method: oauth
# project: my-gcp-project
# dataset: dbt_dev
# threads: 4
# timeout_seconds: 300
-- models/staging/stg_orders.sql
-- 生データを正規化するステージングモデル
WITH source AS (
SELECT * FROM {{ source('raw', 'orders') }}
),
renamed AS (
SELECT
order_id,
customer_id,
CAST(order_date AS DATE) AS order_date,
LOWER(status) AS status,
total_amount / 100.0 AS total_amount_jpy,
_loaded_at AS created_at
FROM source
WHERE order_id IS NOT NULL
)
SELECT * FROM renamed
-- models/marts/fct_daily_revenue.sql
-- ビジネスロジックを含むファクトテーブル
{{ config(
materialized='table',
partition_by={
"field": "order_date",
"data_type": "date",
"granularity": "day",
},
cluster_by=["customer_segment"]
) }}
WITH orders AS (
SELECT * FROM {{ ref('stg_orders') }}
),
customers AS (
SELECT * FROM {{ ref('stg_customers') }}
),
final AS (
SELECT
o.order_date,
c.customer_segment,
COUNT(DISTINCT o.order_id) AS order_count,
SUM(o.total_amount_jpy) AS total_revenue_jpy,
AVG(o.total_amount_jpy) AS avg_order_value_jpy
FROM orders o
LEFT JOIN customers c USING (customer_id)
WHERE o.status = 'completed'
GROUP BY 1, 2
)
SELECT * FROM final
# models/staging/sources.yml
version: 2
sources:
- name: raw
database: my-gcp-project
schema: raw_data
tables:
- name: orders
description: "Stripeから同期された注文データ"
freshness:
warn_after: {count: 12, period: hour}
error_after: {count: 24, period: hour}
loaded_at_field: _loaded_at
columns:
- name: order_id
tests:
- unique
- not_null
- name: status
tests:
- accepted_values:
values: ['completed', 'pending', 'cancelled', 'refunded']
# dbtコマンドの実行
dbt run # 全モデルを実行
dbt test # データテストを実行
dbt docs generate # ドキュメントを生成
dbt docs serve # ブラウザでドキュメントを表示
# 特定モデルとその依存関係を実行
dbt run --select +fct_daily_revenue # fct_daily_revenueと全上流モデル
dbt run --select fct_daily_revenue+ # fct_daily_revenueと全下流モデル
Apache Airflow
Pythonでワークフロー(DAG)を定義し、スケジューリング・依存関係管理・監視を行うオーケストレーターです。ETL・MLパイプライン・データ処理の自動化に広く使われています。
# AirflowをDockerで起動(公式compose)
curl -LfO 'https://airflow.apache.org/docs/apache-airflow/stable/docker-compose.yaml'
mkdir -p ./dags ./logs ./plugins ./config
echo -e "AIRFLOW_UID=$(id -u)" > .env
docker compose up airflow-init
docker compose up -d
# UIは http://localhost:8080 でアクセス(airflow/airflow)
# dags/etl_orders_pipeline.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
from airflow.providers.http.operators.http import SimpleHttpOperator
from datetime import datetime, timedelta
import httpx
import json
default_args = {
"owner": "data-team",
"depends_on_past": False,
"start_date": datetime(2026, 1, 1),
"email_on_failure": True,
"email": ["data-team@example.com"],
"retries": 3,
"retry_delay": timedelta(minutes=5),
}
def extract_stripe_orders(**context):
api_key = context["var"]["value"]["STRIPE_API_KEY"]
execution_date = context["ds"] # 実行日付(YYYY-MM-DD)
all_charges = []
has_more = True
starting_after = None
while has_more:
params = {"limit": 100, "created[gte]": execution_date}
if starting_after:
params["starting_after"] = starting_after
response = httpx.get(
"https://api.stripe.com/v1/charges",
headers={"Authorization": f"Bearer {api_key}"},
params=params,
)
data = response.json()
all_charges.extend(data["data"])
has_more = data["has_more"]
if has_more:
starting_after = data["data"][-1]["id"]
context["ti"].xcom_push(key="charges", value=all_charges)
return len(all_charges)
with DAG(
"etl_orders_pipeline",
default_args=default_args,
description="StripeからBigQueryへの毎日の注文データETL",
schedule_interval="0 2 * * *", # 毎日2時に実行
catchup=False,
tags=["etl", "stripe", "bigquery"],
) as dag:
extract = PythonOperator(
task_id="extract_stripe_orders",
python_callable=extract_stripe_orders,
)
load_to_bq = BigQueryInsertJobOperator(
task_id="load_to_bigquery",
configuration={
"query": {
"query": "SELECT * FROM UNNEST(@charges)",
"useLegacySql": False,
}
},
gcp_conn_id="google_cloud_default",
)
run_dbt = PythonOperator(
task_id="run_dbt_models",
python_callable=lambda: __import__("subprocess").run(
["dbt", "run", "--select", "+fct_daily_revenue"], check=True
),
)
extract >> load_to_bq >> run_dbt
Meltano
Singer仕様のコネクター(Tap/Target)とdbt・Airflowを統合したELTプラットフォームです。meltano.ymlで全パイプラインをコードとして管理し、GitOpsなデータエンジニアリングを実現します。
# Meltanoをインストールしてプロジェクトを初期化
pip install meltano
meltano init my-data-pipeline
cd my-data-pipeline
# Salesforce Tapを追加
meltano add extractor tap-salesforce
meltano config tap-salesforce set username "user@example.com"
meltano config tap-salesforce set password "password"
meltano config tap-salesforce set security_token "token"
# BigQuery Targetを追加
meltano add loader target-bigquery
meltano config target-bigquery set project_id "my-gcp-project"
meltano config target-bigquery set dataset_id "raw_salesforce"
# dbtを追加して変換パイプラインを統合
meltano add transformer dbt-bigquery
meltano add orchestrator airflow
# 全パイプラインを実行
meltano run tap-salesforce target-bigquery dbt-bigquery:run
機能比較表
| 比較項目 | dbt Core | Apache Airflow | Meltano |
|---|---|---|---|
| 主な用途 | データ変換(T) | ワークフローオーケストレーション | ELT統合プラットフォーム |
| ETLのE(Extract) | ❌ | ✅ | ✅ Singerコネクター |
| ETLのT(Transform) | ✅ SQL | ✅ Python | ✅ dbt統合 |
| ETLのL(Load) | ❌ | ✅ | ✅ Singerコネクター |
| スケジューリング | ❌(外部ツール要) | ✅ 高機能 | ✅(Airflow統合) |
| UIダッシュボード | ✅ ドキュメント | ✅ 監視UI | ✅ |
| データテスト | ✅ 組み込み | ❌(別途) | ✅ |
| データリネージ | ✅ | ✅ | ✅ |
| SQLウェアハウス対応 | BigQuery/Snowflake/Redshift/DuckDB | 全DB | 全DB |
| 学習コスト | ★★★☆☆ | ★★☆☆☆ | ★★★☆☆ |
| 実装言語 | Python(CLIはGo) | Python | Python |
| ライセンス | Apache 2.0 | Apache 2.0 | MIT |
| GitHub Stars | 10k+ | 37k+ | 2k+ |
データ変換・分析ツールはknowledgeカテゴリ(/categories/knowledge)で一覧でき、データ可視化・BIツールはDevOpsカテゴリ(/categories/devops)でも探せます。
FAQ
Q. dbt・Airflow・Meltanoはどう組み合わせて使いますか?
A. 典型的な構成: Meltano(またはAirbyte)でSaaS→ウェアハウスへのデータ取り込み(E・L)→dbtでウェアハウス内のデータ変換(T)→Airflowまたはdbtのスケジューラーでパイプラインをオーケストレーション。Meltanoは内部でdbtとAirflowを統合できるため、3ツールすべてをmeltano.yml一つで管理することも可能です。小規模チームはdbt+GitHubActions(スケジュール)で始め、複雑なパイプラインが増えたらAirflowを追加するのが現実的です。
Q. Fivetranの代替としてMeltanoを使う際の注意点は?
A. MeltanoはSinger仕様の350以上のコネクター(Tap)を使えますが、Fivetranのコネクターほどメンテナンスが行き届いていないものがあります。特に: 主要なコネクター(Salesforce・HubSpot・Stripe・Google Analytics)は比較的安定しているが、マイナーなSaaSは動作しない場合がある。本番移行前に必ずテスト環境で動作確認を。Airbyte(別のオープンソースELTツール)はコネクターのメンテナンスがより充実しており、Fivetran代替としてはAirbyteとMeltanoを比較検討することを推奨します。
Q. dbtのモデルを本番環境で安全に実行するには?
A. 本番環境でのdbtのベストプラクティス: 1) Slim CI(変更されたモデルのみを実行--select state:modified+)でCIコストを削減、2) --target prodで本番用のBigQuery/Snowflakeデータセットに向ける、3) dbt testをdbt runの前後に実行して異常値・null・重複を検出、4) dbt source freshnessでソースデータの鮮度を定期チェック、5) dbt build(run+test+snapshot+seed)でモデル・テスト・スナップショットを一括実行。GitHub Actions + dbt Cloudの組み合わせが最も一般的な本番構成です。
Q. Apache AirflowのDAGをKubernetesで動かすメリットは?
A. KubernetesExecutorまたはCeleryExecutorを使うと、各タスクが独立したKubernetes Podで実行されます。メリット: リソース分離(重いタスクが他のタスクに影響しない)・自動スケーリング(タスク数に応じてPodが増減)・隔離された依存関係(タスクごとに異なるPythonパッケージのDockerイメージを使える)。Helmチャート(apache-airflow/airflow)でKubernetesにデプロイでき、Google Cloud Composerはマネージドなど多くのクラウドがAirflowマネージドサービスを提供しています。
まとめ
| ユースケース | 推奨ツール |
|---|---|
| SQLによるデータ変換・テスト | dbt Core |
| 複雑なワークフロー・スケジューリング | Apache Airflow |
| SaaS→ウェアハウスのELT全工程 | Meltano |
| Fivetran代替 | Meltano または Airbyte |