結構頑張って進んできたな
Data Ingestionとは
Data Ingestion(取り込み)は、データをソースシステムから別のストレージ/システムへ移すこと。
つまり基本的には、
ソースシステム ↓ データ取り込み(ingestion) ↓ ストレージ / 最終地点
という A地点からB地点へのData Movement を指します。
- データ取り込みとデータ統合(integration)は違う
データ取り込み → データを移動する
データ統合 → 複数ソースのデータを組み合わせて、新しいデータセットを作る
例えば、
CRM 広告Data Web Analytics ↓ 組み合わせる ↓ 顧客データ
はデータ統合。
CRMからDWHへデータを運ぶ部分はデータ取り込みです。
- Internal Ingestion(内部統合)はここでは別扱い
同じシステム内部で、
Table A → Table B ストレーム→ キャッシュ のようにデータを移すこともあるが、この本ではそれを主にデータ変換の一部 として扱う。
Chapter 7のデータ取り込みは、主にソースシステムからデータエンジニアリング側へ持ってくる部分。
*データパイプラインはデータ取り込みより広い概念
データ取り込みはパイプラインの一部。
データパイプラインは、
アーキテクチャ + システム + プロセス
を組み合わせて、データをデータエンジニアリングライフサイクル全体に流す仕組み。
つまり、
Source ↓ Ingestion ↓ Storage ↓ Transformation ↓ Analytics / ML / Serving
全体がデータパイプライン。
- ETL / ELT / Reverse ETLも全部パイプラインのパターン
ETLやELTを別物として強く区別しすぎるのではなく、現代では
「目的に応じて適切なツールとパターンを組み合わせる」
ことが重要。
Reverse ETLやデータシェアリングも、広い意味ではデータパイプラインの一部として扱える。
- モダンデータパイプラインは柔軟であるべき
昔はMonolithic ETLのように、 決まったシステム・決まったフロー が中心だった。
今はクラウドサービスやツールをLEGOのように組み合わせ、
100 Sourceから取得 ↓ 20 Tableに統合 ↓ ML ModelをTraining ↓ ProductionへDeploy ↓ Monitor
のような複雑なワークフローもデータパイプラインに含まれる。
ingestの軸
まずユースケースと目的地を決める
Ingestionでは最初に、 何のためのデータか / どこへ送るか / どれくらいの頻度で更新するか / どれくらいの量か を考える。
さらに、 フォーマット / データの質 / 下流との互換性 / ストリーミング時の前処理 も確認する。
つまり「とりあえず全部取り込む」のではなく、下流でどう使うかから逆算してIngestionを設計する。
- Bounded vs Unbounded Data
Unbounded Data → 終わりのない継続的なEvent Stream。
例:
Web click IoT sensor Orders Logs
Bounded Data → 時間などの境界で区切られた有限データセット。
例:
2026-09-17の注文 1時間分のLog
本文の重要な考え方は、
“All data is unbounded until it’s bounded.”
実世界ではデータは継続的に発生していて、Batchとはそれを人工的に区切っているだけ、という発想。
Frequency:Batch / Micro-batch / Streaming
Ingestion頻度は連続的なスペクトラム。
Batch → 1日1回、1時間1回などまとめて処理
Micro-batch → 数秒〜数分程度の小さなBatch
Streaming / Near Real-time → イベントが来たらほぼ即時処理
完全なReal-timeは存在せず必ずレイテンシがあるので、正確には Near Real-time。
またStreamingで取り込んでも、下流がBatchなら、そこがパイプライン全体のレイテンシボトルネックになる。
- Synchronous vs Asynchronous Ingestion
Synchronous
A完了 ↓ B開始 ↓ C開始
各ステージが強く依存している。
一部がfailすると下流が止まり、場合によっては最初から再実行になる。
Asynchronous
イベント ↓ キュー / ストリーム ↓ 各ステージが独立して処理
データが到着したものから並列で処理できる。
キューやストリームが バッファ / ショック吸収 になり、Burst Traffic(一時的な大量データ)も吸収できる。
本文では基本的に、強いSynchronous Couplingを避けてAsyncに寄せる方が柔軟という考え方。
- Serialization(シリアライゼーション) / Deserialization(デシリアライゼーション)
Ingestionではソースデータをネットワークやストレージへ送れる形式にシリアライズ し、データ受け取り側で デシリアライズ する。
重要なのは、 データ受け取り側がそのフォーマットをちゃんと読めるか を確認すること。
データは届いたのにデシリアライズできず使えない、という状態は避ける。
- スループット / スケーラビリティ
Ingestion Systemがどれだけのデータ量を処理できるか。
通常時だけでなく、 Backfill / Burst / Source復旧後の大量流入 に耐えられるかを見る。
例えばソースDBが1時間落ちて、その後溜まったデータが一気に流れてくる場合、
通常 1,000 events/s ↓ 復旧後 10,000 events/s
のような状況が起こる。
そのため バッファ + 水平スケーリング が重要。
可能ならマネージドサービスでオートスケーリングさせる。
- 信頼性 vs 耐久性
信頼性 → Ingestion System自体が安定して動くか
耐久性 → データを失わないか
例えばIoT Deviceがイベントを再送できない場合、Ingestion Systemが落ちるとデータそのものが永久に失われる。
そのためIngestionの信頼性が、そのままデータの耐久性に直結する。
ただしMulti-AZ / Multi-Region / 24時間On-callなどを増やすほどコストも上がるので、どこまで守るかはトレードオフ。
- ペイロード
ペイロードとは 実際に運ぶデータそのもの。
主に見るのは、 Kind / データの形 / サイズ / スキーマ / データの型 / メタデータ。
例えば、
Kind → Table / Image / Video / Text
Shape → Rows × Columns、JSON Nesting、画像Resolutionなど
Size → 数KBなのか数TBなのか
Schema → カラム名・データタイプ・Nested Structureなど
ペイロードの性質によって、最適なIngestion方法が変わる。
大きなPayloadはChunking(かたまりに分割)する
巨大ファイルはそのまま送るより、
huge_file ↓ chunk1 chunk2 chunk3 ↓ データ受け取り側で再構築
のように分割することがある。
ネットワーク転送しやすくなり、並列処理もしやすい。
- スキーマ変更に備える
ソース側では、 列追加 / 型変更 / テーブル追加 / 列のリネーム が普通に起こる。
Ingestionツールが自動検知・自動反映できても、下流のレポートやモデルが壊れる可能性がある。
なので、
オートメーション + アラート + 人間のコミュニケーション
が重要。
「自動で通ったから問題なし」ではない。
- Schema Registry(スキーマレジストリ)
ストリーミングでは生成者と消費者の間でメッセージスキーマを共有する必要がある。
スキーマレジストリは、 スキーマ / データ型 / バージョン / 履歴 を管理するメタデータリポジトリ。
これにより生成者と消費者で、 同じシリアライゼーション / デシリアライゼーション前提 を保ちやすくなる。
- メタデータ
ペイロード本体だけでなく、 スキーマ / ソース / 日付 / 所有者 / リネージュ などのメタデータも重要。
メタデータがないと、データレイクが「何のデータかわからないData Swamp」になりやすい。
- Push / Pull / Poll
Push → ソースがデータ受け取り側へデータを送る
Source → Destination
例:Webhook、Event Stream
Pull → データ受け取り側がソースから取りに行く
Destination → Source
例:APIからData取得
Poll → データ受け取りが定期的にソースを確認し、変更があればPullする
1分ごとに確認 ↓ 変更あり? ↓ Pull
Pollingは簡単だが、頻繁すぎると無駄なリクエストが増え、遅すぎるとレイテンシが増える。
Batch取り込みの基本
- Batch Ingestionの基本
Batch Ingestionは、データをあるまとまりで一括して取り込む方式。
区切り方は主に2つ。
Time-based → 1時間ごと、1日ごとなど時間で区切る
Size-based → 100MBごと、10万イベントごとなど量で区切る
特にストリーミングデータをオブジェクトストレージへ落とす場合、最終的にはFile/Object単位にまとめる必要があるため、Size-based Batchがよく使われる。
- Full Snapshot vs Differential / Incremental
Full Snapshot → 毎回ソースシステム全体の現在状態を取得する。
実装は単純だが、データ量・ネットワークコスト・ストレージコストが大きくなりやすい。
Differential / Incremental → 前回以降に追加・変更されたデータだけ取得する。
効率はよいが、どこまで取得済みかを正しく管理する必要がある。
つまり、
Snapshot = Simpleだが重い Incremental = Efficientだが管理が複雑
- File-based Export / Ingestion
ソースDBへ直接接続せず、ソース側でデータをファイル化して渡す方式。
ソースDB ↓ export CSV / Parquetなど ↓ S3 / SFTP / SCP ↓ Destination
ソース側が「何をExportするか」をControlできるため、セキュリティやプロダクトDBへの負荷管理でメリットがある。
本文では、ソース側からファイルを用意して送るので Push Pattern として扱っている。
- ETL vs ELT
共通するのは、
Extract → ソースからデータを取得
Load → データ受け取り側へデータを保存
違いはTransformationのタイミング。
ETL Extract → Transform → Load
ELT Extract → Load → Transform
Load時には、Destinationのスキーマや性能特性を意識する必要がある。
- バッチサイズはストレージ特性に合わせる
バッチシステムでは、小さいWriteを大量に行うと性能が悪化することがある。
特に列指向データベースでは、
1 row insert 1 row insert 1 row insert ...
のような処理は、小さいFile/Objectを大量に作ってしまい非効率。
小さいUpdateを大量に行うのも、既存列ファイルのスキャンが必要になって重くなりやすい。
そのため、 ストレージエンジンに適したバッチサイズとWrite Patternを理解することが重要。
- ツールによって得意なWrite Patternが違う
例えば本文では、
Druid / Pinot → High Insert Rateに強い
SingleStore → OLTP + OLAPのHybrid
BigQuery → SQLで大量のSingle-row Insertは苦手だが、Streaming Buffer経由なら強い
といった違いがある。
つまり、同じ「Insert」でもシステムごとに正しい入れ方が違う。
- データ移行
データベースや環境を移行するときは、大量データをバルクで移す必要がある。
重要なのはデータ量だけでなく、スキーマの互換性。
SQL Server → Snowflakeのように似たデータベースでも、 データタイプやスキーマの扱いに細かな差がある。
そのため、いきなり全量を移すのではなく、サンプルデータで先にテストするのが重要。
- 移行ではデータよりパイプライン接続の移行が難しいこともある
データベースそのものを移せても、
Old DB ↑ ↓ Pipeline A Pipeline B Dashboard Application
など、Old Systemにつながっている依存関係をNew Systemへ切り替える必要がある。
そのためデータ移行では、データコピーだけでなく依存関係 / データ間の関係の移行まで考える。
- Direct Database Connection:ODBC(Open Database Connectivity,アプリケーションから異なる種類のデータベースに共通の書き方でアクセスできるようにする) / JDBC(Java Database Connectivity,)
データベースへ直接接続してクエリし、データを取り出す方法。
ODBC / JDBC はデータベースごとの差をドライバが吸収して、標準的なインタフェースで接続できる。
JDBCはJVM上で動くため携帯性が高く、Sparkなどでもよく使われる。
ただし大量データ取得では、並列クエリを増やすほどソースデータベースへの負荷も増えるので注意。
- JDBC / ODBCの限界
これらは基本的に 行形式 でデータを送るため、列指向データベースやネストデータとの相性がよくない。
そのため最近では、 Parquet / ORC / AvroへのDirect Export やREST APIなどを使う場合も多い。
実務では、
ソースデータベース ↓ JDBC リーダー ↓ オブジェクトストレージ ↓ ターゲットDWH
のように、JDBCだけで完結せず他のデータ取り込み方式と組み合わせることが多い。
- CDC(変更データキャプチャ):変更だけを取り込む
CDCはソースデータベース INSERT / UPDATE / DELETE を取得する仕組み。
バッチCDC → updated_at などで前回以降の変更行を取得
継続的CDC → データベースログを読んで、変更をイベントとして継続的に取得
継続的CDCならニアリアルタイムに近い複製やストリーミング分析ができる。
バッチCDCの弱点
updated_at だけを見るバッチCDCでは、 途中で何回変更されたかという履歴を失う。
例えば残高が1日に5回変わっても、最後の残高しか取れない。
完全な履歴が必要なら、 Insert-only設計やログベースCDC の方が向いている。
- CDCと複製
Synchronous Replication(同期レプリケーション) → プライマリーとレプリカを完全同期 → リードレプリカとして使える
Asynchronous CDC(非同期CDC) → 少し遅延してもよい代わりに疎結合
プライマリDB ↓ CDC ストリーム ├→ レプリカ ├→ オブジェクトストレージ └→ リアルタイム分析
のように複数データ出力へ流せる。
分析用途では、この柔軟性が大きい。
- API Ingestion
SaaSなどからデータを取る一般的な方法。
ただしAPIは標準化が弱く、 認証 / ページネーション / レートリミット / スキーマ / エラーハンドリング などを個別に理解する必要がある。
そこで、 Client Library / マネージドコネクタ / データシェアリング をできるだけ使って、カスタムコネクタ開発を減らす。
- マネージドデータコネクタ
データベースやAPIごとのコネクタをVendorやOSSに任せる方法。
通常は、 ソース / Destination / Credential / Sync Frequency / CDC方式 を設定するだけでデータシンクできる。
エラー時の監視やアラートも提供される。
本文の立場はかなり明確で、 コネクタ作成は無駄な労働からの解放がメインなので、可能ならマネージドサービスを使うべき というもの。
メッセージキュー/ イベントストリーミング
リアルタイム取り込みではキューやストリームを使う。
メッセージキュー → メッセージ単位で処理し、ACK後は消える
ストリーム → 順番付きログとして保持し、再読込や再処理ができる
ストリーミングでは、 Publish → Consume → Transform → Republish のようにデータフローが非線形になりやすい。
またスループット / パーティション / CPU / メモリー / オートスケーリングを考える必要がある。
- オブジェクトストレージを使ったデータトランスファー
大量ファイルをやり取りするなら、 S3 / GCS / Azure Blob のようなオブジェクトストレージが非常に相性がよい。
特徴は、 高スケーラビリティ / セキュリティ / 信頼性 / 任意ファイルフォーマット対応。
チーム間や企業間でのファイル交換にも向く。
- データベースファイルエクスポート
ソースDBからCSVやParquetなどへバルクエクスポートして取り込む方式。
大量スキャンは製品DBに負荷をかけるため、 時間帯をずらす / パーティション単位でエクスポート / 読み取り用複製DBを使う などの対策が必要。
クラウドデータウェアハウスではオブジェクトストレージへの直接エクスポートがかなり最適化されている。
- CSVは便利だが危険
CSVは非常に普及しているが、 Delimiter / Quote / Escape / Encoding / Schema が曖昧。
スキーマも内包しないため、製品では自動検出に頼りすぎない方がよい。
一方、 Parquet / Avro / ORC / Arrow / JSON はスキーマを持てて、Nested Dataも扱いやすい。
特に列指向分析ではParquet / ORC / Arrowが相性がよい。
Shell / SSH / SFTP / SCP
小規模なIngestionならシェルスクリプトでも十分実用的。
例えば、
DBからExtract ↓ ファイル変換 ↓ S3 Upload ↓ Target Load
をCLIだけで実装できる。
SSHはデータベースへのセキュアトンネルやSCPに使える。
SFTPも古いが、企業間データ移行では今でも現役。
Webhook
普通のAPIは消費者がソースへ取りに行くが、Webhookは逆。
Source ↓ HTTP POST Consumer Endpoint
なので Reverse API と呼ばれる。
実用的なアーキテクチャでは、
Webhook ↓ Lambda ↓ Kinesis ↓ Flink ↓ S3
のように、受信・Buffer・Proc。essing・ストレージを分離することが多い。
- Web UI / Webスクレイピング
APIがない場合、Web UIからFileを手動ダウンロードすることもあるが、 人依存なのでAutomateできるならAutomateするべき。
Webスクレイピングは最後の手段に近く、 DoS防止 / Rate Control / ToS / Legal Risk / HTML変更による保守負荷 を考える必要がある。
- Transfer Appliance
100TB以上の非常に大量データでは、インターネット転送より 物理ディスクをCloud Vendorへ送る方が速く安い 場合がある。
AWS Snowball / Snowmobileなど。
これは継続Ingestionではなく、One-time Migration向け。
- データシェアリング
Snowflake / BigQuery / Redshiftなどでは、データを物理コピーせずに共有できる。
厳重な意味ではIngestionではないが、 データを使える状態にする手段 として非常に重要。
ただし自分がデータを所有しているわけではないので、プロバイダにアクセスを切られると使えなくなる。
- ingestのプロセスについて
上流ステークホルダとの連携
ソースデータを作るのは多くの場合ソフトウェアエンジニアで、データエンジニアとは別チームにいる。
問題は、ソフトウェアエンジニア側がデータエンジニアを単なる下流消費者として見がちなこと。
そこでデータエンジニアは、 スキーマ変更の共有 / イベント設計 / データの質改善 /リアルタイムアーキテクチャ などを一緒に進める必要がある。
特に重要なのは、ソース側で質を改善すること。
- 下流ステークホルダをカスタマーと考える
データサイエンティストやアナリストだけでなく、 マーケティング / サプライチェーン / 管理職 などビジネスユーザーもカスタマー。
高度なストリーミングプラットフォームを作ることより、 Google Ads Reportの手動ダウンロードを自動化する 方がビジネスバリューが高い場合もある。
つまり、技術的な熟練度より 実際のビジネスへのインパクトを優先する。
- コミュニケーションが中心
上流にも下流にも共通する重要語がコミュニケーション。
スキーマ変更、データ定義、パイプライン変更などを早期に共有することで、手戻りを減らせる。
データエンジニアはシステム間だけでなく、組織間のインターフェース設計でもある。
- セキュリティ:Data in Motionを守る
Ingestionではデータをネットワーク越しに動かすため、セキュリティリスクが増える。
基本は、 VPC内はPrivate Endpoint On-prem ↔ CloudはVPN / Private Connection パブリックなインターネットを通るなら暗号化 。
ストレージ時だけでなく、Transfer中の暗号化も重要。
スキーマ変更は「厳しすぎても緩すぎてもダメ」
変更承認に半年かかるようなCommand-and-Controlでは機敏性が死ぬ。
逆にソーススキーマ変更をそのままTargetに自動反映すると、下流が壊れる。
本文ではGitのBranchに近い考え方として、 Development Tableで新スキーマを試してからメインへ反映 するような方式を示している。
つまり、スキーマにもVersioning / Staging / Controlled Promotionを持たせる。
- センシティブなデータは「そもそも取る必要があるか」を考える
Privacyで最も強い対策は、 不要なSensitive DataをIngestしないこと。
取る ↓ Encryptする
より、
不要なら最初から取らない
方が安全。
必要ならTokenization / Hashing / Maskingなどを使う。
- Touchless Production / Broken-glass
センシティブなデータを扱う製品では、人が直接データに触らない Touchless Production が理想。
Development / StagingではSyntheticやCleansed Dataを使い、ProductionへのDeploymentは自動化する。
どうしても製品データを見る必要がある場合は、 複数人承認 / スコープ限定 / 期限付きアクセス のBroken-glass Processを使う。
- EncryptionだけではPrivacy問題は解決しない
実際には多くのCloud DBはAt Rest / In Transit Encryptionを標準で備えている。
本質的な問題は、 誰がそのデータへアクセスできるか の方であることが多い。
ハッシュも単純に使えば安全とは限らず、既知のEmailをハッシュして照合できる場合もある。
つまり、セキュリティコントロールは技術を置くだけでなく攻撃シナリオまで考える。
- DataOps:Ingestionは特に監視が重要
Ingestionが止まると、 DWH / Data Lake / Report / ML Model 全部の更新が止まる。
最低限見るべきなのは、 Uptime / Latency / Data Volume / Event Rate / Event Size。
さらに、 Event Time / Ingestion Time / Process Time / Processing Time を追跡すると、どこでLatencyが起きているか分かる。
Third-party Dependencyも監視する
マネージドサービスを使えばオペレーション負荷は下がるが、そのシステムは自分ではコントロールできない。
そのため、 Outage Alert / Failover / Incident Response Plan を考える必要がある。
「マネージドだから止まらない」ではない。
- データ質のテスト
データはコードと違い、自分たちが何もデプロイしていなくても壊れる。
例えば、 Null増加 / Category追加 / Distribution変化 / Bot Traffic増加 など。
そのため、 バイナリーチェック + 統計的な監視 が必要。
特に危険なのはパイプラインが落ちることではなく、間違ったデータが正常そうに流れ続けること。
- データの質はソースで直すのが理想
下流でクリーンし続けるより、 ソフトウェアエンジニアと協力してソース側で、 Validation / Logging / Exception Handling を入れる。
データの質はデータチームだけの責任ではない。
- オーケストレーション
IngestionはData Pipelineの一番上流にあり、その後に多数のTaskが依存する。
小規模ならCronでも動くが、複雑になると脆い。
本格的には、 Task Graph / Dependency / Retry / Scheduling を扱えるOrchestratorが必要。
Ingestion完了 ↓ Transformation ↓ Aggregation ↓ Serving
のように依存関係を管理する。
- ソフトウェアエンジニア
Ingestionは外部システムとの接点が多く、かなりエンジニアリング的に負担が重い。
マネージドコネクタを使えるところは使い、 カスタムコードを書くなら、 Version Control / Code Review / Testing / CI/CD を適用する。
またソースやデータ取り込み側に強く依存したモノリスな設計を避け分離された設計にする。