OSテーマ

データ加工に使う製品・ライブラリの全体像

データ加工に使う製品・ライブラリの全体像

まず、「取り込み」「保存」「加工・計算」「加工手順の管理」「実行管理」に分けると整理しやすくなります。以下は主要カテゴリを網羅した地図で、製品の全件一覧ではありません。

例えば、S3に保存したデータをSparkで加工し、Icebergでテーブル管理し、Trinoから検索する構成があります。これらは同時に使う道具です。

1. 1台のPC・サーバーで加工する

  • pandas:Pythonの表形式データ処理。探索的な分析、細かな前処理、Pythonの他ライブラリとの連携に向きます。既存コードの豊富さも利点です。
  • Polars:並列処理と遅延実行による最適化を活用するDataFrameライブラリ。繰り返す変換・結合・集計を効率よく処理したい場合の候補です。
  • DuckDB:手元で動く分析向けSQLデータベース。CSV・ParquetなどをSQLで直接扱え、別のDBサーバーを用意する必要がありません。

目安は「Pythonの既存資産ならpandas」「DataFrame処理の効率ならPolars」「SQLで考えるならDuckDB」。併用もできます。

PolarsやDuckDBは、処理によってはメモリに収まらないデータも扱えます。ただし中間結果や演算の種類に依存します。なお、Polarsのstreamingは主に小分けで処理する仕組みで、後述の常時イベント処理とは意味が異なります。

GUI中心なら、Power Query、KNIME、Alteryxも候補です。

出典:pandas/Polars/Polarsのstreaming/DuckDB/Power Query

2. DWH内でSQLを使って加工する

Snowflake・BigQuery・Redshiftは、保存だけでなく、結合・集計・変換を実行する計算機能も持っています。データが既にDWH内にあれば、その場で加工するのが有力です。

  • Snowflake:用途別に計算資源を分けながら、分析・加工・共有などを管理する基盤。
  • BigQuery:Google CloudのサーバーレスDWH。サーバー管理を抑えて大規模なSQL分析を実行できます。
  • Redshift:AWSのDWH。AWS上のデータやサービスと組み合わせる場合の候補です。

選択はデータ量だけでは決まりません。既存クラウド、接続先、権限、同時利用、実際のクエリ性能、総費用を比較します。

出典:Snowflake/BigQuery/Redshift

3. 大量データや重い処理を分散実行する

  • Apache Spark:複数台へ処理を分担させる計算エンジン。大規模な結合・集計、ログ加工、定期バッチなどに使います。PythonからはPySparkを使え、SQLも利用できます。1台でも動きます。
  • Dask:Pythonの処理を並列・分散化する仕組み。pandasに近い操作やPythonの計算を拡張したい場合に検討します。
  • Ray/Ray Data:Python・AI向けの分散処理基盤。画像・文章の前処理、GPUでのバッチ推論、学習データ準備などが代表用途です。

分散化には通信、処理の調整、障害対応の負担もあります。まず1台で要求時間内に終わるかを調べ、不足する理由を確認してから分散化を考えます。

出典:Spark/Dask/Ray Data

4. 流れ続けるイベントを加工する

クリック、決済、センサー情報などを継続的に処理する領域です。過去の状態、時間帯ごとの集計、遅れて届くデータを扱います。

  • Apache Flink:状態を持つ分散ストリーム処理エンジン。時間窓やイベントの結合などに対応します。バッチ処理も可能です。
  • Kafka Streams:Kafkaのデータ加工をJavaアプリに組み込むライブラリ。Kafka中心のイベント駆動アプリで使います。
  • Spark Structured Streaming:SparkのSQL・DataFrameの考え方で継続処理する仕組み。標準では小さなバッチを繰り返します。
  • Apache Beam:バッチ・ストリーム処理を書く統一モデルとSDK。実行にはDataflow・Flink・Sparkなどのrunnerが必要です。runner間で対応機能も異なります。

ここでは速さだけでなく、重複・遅延・順序・再実行の扱いが重要です。「exactly-once」も、外部DBやAPIへの書き込みまで無条件に保証する言葉ではありません。

出典:Flink/Kafka Streams/Structured Streaming/Beam

5. データを取り込む・変更を運ぶ

  • Airbyte:DBやSaaSからデータを取り込むコネクタ基盤。
  • Fivetran:データ転送・同期を管理するサービス。接続や同期の運用負担を抑えたい場合の候補です。
  • Debezium:DBの追加・更新・削除をイベントとして取り出すCDCの仕組み。
  • Kafka:イベントを保存し、複数の処理へ受け渡す基盤。イベントの再読み込みにも使います。加工はFlinkやKafka Streamsなどが担う構成が典型です。

取り込みでは、接続先の有無に加えて、削除の反映、履歴、スキーマ変更、障害からの復旧を確認します。

出典:Airbyte/Fivetran/Debezium/Kafka

6. 保存場所・ファイル形式・テーブル形式

