Instiq
第1章 · データの取り込みと変換·v2.1.0·更新 2026/6/14·読了目安 約9分

変更要約: in-scope サービス網羅: 取り込み・移行(AppFlow/Data Exchange/API Gateway/DataSync/Transfer Family/MGN/Discovery/Snow Family)を追加

1.3オーケストレーションとストリーミング処理

この節の要点

パイプラインのオーケストレーション(Step Functions / Managed Workflows for Apache Airflow)、ストリーミング処理(Managed Service for Apache Flink)、Lambda による軽量変換を理解します。

複数の取り込み・変換ステップは順序立てて自動実行(オーケストレーション)します。また、ストリームをその場で集計する処理も重要です。

1.3.1オーケストレーションと処理

Kinesis Data Streams(耐久・再処理可・シャード・独自コンシューマー)、Data Firehose(S3/Redshift へ配信・コード不要・バッファリング)、MSK(マネージド Apache Kafka・Kafka 互換・既存 Kafka 移行)という3つのストリーミングサービスを並べた図。
ストリーミングサービスの使い分け
  • Step Functions:ステートマシンでパイプラインの順次/分岐/並列/リトライを編成する。
  • Managed Workflows for Apache Airflow(MWAA)Airflow の DAG でデータパイプラインを編成する(既存 Airflow 資産に向く)。
  • Managed Service for Apache Flink:ストリームをその場で集計/変換する(ウィンドウ集計など)。
  • Lambda:軽量な変換やイベント駆動の処理に使う(Firehose の変換にも組み込める)。
試験ポイント

「データパイプラインの編成=Step Functions または MWAA(Airflow)」「ストリームのリアルタイム集計=Managed Service for Apache Flink」「軽量変換=Lambda」 は DEA で頻出です。既存 Airflow 資産があれば MWAA を選びます。

パイプラインは「編成(オーケストレーション)」と「ストリーム処理」を分けて考えます。編成には Step Functions(ステートマシンで順次/分岐(Choice)/並列(Map)/リトライ/エラー処理を宣言的に記述・サーバーレス・多数の AWS サービスと直接連携)と MWAA(Managed Workflows for Apache Airflow)(Python の DAG で記述・既存 Airflow 資産やコミュニティ Operator を活かす)があり、軽量な単純連鎖なら Glue ワークフロー、イベント駆動の起動には EventBridge(スケジュール/イベントルール)も使えます。ストリームをその場で集計・変換するなら Managed Service for Apache Flinkウィンドウ集計・SQL/Java・イベント時間処理)、軽量な1件変換や Firehose 内変換は Lambda です。判断軸は「複雑な分岐や多サービス連携=Step Functions」「Airflow 互換/既存 DAG=MWAA」「ストリームのリアルタイム集計=Flink」「短い変換=Lambda」。失敗時の再実行・可観測性(CloudWatch ログ/メトリクス・X-Ray)も運用設計の要点です。

やりたいこと使うもの
複雑な分岐・多サービス連携の編成Step Functions
既存 Airflow DAG で編成MWAA
ストリームのリアルタイム集計Managed Service for Apache Flink
軽量な1件変換Lambda

シナリオ:毎日の取り込み→Glue 変換→Redshift ロードを、失敗時リトライ付きで自動化したい。 EventBridge のスケジュールで Step Functions を起動し、Glue ジョブ→品質チェック→Redshift COPY を順次/分岐で編成(各ステップにリトライとエラー通知)。一方、到着イベントのリアルタイム集計は Managed Service for Apache Flink、Firehose 配信前の軽い整形は Lambda で行います。

補足

Q. 複雑な分岐の編成? Step Functions。Q. 既存 Airflow? MWAA。Q. ストリームの集計? Managed Service for Apache Flink。Q. 短い変換? Lambda。Q. スケジュール起動? EventBridge。

注意

混同に注意:
編成(Step Functions/MWAA)とストリーム処理(Flink)は別物——Flink は「順番に動かす」道具ではない。
②MWAA は既存 Airflow 資産がある時に有利だが運用負荷とコストは Step Functions より高め。
③Lambda は実行時間・ペイロード上限があり大規模変換には不向き(Glue/EMR)。
④Glue ワークフローは単純連鎖向けで、複雑分岐は Step Functions。

補足

Glue にもワークフロー機能がありますが、複雑な分岐や他サービス連携を含む編成は Step Functions、Airflow 互換が要件なら MWAA が適します。

1.3.2取り込み・移行の in-scope サービス

DEA-C01 では、データの入口として多様な取り込み・移行サービスが問われます。外部データの取り込みでは、Amazon AppFlow が SaaS(Salesforce/Zendesk 等)と AWS の間でノーコードのデータ連携を行い、AWS Data Exchange はサードパーティのキュレーション済みデータ製品を購読して分析基盤へ取り込みます。アプリからの送信は Amazon API Gateway をマネージドな HTTP/REST の入口にし、スロットリングや認証を効かせて Lambda/Kinesis へ流します。 オンプレミスからの移行では使い分けが重要です。AWS DataSync は NFS/SMB ファイルを S3/EFS/FSx へ高速・増分・スケジュールで同期し、AWS Transfer Family は SFTP/FTPS/FTP のファイル授受を S3/EFS へサーバー運用なしで受け止めます。サーバー全体のリフト&シフトは AWS Application Migration Service(MGN)がブロックレベルの継続レプリケーションで担い、その前段の計画には AWS Application Discovery Service でサーバーの構成・依存・利用状況を収集します。回線が細い・ペタバイト級のデータは AWS Snow Family で暗号化物理デバイスによりオフライン輸送し、Snowball Edge ではエッジ前処理も行えます。

用途サービス使いどころ
SaaS からの取り込みAmazon AppFlowノーコードの SaaS↔AWS 連携
外部データの購読AWS Data Exchangeサードパーティのデータ製品
ファイルの高速同期AWS DataSyncNFS/SMB → S3/EFS の増分転送
SFTP の受け口AWS Transfer FamilySFTP/FTPS でサーバーなし授受
サーバー移行Application Migration Service / Discovery ServiceMGN で移行・Discovery で計画
大容量オフライン移送AWS Snow Family細い回線・PB 級

1.3.3この節のまとめ

  • 編成=Step Functions / MWAA(Airflow)
  • リアルタイム集計=Flink、軽量変換=Lambda
  • 取り込み=AppFlow/Data Exchange/API Gateway、移行=DataSync/Transfer Family/MGN/Snow Family

進捗の記録にはログインが必要です。

理解度チェック

(軽い確認用)

Q1. 既存の Apache Airflow の DAG を使ってデータパイプラインを編成したい。最も適したサービスはどれですか?

Q2. ストリーミングデータをリアルタイムにウィンドウ集計したい。最も適したサービスはどれですか?

Q3. 複数ステップのデータパイプラインを順次・分岐・リトライ付きで編成したい。最も適したサービスはどれですか?

理解度を確認第1章「データの取り込みと変換」の問題を解く