別名・略称:(なし)
「etl tool」は統計データ分析の文脈で扱う重要概念のひとつ。 本ページでは「etl tool」を取り巻く中核キーワードを以下にチップで一覧化する。 各キーワードは関連する概念・手法・道具立てを含み、 文献検索や学習計画の起点になる。
これらのキーワードは「etl tool の理解 → 適用 → 検証」のプロセスを構成する。 各章で詳しく解説する。
🍰 まずはやさしく
データの整理を自動でやる道具です。
分析しやすい形に整えるために使います。
バラバラなメモをノートにまとめる感じです。
この章ではツールの基本について読みます。
ETLツール(ETL Tool):データ抽出・変換・ロードの自動化ツール
🍰 まずはやさしく
データの準備を助ける仕組みです。
バラバラな形式のデータを揃えるために使います。
スマホやアプリから情報を集める時に似ています。
なぜこのツールが必要なのかを読みます。
🍰 まずはやさしく
データの運び屋のようなものです。
集めて、変えて、入れる作業を自動化します。
部活の出欠表をまとめて集計する感じです。
具体的な3つのステップについて読みます。
| ツール | タイプ | 特徴 |
|---|---|---|
| Apache Airflow | OSS / Python | DAG でジョブ定義、 デファクト |
| dbt | OSS / SQL | DWH 内変換に特化(ELT) |
| Embulk | OSS / Java | 並列バルクロードに強い |
| Talend | 商用 / GUI | 大企業向け統合 |
ETL ツールを運用する場面では、 「データを動かした結果が正しいか」 を継続的に統計可視化する必要があります。 ここでは SSDSE-B-2026 を ETL パイプラインの題材として扱い、 拡張の図 3 点・表 3 点・Python 実装を提示します。
| ツール | Extract | Transform | Load |
|---|---|---|---|
| Airflow | Operator 経由 | Python 自由 | DB Operator |
| dbt | DWH 前提 | SQL モデル | DWH 内テーブル |
| Prefect | Task 関数 | Python/SQL | 任意 |
| Dagster | Asset 宣言 | Asset 関数 | IO Manager |
| Talend | GUI コンポーネント | GUI フロー | 多種 DB 接続 |
| Embulk | プラグイン | filter プラグイン | プラグイン |
| タスク | 頻度 | 処理時間目安 |
|---|---|---|
| CSV ダウンロード | 年 1 回 | 数秒 |
| エンコーディング統一 | 毎回 | 数秒 |
| 欠損値検査 | 毎回 | 数秒 |
| 単位変換 | 毎回 | 数秒 |
| DB ロード | 毎回 | 10-30 秒 |
| 集計マート再生成 | 毎回 | 1-5 分 |
| 観点 | 手法 | 合否基準 |
|---|---|---|
| 行数 | COUNT | 47 × 年数 |
| 欠損 | IS NULL | 0 件想定 |
| 範囲 | MIN/MAX | 負値ないか |
| 分布 | ヒストグラム | 前回と類似か |
| 一意性 | 主キー重複 | 0 件 |
| 時系列連続性 | 年抜けチェック | 飛びなし |
このコードでやること:SSDSE-B-2026 を読み込み、 ETL の Transform 工程で「行数」「欠損率」「範囲外」 をチェックし、 異常があれば例外を投げる。
📥 入力データ (SSDSE-B-2026 抜粋):
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 | import pandas as pd def etl_quality_check(df): # 1) 行数チェック assert len(df) >= 47, f"行数不足: {len(df)}" # 2) 欠損率チェック missing_rate = df['総人口'].isna().mean() assert missing_rate < 0.05, f"欠損率過大: {missing_rate:.2%}" # 3) 範囲チェック assert df['総人口'].min() > 0, "総人口に負値" assert df['総人口'].max() < 20_000_000, "総人口が想定外" # 単位は「人」 return True # skiprows=1 で読むと列名は日本語('人口総数' ではなく '総人口') df = pd.read_csv('data/raw/SSDSE-B-2026.csv', skiprows=1, encoding='cp932') df = df[df['地域コード'].astype(str).str.match(r'^R\d{5}$', na=False)].copy() df['年度'] = pd.to_numeric(df['年度'], errors='coerce') df['総人口'] = pd.to_numeric(df['総人口'], errors='coerce') # SSDSE-B-2026 の最新年度は 2023(2024 年のデータは入っていない) latest = int(df['年度'].max()) df_latest = df[df['年度'] == latest] print(f"品質チェック開始(対象年度 {latest}, {len(df_latest)} 行)") etl_quality_check(df_latest) print("✅ 品質チェック合格、 DB ロードへ") |
📤 実行すると次の出力が得られる:
💬 結果の読み方:ETL ツールに「品質チェック関数」 を組み込むことで、 異常データが DWH に流れ込むのを未然に防げる。 SSDSE-B-2026 のような統計データでも、 元データの更新ミスや CSV 破損は発生しうるため、 こうしたガードレールは不可欠。
ETL ツールは「データを動かす道具」 だが、 同時に「データ品質を担保する仕組み」 でもある。 SSDSE-B-2026 を題材に、 図1〜3 の可視化と品質チェック関数を毎ジョブに組み込めば、 安心して下流の分析・可視化に進める。
ETL は「動けばよい」では済まない。 失敗時の復旧・再実行・冪等性・スキーマ進化対応など、 運用面の機能差が大きい。 ここでは実務で頻出する 4 つの選定軸を、 SSDSE-B-2026 規模(数十 KB〜数 MB / 毎年 1 回更新)と Web 行動ログ規模(数 GB / 毎日更新)の両極で対比して整理する。
| ツール | 特徴 | SSDSE 規模での適性 | 学習コスト |
|---|---|---|---|
| pandas + cron | Python スクリプト直書き | ◎ 数 MB なら最速・最簡素 | ⭐ 低 |
| Airflow | DAG ベース・スケジューラ・UI 監視 | ◎ 教材用途・複数ジョブ並走 | ⭐⭐⭐ 中 |
| dbt | SQL 中心・テスト機能内蔵 | ○ DWH 既存環境で本領発揮 | ⭐⭐ 低〜中 |
| Prefect / Dagster | Python ネイティブ・型付き | ○ 中規模に最適、 SSDSE は overkill | ⭐⭐⭐ 中 |
| Talend / Informatica | 商用 GUI 中心 | △ ライセンスコスト過大 | ⭐⭐⭐⭐ 高 |
| AWS Glue / Azure Data Factory | クラウドネイティブ・サーバレス | ○ クラウド DWH 連携時 | ⭐⭐⭐ 中 |
「同じジョブを 2 回流しても、 結果が変わらない」状態を冪等性と呼ぶ。 ETL では 必ず冪等にすることが運用の鉄則である。 SSDSE のように年 1 回しか更新されないデータでも、 「途中で失敗 → 再実行」のシナリオは頻発する。
このコードでやること:SSDSE-B-2026 の 47 都道府県データを「年度」キーで INSERT ... ON CONFLICT DO UPDATE(PostgreSQL)で冪等にロードする手順を Python で示す。
📥 入力:SSDSE-B-2026.csv (都道府県・年度・指標列 多数)。 出力:fact_pref_yearly テーブル(year, pref_code, indicator, value)に UPSERT。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 | import pandas as pd df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df = df[df['年度'] == 2023] # 最新年度の 47 都道府県だけにする # 冪等な UPSERT SQL(PostgreSQL 構文) upsert_sql = """ INSERT INTO fact_pref_yearly(year, pref_code, indicator, value) VALUES (%s, %s, %s, %s) ON CONFLICT (year, pref_code, indicator) DO UPDATE SET value = EXCLUDED.value, updated_at = now(); """ rows = [] for _, r in df.iterrows(): for col in ['総人口', '合計特殊出生率']: if pd.notna(r[col]): rows.append((int(r['年度']), r['地域コード'], col, float(r[col]))) print(f'UPSERT 件数 = {len(rows):,} 行') print(f'最初の 1 件 = {rows[0]}') # 実際の DB 接続は psycopg2.executemany(cur, upsert_sql, rows) |
📤 実行結果:
💬 結果の読み方:47 県 × 2 指標 = 94 行を UPSERT。 同じ CSV を 2 回流しても ON CONFLICT DO UPDATE により重複は発生せず、 既存値が「最新ファイルの値」に書き換わる。 これが冪等な ETL の最小単位。
関連: ETL / データウェアハウス / データレイク / API / データガバナンス
1990 年代に Informatica や Ascential(後の IBM DataStage)が登場した頃、 ETL は「データウェアハウスを構築するための専用 GUI ツール」という位置づけだった。 GUI 上でノードを繋ぐ宣言的な設計が主流で、 SQL や Python は補助的な扱いであった。 2000 年代後半に Hadoop が普及すると ELT パターン(先にデータレイクへ生データを投入し、 後でクエリ時に変換)が台頭し、 ETL は「重い前処理」から「軽い orchestration」に役割を変えた。 2014 年に Airbnb が開発した Airflow が OSS 化されたことで、 Python DAG ベースのオーケストレーションが標準化された。 2020 年代に入ると dbt が SQL 中心の宣言的 transform 層として急速に普及し、 ETL は「データの抽出・ロード」と「SQL 変換」の二段構えに分離する設計が一般的になった。
この歴史的経緯を踏まえると、 SSDSE-B-2026 のような少量・年次更新のデータでは、 重厚な Airflow + dbt スタックは過剰投資である。 単一 Python スクリプト(pandas.read_csv → 変換 → to_parquet)と cron 起動で十分機能する。 一方、 数百テーブル × 毎日更新の大企業 DWH 環境では Airflow + dbt + Great Expectations(データ品質チェック)+ Lineage Tracker(系譜追跡)の組み合わせが事実上の標準になっている。
「ジョブが流れた = データが正しい」ではない。 ETL の本質は 「不正なデータが下流に流れ込むのを防ぐ砦」にある。 そのための具体的な実装パターンは、 以下の 5 層で整理できる。
| 層 | チェック項目 | 実装ライブラリ | SSDSE での例 |
|---|---|---|---|
| 1. 入力 | ファイル存在 / サイズ / ハッシュ | pathlib / hashlib | CSV 行数が前年より >50% 減なら警告 |
| 2. スキーマ | カラム名 / 型 / 必須列 | pandera / great_expectations | A1101 が int 型かチェック |
| 3. 値域 | 最小値 / 最大値 / 欠損率 | pandera / 自前 SQL | 人口が 0 や負値でないこと |
| 4. 整合性 | 外部キー / ユニーク制約 | dbt tests / SQL | 都道府県コードがマスタに存在 |
| 5. 統計的 | 分布の急変・外れ値率 | evidently / scipy.stats | 前年比 ±20% を超える県は要確認 |
このコードでやること:SSDSE-B-2026 を pandera スキーマで宣言的に検証する。 「人口は正の整数」「高齢化率は 0〜100」などの制約をスキーマで定義し、 違反があれば即座に例外で停止する。
📥 入力:data/raw/SSDSE-B-2026.csv。 期待スキーマは year (int)、 pref_code (str)、 A1101 (int > 0)、 A4103 (float 0..100)。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 | import pandas as pd import pandera.pandas as pa from pandera.pandas import Column, DataFrameSchema, Check schema = DataFrameSchema({ '年度': Column(int, Check.ge(2000)), '地域コード': Column(str, Check.str_matches(r'R\d{5}')), 'A1101': Column(int, Check.gt(0), nullable=True), 'A4103': Column(float, Check.in_range(0, 100), nullable=True), }) df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932') try: validated = schema.validate(df, lazy=True) print(f'OK: {len(validated)} 行を検証しました') except pa.errors.SchemaErrors as e: print(f'NG: {len(e.failure_cases)} 件の違反') print(e.failure_cases.head()) |
📤 実行結果(クリーンな SSDSE 想定):
💬 結果の読み方:pandera の lazy=True は全ての違反を集めて一度に報告する。 命令的なチェックコードを書く代わりに、 スキーマを 1 箇所に集約することで「データ契約 (data contract)」として運用できる。 不正データは ETL ジョブで早期に弾かれ、 DWH の汚染を防げる。
現代の ETL/ELT スタックでは、 「データを動かす責務」と「SQL で整形する責務」を別レイヤーに分離するのが主流である。 Airflow は前者(DAG オーケストレーション)、 dbt は後者(SQL 変換)を担当する。 この分離により、 データエンジニア(Airflow 担当)とアナリティクスエンジニア(dbt 担当)の責任分界が明確になる。 SSDSE-B-2026 のような単一データソースでも、 「raw → staging → mart」の 3 層構造を dbt で明示的に管理することで、 変更履歴と依存関係が追跡可能になる。
具体的には、 raw 層に「SSDSE-B-2026 を CSV のまま投入」、 staging 層で「列名を英語にリネーム・型変換・欠損補完」、 mart 層で「分析用の集計テーブル(都道府県別・年度別の指標サマリ)」を作る。 dbt の {{ ref('staging_ssdse') }} 記法でテーブル間の依存を宣言すれば、 dbt が自動で実行順序を解決し、 並列実行可能な部分は並列化してくれる。
ETL ジョブは 必ず失敗する前提で設計しなければならない。 ネットワーク断・ディスクフル・スキーマ変更・上流データ欠損・タイムアウト・OOM (out of memory) など、 失敗要因は多岐にわたる。 これらに対する対応策は、 (1) リトライ + 指数バックオフ、 (2) サーキットブレーカー、 (3) チェックポイント保存、 (4) アラート通知、 (5) 手動再実行コマンドの提供、 という 5 層で考えると整理しやすい。 Airflow であれば retries=3, retry_delay=timedelta(minutes=5), retry_exponential_backoff=True をタスクパラメータに設定するだけで指数バックオフリトライが実現できる。 さらに sla=timedelta(hours=1) で SLA を宣言すれば、 規定時間内に完了しない場合に自動通知が飛ぶ。
SSDSE-B-2026 のような 毎年 1 回しか更新されないデータでは、 失敗時の影響は「次のジョブまでに修正すれば済む」程度に小さい。 しかし大企業 DWH で「毎日 24:00 に開始するバッチ」が失敗した場合、 翌朝 9:00 までに復旧できないと業務に支障が出る。 こうした SLA を意識した 「失敗予算」の概念(Site Reliability Engineering で言う Error Budget)を ETL 運用にも持ち込むべきである。 月間 99.9% の SLA を約束するなら、 失敗予算は月 43 分しかなく、 1 回 1 時間止まれば SLA 違反となる。 こうした観点から「リトライ戦略」「フェイルオーバ」「代替経路」を事前に設計する。
ETL ツール選定で見落とされがちなのが 「ベンダーロックインのコスト」である。 商用ツール(Informatica、 Talend、 AWS Glue)を選ぶと、 GUI で組んだジョブが他ツールに移植できず、 数年後の移行で全工程を作り直すコストが発生する。 一方で OSS の Airflow + dbt はコードベース(Python + SQL)なので、 GitHub で履歴管理でき、 他クラウドへの移行も比較的容易である。 「初期構築は GUI ツールが速い」が、 「5 年後の移行で OSS の 5 倍のコスト」となるケースが頻発する。 SSDSE-B-2026 のような教育・研究目的であれば、 学習コストの低い pandas + cron から始め、 規模拡大に応じて Airflow + dbt へ段階的に移行する戦略が現実的である。
また、 ETL ツールのもう一つの選定軸は 「観測可能性 (observability)」である。 ジョブの実行履歴・所要時間・処理レコード数・エラー詳細が、 ダッシュボードで一目で把握できるかが運用品質を左右する。 Airflow には UI が標準で付属しているが、 観測可能性に特化した Dagster や Prefect Cloud は更に高度な可視化を提供する。 観測可能性が低いと、 「ジョブが遅い」「結果が変わった」という現象に気付くのが遅れ、 業務影響が拡大する。 SSDSE-B-2026 では年に 1 回の実行なので Slack 通知だけで十分だが、 本番 DWH では Datadog や Grafana と連携したリアルタイム監視が標準である。
ETL ジョブのパフォーマンスは、 多くの場合 I/O ボトルネックに起因する。 計算は CPU の数 GHz で進むが、 ディスク I/O は数百 MB/s、 ネットワーク I/O は更に遅い。 そのため、 「読み込み回数を減らす」「メモリ効率を上げる」「並列化する」の 3 軸が最適化の主戦場である。 具体的なテクニックとしては、 (1) カラムナフォーマット(Parquet、 ORC)の採用で必要列のみ読み込む、 (2) partition pruning で日付やキーで物理的に分割し、 不要なファイルを読み飛ばす、 (3) predicate pushdown で読み込み段階で WHERE 句を適用、 (4) columnar compression(Snappy、 Zstd)で I/O 量を圧縮、 (5) 並列パイプラインで複数ステージを同時実行、 などがある。
SSDSE-B-2026 のような数 MB 規模では上記の最適化はほぼ不要だが、 本番 DWH で数 TB を扱う場合は、 Parquet 化と partition 設計だけでジョブ時間が 1/10 になるケースが珍しくない。 例えば、 「過去 5 年の毎日の売上を集計」というジョブで、 CSV のままだと数時間かかるが、 日付別 Parquet に変換すれば直近 1 ヶ月の読み込みだけで済むので数分で完了する。 また、 メモリ効率の観点では ストリーミング処理(pandas でなく Polars や DuckDB を使う、 PyArrow のチャンク読み込み)が有効である。 数 GB の CSV を pandas で一括読み込みすると OOM になるが、 DuckDB なら SQL クエリで自動的にメモリ管理してくれる。
現代の ETL 運用では、 ジョブ定義・データ品質ルール・スケジュール・通知設定の全てを コードベースで管理し、 Git で履歴を残すのが標準となっている。 これを Infrastructure as Code (IaC) あるいは Data as Code と呼ぶ。 具体的には、 Airflow DAG を Python ファイルで定義、 dbt のモデルを SQL ファイルで定義、 データ品質ルールを schema.yml で宣言、 という形で全てを GitHub に commit する。 これにより、 「いつ・誰が・なぜ」変更したかが Pull Request の履歴に残り、 レビュアーが事前に変更を確認できる。 SSDSE-B-2026 を扱う研究プロジェクトでも、 たとえ単純なスクリプトであっても GitHub で管理することで再現性が確保される。
IaC のもう一つの利点は 環境の再現性である。 開発環境・ステージング環境・本番環境を Docker Compose や Terraform で記述し、 「どの環境でも同じジョブが動く」状態を担保できる。 これにより、 「開発では動いたのに本番で動かない」というよくある問題が解消される。 SSDSE-B-2026 のような教育・研究プロジェクトでは、 Docker で Python + pandas 環境を統一すると、 学生間の環境差異による混乱を防げる。 また、 学生が自宅 PC で実行した結果と、 教員が確認した結果が完全一致するため、 課題の採点や評価が客観的になる。
ETL 運用の品質を定量化するには、 以下の指標を継続的に計測することが重要である:
(1) ジョブ成功率: 過去 N 日で成功したジョブ数 / 総ジョブ数。 99% 以上が業界標準。
(2) 平均処理時間 (AHT, Average Handling Time): ジョブ開始から完了までの平均時間。 トレンドを見ることで性能劣化を検知。
(3) データ鮮度 (Data Freshness): 最新データが DWH に到達するまでの遅延時間。 リアルタイム要求の高い分析では数分単位が望まれる。
(4) データ品質スコア: 期待スキーマに対するルール違反数 / 検証ルール数。 90% 以上を維持。
(5) MTTR (Mean Time To Recovery): 障害発生から復旧までの平均時間。 30 分以下が理想。
これらの指標を Datadog や Grafana の ダッシュボードに常時表示し、 SLO(Service Level Objective)として組織で合意することが、 健全な ETL 運用の第一歩である。 さらに毎月の振り返り会議で「最も失敗が多かったジョブ」「最も時間がかかったジョブ」を共有し、 改善計画に落とし込むサイクルを継続することが、 長期的な ETL 品質向上の鍵となる。 SSDSE-B-2026 のような教育プロジェクトであっても、 こうした「数値で振り返る」習慣を学生時代から身につけておくと、 実務に出てからの即戦力につながる。
ETLツールの比較では、 接続先、変換の複雑さ、スケジューリング、監視、再実行、権限管理を分けて評価します。 ツール名だけでなく、 失敗時にどこから再開できるかを説明すると実務性が高まります。
🍰 まずはやさしく
データの流れを決めるルールです。
処理にかかる時間を計算するために使います。
買い物のレシートを整理して記録する感じです。
データの流れを図と式で詳しく読みます。
ETL は概念図で表すのが分かりやすい:
$T_{\text{overhead}}$ にはスケジューラ起動、 接続確立、 メタデータ更新が含まれる。 小さなタスクが大量にある場合、 オーバーヘッドが支配的になる(small file problem)。
100 万行のテーブルに 1000 行追加された場合、 増分ロード 1.5 分・全件再ロード 15 分なら $\eta = 1 - 1.5/15 = 0.9$(90% 効率化)。 ただし増分追跡のオーバーヘッド・更新検出ミスのリスクとのトレードオフ。
並列実行可能な DAG では、 全体時間は最長依存経路で決まる。 Airflow / Dagster の DAG 設計では、 タスク分割を細かくし並列度を上げるのが基本戦略。
dbt は SQL ベースで Transform を書く。 SSDSE-B-2026 を BigQuery にロード後、 次の models/ssdse_b_summary.sql のような変換を Git 管理する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 | -- models/ssdse_b_aged.sql
{{ config(materialized='table') }}
with src as (
select 都道府県 as prefecture,
年度 as year,
総人口 as population,
65歳以上人口 as aged_pop
from {{ source('raw','ssdse_b_2026') }}
where 年度 = 2023
)
select prefecture, population, aged_pop,
round(aged_pop * 100.0 / population, 2) as aged_ratio
from src
order by aged_ratio desc
|
ref() でモデル間依存を解決し DAG を自動生成tests: not_null, unique, accepted_values で品質保証snapshots: SCD2 で履歴管理docs: 自動ドキュメント+カラム説明ETL ツール選定が「Airflow か dbt か Prefect か」で揉めるのは、 みんなが暗黙に違うコスト関数を見ているから。 ここでは パイプライン総コスト を式で固定して比較する。
$$ C_{\text{total}} = \sum_{i=1}^{n} \underbrace{c_i \cdot s_i}_{\text{計算}} + \underbrace{\tau \cdot s_i \cdot d_i}_{\text{転送}} + \underbrace{\eta \cdot s_i}_{\text{ストレージ}} $$
s_i = タスク i のデータサイズ (GB)、 c_i = 計算単価 ($/GB)。 ETL の T (Transform) は マシン上で実行するので c_i × s_i が支配的、 ELT は DWH 上で実行するので c_i ≈ クエリ単価 × 圧縮率。τ = ネットワーク転送単価、 d_i = 転送距離 (リージョン間=1, リージョン内=0.1)。 ETL は 抽出マシン → 加工マシン → DWH と 2 ホップなので τ·s·d が 2 倍。 ELT は 1 ホップ。η = ストレージ単価 ($/GB/月)。 中間ファイル (ETL の I/O 領域) は η·s_i が積み上がる。 ELT は DWH 内 view で完結すれば η ≈ 0。SSDSE データを使った ETL の例:
prefecture_stats テーブルに INSERTSSDSE-B-2026 のような公的データを定期更新する場合、 手書きスクリプトより専用 ETL ツールを使うと再実行性・モニタリング・依存関係管理が一気に楽になる。
| ツール | タイプ | 強み | SSDSE-B での使い所 |
|---|---|---|---|
| Apache Airflow | ワークフローエンジン | Python DAG, スケジュール | 年次データの自動取得→変換→ロード |
| dbt | SQL ベース変換 | Git連携, テスト | DWH に入った後の集計クエリ管理 |
| Fivetran / Airbyte | コネクタ | API 連携豊富 | e-Stat API → BigQuery |
| Talend / Informatica | GUI ETL | ノーコード | 業務側がメンテ |
| Apache NiFi | ストリーミング | フロー視覚化 | センサー・ログのリアルタイム取り込み |
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 | import pyarrow # parquet の読み書きに必要(ブラウザには無い) # Airflow DAG で SSDSE-B-2026 を ETL する例 from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import pandas as pd import sqlite3 def extract(**ctx): df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df.to_parquet('/tmp/ssdse_raw.parquet') def transform(**ctx): df = pd.read_parquet('/tmp/ssdse_raw.parquet') df = df[df['年度']==2023].copy() df['高齢化率'] = df['65歳以上人口']/df['総人口']*100 df.to_parquet('/tmp/ssdse_t.parquet') def load(**ctx): df = pd.read_parquet('/tmp/ssdse_t.parquet') con = sqlite3.connect('ssdse.db') df.to_sql('ssdse_b', con, if_exists='replace', index=False) con.close() with DAG('ssdse_b_etl', start_date=datetime(2026,1,1), schedule='@yearly', catchup=False) as dag: t1 = PythonOperator(task_id='extract', python_callable=extract) t2 = PythonOperator(task_id='transform', python_callable=transform) t3 = PythonOperator(task_id='load', python_callable=load) t1 >> t2 >> t3 |
SSDSE-B-2026 は cp932 エンコード + ヘッダ 2 行(英語+日本語)という、 公的データあるあるの「読み込み難易度高い」CSV。 これを pandas で安全に Extract する。
🎯 このコードでやること: SSDSE-B-2026 を cp932 エンコードで読み込み、 2 行目(日本語ヘッダ)をスキップしつつ年度 2023 行を抽出する。
1 2 3 4 5 6 7 8 9 10 11 12 | import pandas as pd import time t0 = time.time() df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1], # 2 行目(日本語列名)をスキップ low_memory=False ) t_extract = time.time() - t0 print(f'抽出 {len(df)} 行 / {df.shape[1]} 列 / {t_extract*1000:.1f} ms') print(f'年度 unique: {sorted(df["SSDSE-B-2026"].unique())}') |
📤 実行結果(参考値):
💬 結果の読み方: 12 年分 × 47 都道府県 = 564 行を 87 ms で抽出。 cp932 + skiprows でヘッダ 2 段構造を正しく処理。 大規模データなら chunk 読み込み (chunksize=10000) や Polars を検討。
SSDSE-B-2026 から「高齢化率」「人口密度ランク」「都市圏フラグ」を作る。
🎯 このコードでやること: 派生列を 3 個追加し、 各列の処理時間を測定する。
1 2 3 4 5 6 7 8 9 10 11 12 13 | import pandas as pd import time df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023].copy() t0 = time.time() df['高齢化率'] = df['A1303'] / df['A1101'] * 100 df['人口密度ランク'] = df['A1101'].rank(ascending=False).astype(int) df['都市圏フラグ'] = df['Prefecture'].isin(['東京都','神奈川県','埼玉県','千葉県','大阪府','兵庫県','京都府','愛知県']) t_transform = time.time() - t0 print(f'変換 3 列追加 / {t_transform*1000:.2f} ms') print(df[['Prefecture','高齢化率','人口密度ランク','都市圏フラグ']].head(5)) |
📤 実行結果(参考値):
💬 結果の読み方: 47 行 × 3 列の変換は 3.4 ms。 高齢化率 33% (北海道)、 38% (秋田)が見える。 都市圏 8 都府県を boolean で抽出。 これを ETL ツールで定期実行すれば、 BI で常に最新指標が見られる。
🎯 このコードでやること: 変換後 DataFrame を SQLite に書き出し、 メタデータ(テーブル名・行数・実行時刻)を別テーブルに記録する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 | import pandas as pd import sqlite3 import time from datetime import datetime df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023].copy() df['高齢化率'] = df['A1303'] / df['A1101'] * 100 t0 = time.time() conn = sqlite3.connect('data/dwh.sqlite') df.to_sql('prefecture_demographics_2023', conn, if_exists='replace', index=False) t_load = time.time() - t0 # メタデータ記録 meta = pd.DataFrame([{ 'table': 'prefecture_demographics_2023', 'rows': len(df), 'cols': df.shape[1], 'load_time_ms': round(t_load*1000, 2), 'loaded_at': datetime.now().isoformat(timespec='seconds'), 'source': 'data/raw/SSDSE-B-2026.csv', }]) meta.to_sql('etl_metadata', conn, if_exists='append', index=False) print(meta.iloc[0].to_dict()) conn.close() |
📤 実行結果(参考値):
🕐 load_time_ms と loaded_at は実行のたびに変わります(行数 47・列数 113 は変わりません)。
💬 結果の読み方: SQLite を DWH 代用にして 47 行 × 113 列を 10 ms 前後で書き込み。 メタデータ別テーブルに「いつ」「どこから」「どれだけ」を残すのが ETL の鉄則 ── これがないと障害時の re-run・行数チェックが地獄。 BigQuery では INFORMATION_SCHEMA が同等機能を提供。
🎯 このコードでやること: SSDSE-B-2026 の ETL を Airflow DAG に落とし、 「extract → transform → load → notify」の 4 タスクを定義する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 | import pyarrow # parquet の読み書きに必要(ブラウザには無い) # dags/ssdse_etl_dag.py from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime, timedelta import pandas as pd import sqlite3 def extract(**ctx): df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df.to_parquet('/tmp/ssdse_raw.parquet') return len(df) def transform(**ctx): df = pd.read_parquet('/tmp/ssdse_raw.parquet') df = df[df['SSDSE-B-2026'] == 2023].copy() df['高齢化率'] = df['A1303'] / df['A1101'] * 100 df.to_parquet('/tmp/ssdse_curated.parquet') return len(df) def load(**ctx): df = pd.read_parquet('/tmp/ssdse_curated.parquet') conn = sqlite3.connect('data/dwh.sqlite') df.to_sql('prefecture_2023', conn, if_exists='replace', index=False) conn.close() return len(df) def notify(**ctx): ti = ctx['ti'] n = ti.xcom_pull(task_ids='load') print(f'Slack: ETL 完了 {n} 行ロード') with DAG( 'ssdse_etl', schedule_interval='@daily', start_date=datetime(2026, 1, 1), catchup=False, default_args={'retries': 3, 'retry_delay': timedelta(minutes=5)}, tags=['ssdse', 'demographics'], ) as dag: e = PythonOperator(task_id='extract', python_callable=extract) t = PythonOperator(task_id='transform', python_callable=transform) l = PythonOperator(task_id='load', python_callable=load) n = PythonOperator(task_id='notify', python_callable=notify) e >> t >> l >> n |
📤 実行結果(Airflow UI 上):
💬 結果の読み方: 4 タスクが順序通り実行され、 各タスクの行数を XCom で次タスクに渡している。 retries=3 + retry_delay=5min により一時的なネットワーク障害は自動回復。 失敗時に Slack 通知する on_failure_callback を設定すれば運用性向上。
ELT モデルでは Transform は「DWH 内の SQL」で書く。 dbt (data build tool) は SQL に Jinja を加えてモジュール化・テスト・lineage 自動生成を可能にする。 SSDSE-B-2026 を BigQuery にロード済みと仮定した dbt モデル例。
🎯 このコードでやること: dbt の SQL モデルで「都道府県別高齢化率」を計算し、 unique / not_null のテストを定義する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 | -- models/marts/prefecture_aging_2023.sql {{ config(materialized='table', cluster_by=['region']) }} SELECT Prefecture AS prefecture, A1101 AS total_population, A1303 AS elderly_population, ROUND(A1303/A1101*100, 2) AS aging_rate_pct, CASE WHEN A1303/A1101 > 0.35 THEN 'high' WHEN A1303/A1101 > 0.28 THEN 'medium' ELSE 'low' END AS aging_level, CURRENT_TIMESTAMP() AS computed_at FROM {{ ref('stg_ssdse_b_2026') }} WHERE year = 2023 |
📤 schema.yml (テスト定義):
💬 結果の読み方: dbt は (1) SQL モデルで Transform、 (2) ref() で依存関係自動構築、 (3) schema.yml で型・値域テスト、 (4) docs で lineage 可視化、 をワンストップで提供。 47 都道府県データなので行数テストを 47 に固定、 高齢化率は 15-45% の範囲に収まることをアサート。 これによりデータ品質が CI/CD と統合される。
合成データで 4 ツールの加重スコアを計算する (性能, 価格, 学習コスト, サポート)。
| ツール | 性能 0.40 | 価格 0.30 | 学習 0.20 | サポート 0.10 | 加重 |
|---|---|---|---|---|---|
| A | 5 | 3 | 4 | 5 | 4.20 |
| B | 4 | 5 | 3 | 4 | 4.10 |
| C | 3 | 4 | 5 | 3 | 3.70 |
| D | 5 | 4 | 3 | 4 | 4.20 |
1 2 3 4 5 6 | import numpy as np weights = np.array([0.40, 0.30, 0.20, 0.10]) tools = np.array([[5,3,4,5],[4,5,3,4],[3,4,5,3],[5,4,3,4]]) scores = tools @ weights print(f"加重スコア: {scores}") print(f"最高 index: {np.argwhere(scores == scores.max()).flatten()}") |
💬 手計算 (Step 2) A=4.20 と Python 出力が完全一致。
SSDSE-B-2026(47 都道府県・2023 年データ)を題材にした最小コード:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 | # pandas で簡易 ETL import os import pandas as pd # Extract: CSV から抽出 df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) # Transform: 列名統一、 欠損除去、 型変換 # 実際の列名は '消費支出(二人以上の世帯)' と '年度' df = df.rename(columns={'地域コード': 'pref_code', '消費支出(二人以上の世帯)': '消費支出', '年度': '年'}) df = df[df['pref_code'].astype(str).str.match(r'^R\d{5}$', na=False)].copy() df['消費支出'] = pd.to_numeric(df['消費支出'], errors='coerce') df = df.dropna(subset=['消費支出']) df['年'] = df['年'].astype(int) # Load: SQLite に書き込み(実運用は PostgreSQL / BigQuery 等) import sqlite3 os.makedirs('data/processed', exist_ok=True) conn = sqlite3.connect('data/processed/etl_demo.db') df.to_sql('prefecture_stats', conn, if_exists='replace', index=False) print(f'{len(df)} 行を prefecture_stats に書き込みました') |
🎯 このコードでやること: Airflow の @task デコレータで SSDSE-B-2026 を medallion architecture (raw → cleaned → mart) に流し、 タスクごとの実行時間とサイズから上の C_total を実測する。
📥 入力例 (SSDSE-B-2026 を Airflow の data/raw/ に置いた状態):
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 | import os os.makedirs('data', exist_ok=True) # 書き出し先を先に作る import os os.makedirs('data/bronze', exist_ok=True) # 保存先のフォルダを作っておく os.makedirs('data/silver', exist_ok=True) # 保存先のフォルダを作っておく os.makedirs('data/gold', exist_ok=True) # 保存先のフォルダを作っておく import pyarrow # parquet の読み書きに必要(ブラウザには無い) from airflow.decorators import dag, task from datetime import datetime import pandas as pd, time @dag(start_date=datetime(2026, 5, 1), schedule='@monthly', catchup=False) def ssdse_medallion(): @task def bronze() -> str: df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=1) df.to_parquet('data/bronze/pop.parquet') # 生データ保全 return 'data/bronze/pop.parquet' @task def silver(p: str) -> str: df = pd.read_parquet(p) df = df.dropna(subset=['総人口']).query('総人口 > 0') # 品質ルール df.to_parquet('data/silver/pop_clean.parquet', compression='snappy') return 'data/silver/pop_clean.parquet' @task def gold(p: str) -> str: df = pd.read_parquet(p).sort_values(['都道府県', '年度']) df['年率'] = df.groupby('都道府県')['総人口'].pct_change() * 100 df.groupby('都道府県')['年率'].mean().to_frame('mean_growth_pct') \ .to_parquet('data/gold/pop_growth_mart.parquet') return 'data/gold/pop_growth_mart.parquet' gold(silver(bronze())) ssdse_medallion() |
📤 実行例 (Airflow UI の Gantt から抜粋):
💬 結果の読み方: 各層のサイズが 4.2 → 1.8 → 0.3 MB と 段階的に縮むのが医療パイプラインの理想形。 gold が η·s_i ≈ 0 なので、 BI ダッシュボード側は gold を直接読めばコストゼロでクエリできる。 もしこの SSDSE が S3 上にあったら、 silver→gold のリージョン間転送 τ·1.8MB·d が新たな支配項に変わる。
🎯 このコードでやること: SSDSE-B-2026 の「行数 47」「Prefecture 列ユニーク」「人口 0 以下なし」を Great Expectations で検証する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 | import pandas as pd import great_expectations as ge df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023].copy() gdf = ge.from_pandas(df) r1 = gdf.expect_table_row_count_to_equal(47) r2 = gdf.expect_column_values_to_be_unique('Prefecture') r3 = gdf.expect_column_values_to_be_between('A1101', min_value=1, max_value=2e7) r4 = gdf.expect_column_values_to_not_be_null('A1303') for name, r in zip(['行数=47','Prefユニーク','人口1-2千万','高齢人口非null'], [r1,r2,r3,r4]): print(f'{name:15s}: success={r["success"]}') |
📤 実行結果(参考値):
💬 結果の読み方: 4 つのアサーションが全て True。 もし行数が 46 や 48 になっていたら即検知。 Airflow DAG の前段に組み込めば、 異常データがそのまま下流に流れ込むのを防止。 Great Expectations は 200+ の expectation 関数を持ち、 統計分布の検証も可能。
Databricks 提唱の 3 層モデル。 ETL の Load 先を 1 つの「ピカピカに整形済みテーブル」にせず、 段階的に磨いていく。
| 層 | 内容 | SSDSE 例 | 用途 |
|---|---|---|---|
| Bronze (Raw) | 無加工、 cp932 のまま | SSDSE-B-2026 そのまま | 監査・再処理用 |
| Silver (Clean) | 型統一、 列名英語化、 欠損処理 | prefecture, year, total_pop, ... | 分析者の探索 |
| Gold (Curated) | 集計済み、 ビジネス KPI 化 | region_aging_summary | BI・ダッシュボード |
Bronze を残すことで「集計ロジックを変えたとき再処理可能」、 Gold で BI が高速に動く、 Silver で分析者が自由にクエリできる、 という三者三様の用途を満たす。 Delta Lake / Iceberg / Hudi はこのアーキテクチャをトランザクション付きで実現する。
「翌日朝までに集計」では遅すぎる用途(不正検知、 在庫リアルタイム、 IoT ダッシュボード)では、 イベント駆動のストリーミング ETL が必要。 SSDSE-B-2026 のような静的データには不要だが、 全体像として:
ウィンドウ処理(5 分 sliding window で平均、 1 時間 tumbling window でカウント)、 遅延データの扱い(watermark)、 状態管理など、 バッチ ETL とは異なる難易度の概念群がある。
100 万 MAR / 10 ソース / 1TB データ / 50 ユーザー想定の月額シミュレーション:
| 構成 | 月額 (USD) | エンジニア工数 |
|---|---|---|
| Fivetran + Snowflake + dbt Cloud | $3,000 - $8,000 | 低(0.3 人月/月) |
| Airbyte OSS + Snowflake + dbt Core | $1,500 - $3,000 | 中(1 人月/月) |
| Airflow OSS + BigQuery + dbt Core | $800 - $2,500 | 高(2 人月/月) |
| 完全 OSS (Airflow + Postgres + dbt) | $200 - $500 | 超高(3+ 人月/月) |
「マネージド = 高額だが運用人件費低」のトレードオフ。 1 人エンジニアが月給 $10K と仮定すると、 Fivetran $5K + 工数 0.3 人月 = $8K vs Airflow + 工数 2 人月 = $21K。 規模 1TB 程度ならマネージドが結果安いケースが多い。
「extract_transform_load を 1 タスクに詰め込む」は失敗時の再実行範囲が広く、 デバッグも困難。 SSDSE-B-2026 を 3 タスクに分けると、 transform で失敗しても extract をやり直さずに済む(中間ファイルが活用される)。
「特定日付の特定タスクだけ再実行」できる設計が必須。 Airflow は execution_date を全タスクに自動注入、 dbt は --vars '{date: "2026-05-24"}' でパラメータ渡し。
「SSDSE-B-2026 のソース CSV がアップロードされたら DAG が起動」のように、 cron + センサで疎結合化できる。 Airflow の S3KeySensor, Dagster の sensor 機能、 Prefect の flow.serve() がこれを担う。
過去 6 年分の SSDSE-B-2026 を一括処理する場合、 並列度を制限しないと DWH が過負荷。 Airflow の max_active_runs, pool で並列度制御。
「7:00 までに完了」を SLA とし、 6:30 までに完了しなければ Slack 警告、 7:00 過ぎたら PagerDuty 起動。 Airflow の sla_miss_callback が標準機能。
「ソース側が予告なく列名変更 → 下流のダッシュボードが崩壊」を防ぐため、 ソースとコンシューマの間で明示的な契約を結ぶ。 Pydantic / Pandera / dbt contracts / OpenLineage がこれを担う。
🎯 このコードでやること: Pandera で SSDSE-B-2026 のスキーマ契約を定義し、 違反時に例外を吐かせる。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 | import pandas as pd import pandera as pa from pandera import Column, DataFrameSchema, Check schema = DataFrameSchema({ 'Prefecture': Column(str, Check(lambda s: s.str.endswith(('都','道','府','県')).all())), 'A1101': Column(int, Check.in_range(1, 20_000_000)), 'A1303': Column(int, Check.greater_than_or_equal_to(0)), 'A4103': Column(float, Check.in_range(0.5, 3.0), nullable=True), }) df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023].copy() try: validated = schema.validate(df[['Prefecture','A1101','A1303','A4103']]) print(f'✅ スキーマ契約合格: {len(validated)} 行') except pa.errors.SchemaError as e: print(f'❌ 契約違反: {e}') |
📤 実行結果:
💬 結果の読み方: 4 列の型と値域を契約として明示し、 ソース更新後の自動検証ができる。 もし「Prefecture」が新たに英語表記の混入で「Tokyo」が来たら例外発生 → CI で即検知。 Data Contract は ETL ツールの中核機能になりつつある。
業務 DB から DWH へリアルタイム同期する場合、 「変更分だけを取り込む」CDC が定番。 主流は 4 方式:
| 方式 | 仕組み | 遅延 | 負荷 |
|---|---|---|---|
| タイムスタンプ比較 | updated_at 列で diff | 10 分〜 | 中 |
| トリガー | DB トリガーで変更ログ書込 | 数秒〜 | 高 |
| スナップショット差分 | 2 時点の全件比較 | 1 時間〜 | 超高 |
| ログベース (WAL) | Postgres WAL / MySQL binlog 解析 | 1 秒以内 | 最小 |
Debezium はログベース CDC の事実上の標準。 Postgres → Kafka → Snowflake のリアルタイムパイプラインが構築できる。 Fivetran も内部的にはログベース CDC を実装している。
「この KPI の数値、 どのソースからどう加工された?」を辿るための仕組み。 障害時のインパクト分析・PII 管理・規制対応に必須。
採用構成: Fivetran (Shopify, Stripe, Salesforce → BigQuery) + dbt Cloud + Looker。 月額 $6K + dbt $200/user × 8 = $7.6K。 エンジニア 1 人で運用、 ビジネス側のアナリスト 7 人が自由にダッシュボード化。 GUI なし・全部 SQL/YAML、 PR レビュー文化が定着。
採用構成: Informatica PowerCenter (オンプレ) + Oracle DWH + Tableau。 機密情報のため社外送出不可、 完全オンプレ運用。 GUI 駆動 ETL ジョブ 3000+ 本、 老朽化により ELT (Snowflake on Azure) への移行を 5 ヵ年計画で実施中。
採用構成: 自前 Python スクリプト + cron + PostgreSQL + Metabase。 SSDSE-B-2026 のような統計データを e-Stat API から取得し、 BI で公開。 予算制約から OSS 中心、 ベテラン職員 1 名が運用。 規模感はこれくらいでも十分回る。
dbt が普及した最大の理由はこの「テスト文化」を SQL の世界に持ち込んだこと。 ETL ツールを選ぶ際は「テストの書きやすさ」を必ず評価項目に入れる。
SSDSE-B-2026 は毎年 1 回更新される程度だが、 業務 DB の取り込みでは「毎時 / 毎分の増分」が必要。 代表 4 パターン:
ログテーブルのように UPDATE/DELETE がないテーブル。 「前回までの最大 ID 以降」を SELECT するだけで増分が取れる。 最もシンプル。
「updated_at > last_sync_time」で差分取得。 ただし NTP ずれ・トランザクション間隔で取りこぼしが起きうるため、 オーバーラップ(5 分前から取り直す)が安全。
主キーで一致したら UPDATE、 なければ INSERT。 BigQuery の MERGE, Snowflake の MERGE, PostgreSQL の ON CONFLICT DO UPDATE が標準。 dbt の incremental マテリアライゼーションは内部でこれを生成。
「東京都の人口が変わったが、 過去時点の値も保持したい」場合、 同じ主キーで複数行を valid_from / valid_to 付きで保持。 BI で「2020 年時点の人口」が正確に出る。 SSDSE-B-2026 の年度列はまさにこの構造。
🎯 このコードでやること: SSDSE-B-2026 の年度を SCD Type 2 風に管理し、 各都道府県 × 年度で「有効期間付き行」を生成する。
1 2 3 4 5 6 7 8 9 10 11 12 13 | import pandas as pd df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) # SCD Type 2: 都道府県×年度で「有効期間付き」テーブル化 df = df.rename(columns={'SSDSE-B-2026': 'year'}).sort_values(['Prefecture', 'year']) df['valid_from'] = df['year'].astype(str) + '-01-01' df['valid_to'] = df.groupby('Prefecture')['year'].shift(-1).fillna(9999).astype(int).astype(str) + '-01-01' df['is_current'] = df['valid_to'] == '9999-01-01' scd = df[['Prefecture','year','A1101','valid_from','valid_to','is_current']] print(scd[scd['Prefecture']=='東京都']) print(f'現在有効行: {scd["is_current"].sum()}') |
📤 実行結果(参考値):
💬 結果の読み方: 東京都の人口推移が 6 行で保持され、 各行に「いつから/いつまで有効か」が記録される。 BI で「2020 年時点の東京都人口」を問い合わせると 14,048,000 が返り、 「現時点の人口」は 14,086,000 が返る。 9999-01-01 は「現在まで有効」を示す sentinel 値。 47 都道府県 × 6 年 = 282 行のうち、 47 行が is_current=True。
同じ「ETL タスク群」でも、 タスクの依存関係を直列 vs 並列で組むと総時間が桁違いに変わる。 SSDSE 系 6 ファイル (SSDSE-A 〜 F) を取り込むケースで:
| パターン | DAG 形状 | 総時間 | 障害時 |
|---|---|---|---|
| 直列 | A→B→C→D→E→F | 6 × 10s = 60s | C 失敗で D-F 未実行 |
| 並列 (扇型) | A,B,C,D,E,F 同時 | max(10s) = 10s | 他のタスクは影響なし |
| バッチ + Merge | [A-F 並列] → Merge | 10s + 5s = 15s | Merge 失敗で全体やり直し |
並列度はDWH の並列クエリ上限とのトレードオフ。 BigQuery は default 100、 Snowflake は warehouse サイズ依存。 Airflow の pool 機能でツール側で抑制可能。
API レスポンスはネストした JSON が典型。 これをリレーショナルテーブルに展開(flatten)する処理は ETL ツールの肝。 e-Stat API から取得した SSDSE 系データを想定した模擬例:
🎯 このコードでやること: SSDSE-B-2026 を JSON 形式にして送る API を想定し、 pandas の json_normalize でネストを展開する。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 | import pandas as pd import json df = pd.read_csv('data/raw/SSDSE-B-2026.csv', encoding='cp932', skiprows=[1]) df = df[df['SSDSE-B-2026'] == 2023].head(3).copy() # SSDSE を「API レスポンス風 JSON」に整形 api_response = { 'meta': {'source': 'SSDSE-B-2026', 'fetched_at': '2026-05-24T10:00'}, 'data': [ { 'pref': {'code': row['Code'], 'name': row['Prefecture']}, 'population': {'total': row['A1101'], 'elderly': row['A1303']}, 'fertility': row['A4103'], } for _, row in df.iterrows() ] } # json_normalize でフラット化 flat = pd.json_normalize(api_response['data'], sep='_') print(flat) |
📤 実行結果:
💬 結果の読み方: ネストキーが pref_code のようにアンダースコア結合され、 リレーショナル形式になった。 BigQuery では UNNEST(JSON_EXTRACT_ARRAY())、 Snowflake では FLATTEN() で同等処理可能。 ETL ツールが Web API 連携を売りにする理由はこの「JSON → テーブル化」を肩代わりしてくれること。
| DWH | 提供形態 | 課金軸 | 強み | 弱み |
|---|---|---|---|---|
| BigQuery | GCP マネージド | スキャンバイト | サーバレス、 ML 統合 | scan 課金で予測難 |
| Snowflake | マルチクラウド | クレジット (時間) | 分離 compute/storage | warehouse 起動オーバヘッド |
| Redshift | AWS マネージド | ノード時間 | AWS 統合、 RA3 で分離可 | vacuum/sort 運用 |
| Databricks SQL | Lakehouse | DBU 時間 | Spark + Delta + ML | 習熟必要 |
| DuckDB | OSS, ローカル | 無料 | 「ローカルで動く分析 DB」 | スケールアウト不可 |
| ClickHouse | OSS / ClickHouse Cloud | 無料/利用量 | 超高速 OLAP | JOIN 弱い |
ETL パイプラインを「動いていること」だけでなく「正しく動いていること」を監視する 4 層モデル:
Monte Carlo / Bigeye / Soda などは「Data Observability」と呼ばれる新興カテゴリ。 ETL ツールとは別レイヤーで、 異常を「早期に・自動で・原因まで」検知する。
老朽化した Informatica / SSIS / Talend を Fivetran + dbt + Snowflake に移行するプロジェクトは、 2020 年代の典型。 段階的アプローチ:
よくある失敗: 「3 ヵ年計画」のはずが 5 年経っても完了しない。 原因は (1) 旧ジョブのドキュメント欠落で機能要件特定不能、 (2) 業務側の検証協力が得られない、 (3) 並走中の二重運用コスト。 「移行スコープを絞る」「業務側の優先順位を逆算」がコツ。
SSDSE-B-2026 のような確定済み静的データを扱うだけなら、 pandas + cron で十分。 業務 DB が増え、 リアルタイム要件が出てきた段階で Airflow / Fivetran 等を検討するのが現実的。
| ツール | カテゴリ | 特徴 | 価格帯 |
|---|---|---|---|
| Apache Airflow | OSS オーケストレーション | Python DAG、 拡張性高、 学習曲線急 | 無料 |
| Prefect | OSS + SaaS | モダン Python、 動的タスク、 Hybrid 実行 | OSS / 有料 |
| Dagster | OSS + SaaS | asset 中心、 型安全、 lineage 強い | OSS / 有料 |
| Fivetran | SaaS マネージド | 300+ ソースコネクタ、 自動スキーマ追随 | 高額(MAR 課金) |
| Airbyte | OSS + SaaS | Fivetran のオープン版、 自分でコネクタ書ける | OSS / 有料 |
| Stitch | SaaS マネージド | Talend 傘下、 Singer 規格 | 中 |
| dbt (Core/Cloud) | SQL 変換 (T) | SQL + Jinja、 lineage 可視化、 テスト | 無料/有料 |
| Informatica PowerCenter | エンタープライズ ETL | 老舗、 GUI 駆動、 大企業実績 | 超高額 |
| Talend Open Studio | OSS GUI ETL | Java ベース、 自動コード生成 | 無料 |
| SSIS | MS SQL Server 付属 | SQL Server エコシステム特化 | SQL Server ライセンス込 |
| AWS Glue | マネージド (Spark) | サーバレス、 Crawler 自動カタログ化 | 従量課金 |
| Google Dataflow | マネージド (Beam) | ストリーミング/バッチ統一 | 従量課金 |
SSDSE-A (個票) / B (時系列, 47 県) / C (人口統計) / D (家計) / E (世帯) / F (家計内訳) を全部取り込み、 統合分析テーブルを作る ETL 演習。
🎯 このコードでやること: 6 ファイルを並列で読み込み、 各データのメタ情報(行数・列数・期間)を集計して 1 テーブルにまとめる。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 | import pandas as pd from concurrent.futures import ThreadPoolExecutor import time SOURCES = { 'SSDSE-A': 'data/raw/SSDSE-A-2025.csv', 'SSDSE-B': 'data/raw/SSDSE-B-2026.csv', 'SSDSE-C': 'data/raw/SSDSE-C-2026.csv', 'SSDSE-D': 'data/raw/SSDSE-D-2023.csv', 'SSDSE-E': 'data/raw/SSDSE-E-2026.csv', 'SSDSE-F': 'data/raw/SSDSE-F-2023v3.csv', } def extract_meta(name_path): name, path = name_path t0 = time.time() df = pd.read_csv(path, encoding='cp932', skiprows=[1], low_memory=False) elapsed = time.time() - t0 return {'name': name, 'rows': len(df), 'cols': df.shape[1], 'load_ms': round(elapsed*1000, 1)} with ThreadPoolExecutor(max_workers=6) as ex: results = list(ex.map(extract_meta, SOURCES.items())) summary = pd.DataFrame(results) print(summary) print(f'合計行数: {summary["rows"].sum()}') print(f'並列ロード時間: {summary["load_ms"].max():.1f} ms (max of 6)') |
📤 実行結果(参考値):
🕐 load_ms と「並列ロード時間」は実行のたびに変わります(行数・列数は変わりません)。
💬 結果の読み方: 6 ファイルを ThreadPoolExecutor で並列ロード。 見るべきは絶対値ではなく「合計」と「最大」の関係で、 上の実行例では 6 本の load_ms の合計が約 152 ms なのに対し、 並列ロード時間は最大値の 46 ms。 直列なら合計時間、 並列なら最も遅い 1 本が律速という並列処理の基本がそのまま出ている。 Airflow DAG ならこの並列度を pool や max_active_tasks で制御。 大規模なら Dask / Polars / Spark に移行を検討。
[シナリオ別推奨構成] case A: 個人プロジェクト / SSDSE 系公開データ取り込み → pandas + cron + SQLite (or DuckDB) で十分 → 規模 100 万行未満なら何もスケール不要 case B: スタートアップ (data team 1-3 人) → Airbyte (OSS) + dbt Core + BigQuery + Looker Studio → 月額 $500-2000、 工数 0.5 人月 case C: 中規模 (data team 5-10 人, 10TB) → Fivetran + dbt Cloud + Snowflake + Looker → 月額 $5K-15K、 工数 0.3 人月 case D: 大企業 (data team 20+, 100TB+) → Airflow + dbt Cloud + Snowflake + Tableau + DataHub → 月額 $50K+、 工数 5+ 人月 case E: 規制業界 (金融/医療, 完全オンプレ) → Informatica or 自作 + Oracle/Postgres + Tableau Server → 年額 $100K+、 工数 10+ 人月 [避けるべきアンチパターン] ❌ 1 つの巨大 Python スクリプトに全 ETL を詰め込む ❌ DB パスワードを Git にコミット ❌ 失敗時に手動で再実行を要する設計 ❌ ドキュメント・lineage なしの「秘伝のタレ ETL」 ❌ テストゼロでプロダクション投入
ETL ツールは単独のソフトではなく、 ソース DB ・API ・ファイル ・DWH ・データレイクを結ぶオーケストレーション基盤である。 Airflow ・dbt ・Talend ・Informatica の選択は組織規模と運用要件で決まる。
ETL ツールは「抽出・変換・ロードを GUI/コードで自動化する」ソフト群で、 上流のソース DB・API・ファイルを統合し、 並列の dbt・Airflow・Talend・Informatica などから組織に合うものを選び、 下流の DWH・データレイクへ定期投入する。
ETL ツール導入の判断は「ジョブ本数・依存関係の複雑さ・運用チーム規模」の 3 軸で決まる。 数本までは Cron + Python、 依存と再実行が増えたら Airflow、 BI 連携重視なら dbt が選択肢。
SSDSE-B-2026 のような月次・年次データを社内基盤に取り込む場合、 まず手作業 Python スクリプトで原型を作り、 運用本番になった時点で Airflow など ETL ツールに昇格させる二段階アプローチが現実的である。