ここは層を分けるのが特に大切です。

  • S3:データの実体を置くオブジェクトストレージ。
  • Parquet:分析向けの列指向ファイル形式。圧縮や必要な列だけの読み取りに適しています。
  • Iceberg:多数のファイルを一つのテーブルとして管理する形式。スナップショットやスキーマ変更などを扱います。
  • Delta Lake:データレイク上のテーブルにトランザクション・履歴管理などを加える仕組み。Spark・Databricksとの組み合わせが代表的です。
  • Hudi:更新・削除・増分処理やテーブル保守を扱う基盤。頻繁なupsertやCDCを重視する場合の候補です。
  • Catalog:テーブル名からメタデータを見つける管理の入口。製品によって権限や監査も担います。

つまり、S3上のParquetファイルを、Icebergのテーブルとして管理することができます。形式を選ぶ際は、使うエンジンが必要な読み書き機能に対応しているかも重要です。

HDFSは分散ファイルシステム、Cassandraは想定したキーによる読み書きを分散して行うNoSQLデータベースです。Sparkはそれらからデータを読んで計算する側、と捉えると整理できます。

出典:S3/Parquet/Iceberg/Delta Lake/Hudi

7. 保存されたデータを問い合わせる・提供する

  • Trino:データレイクや複数DBにSQLを実行する分散クエリエンジン。複数の保存先をまたぐ分析にも使います。接続先やネットワークが性能に影響します。
  • ClickHouse:高速な分析クエリを重視する列指向DB。ログやイベントを継続的に取り込み、素早く集計して画面・アプリから参照する用途で検討します。

DWHもこの役割を担います。単にSQLが使えるかではなく、保存管理を含む基盤なのか、既存の保存先に計算をかけるエンジンなのかを見ます。

出典:Trino/ClickHouse

8. 加工手順と実行を管理する

  • dbt:加工をモデル・依存関係・テスト・文書として管理する変換フレームワーク。SQLで分析用テーブルを積み上げる用途が中心です。基本構成では、実際のデータ計算は接続先のDWHなどが担います。
  • Airflow:実行順、スケジュール、再試行、過去期間の再処理などを管理します。
  • Dagster:テーブルやファイルなどの成果物を中心に、依存関係・更新・品質を管理する考え方が特徴です。
  • Prefect:Python関数をワークフローとして実行・監視する仕組みです。

例えばAirflowからSparkやdbtを起動できます。加工処理そのものと、それをいつ・どの順番で動かすかは別の役割です。

品質管理には、dbtのテストやGreat Expectationsなどを使い、欠損・重複・型・件数などを検査します。

出典:dbt/Airflow/Dagster/Prefect/Great Expectations

9. 複数の機能をまとめたサービス

個別に組むほか、統合基盤を使う選択もあります。

  • Databricks:Spark系加工、SQL分析、AI、テーブル管理、実行管理などを統合。
  • AWS Glue:取り込み・加工・Catalogなどを提供するデータ統合サービス。
  • Google Cloud Dataflow:Beamなどによるバッチ・ストリーム処理を実行するマネージドサービス。
  • Microsoft Fabric:データ統合、Spark、DWH、リアルタイム分析、BIなどを統合。

「どの計算エンジンか」と「運用をどこまで任せるか」は別の選択軸です。

出典:Databricks/Glue/Dataflow/Fabric

10. 実際の選び方

次の順で考えると、製品名から入るより選びやすくなります。

  1. 何を出したいか:日次レポート、即時応答、機械学習の前処理など。
  2. いつまでに必要か:翌朝、数分後、イベント到着後すぐ。
  3. データはどこにあるか:既存の保存先で処理できれば、移動を減らせます。
  4. SQL・Python・GUIのどれが合うか:担当者、既存資産、特殊処理から判断します。
  5. 1台で終わるか:元データ量だけでなく、結合後の膨張、中間結果、メモリ、I/Oを測ります。「何GB以上ならSpark」という固定境界はありません。
  6. 運用・費用・管理はどうか:計算料、保存料、転送料、運用工数、権限、監査、互換性まで含めます。

11. ETL・ELTと生データ保持

  • ETL:Extract → Transform → Load。目的の保存先へ入れる前に加工。
  • ELT:Extract → Load → Transform。目的の基盤に入れてから加工。

生データを残すかは、これとは別の設計です。ETLでもS3などに加工前データを残せますし、ELTでも履歴をどこまで残すかは決める必要があります。

実際には、機密情報を除いてから取り込み、保存後に集計するなど、両方が混ざることもあります。raw領域にも保存期限とアクセス権が必要です。

12. よくある組み合わせ

  • 個人の分析:CSV/Parquet → pandas・Polars・DuckDB
  • 会社の定期レポート:Airbyte/Fivetran → DWH → dbtで整形
  • 大規模データレイク:S3など → Sparkで加工 → Iceberg等で管理 → TrinoやDWHで分析
  • リアルタイム分析:Debezium → Kafka → Flink → ClickHouseなど

全製品をそろえる必要はありません。

理解の軸としては、まず「保存する・計算する・変更を運ぶ・実行を管理する」の4つを押さえるのがおすすめです。DDIAの話ともつながり、製品が入れ替わっても整理しやすくなります。

検索結果がある場合は、上下矢印キーで候補を移動し、Enterキーでメモへ移動します。Escapeキーで検索を閉じます。

キーワードを入力して検索できます。