Amazon Web Services ブログ
Amazon MWAA と Airflow 3.0 によるイベント駆動のパイプラインオーケストレーション
本記事は、2026 年 8 月 6 日 に公開された Event-driven pipeline orchestration with Amazon MWAA and Airflow 3.0 を翻訳したものです。翻訳はクラウドサポートエンジニアの山本が担当しました。
複数の AWS アカウントで Apache Airflow を運用しているデータエンジニアリングチームは、連携という根強い課題を抱えています。チームやビジネスユニットごとに独立した Amazon Managed Workflows for Apache Airflow (Amazon MWAA) 環境を管理している場合、環境間でワークフローを連携させる仕組みが標準では用意されていません。環境をまたぐオーケストレーションは従来、時間ベースのポーリング、複雑なカスタムセンサー、API ベースのトリガーに頼るしかなく、レイテンシーと信頼性の懸念が伴いました。Apache Airflow の Datasets 機能 (バージョン 2.4 で導入) により、単一の Amazon MWAA 環境内で有向非巡回グラフ (DAG。タスクとその実行順序を定義するワークフロー定義) のデータ認識スケジューリングが可能になりました。しかし、複数アカウントで Airflow を運用するチームには、環境間でワークフローを連携させる手段が依然としてありませんでした。
この課題を解決するのが Apache Airflow 3.0 です。このバージョンを利用できるようになったAmazon MWAA 3.0 では、ポーリングの負荷や環境間の密結合なしに、上流のイベント発生に応じて動作するクロスアカウントのイベント駆動オーケストレーションを実現します。Amazon Simple Queue Service (Amazon SQS) をメッセージブローカーとして使い、Asset Watcher がポーリングベースのセンサーをイベント駆動のトリガーに置き換えます。この方式でオーケストレーションのレイテンシーが数分から数秒に短縮され、ポーリングセンサーが占有していたワーカーリソースを解放できます。さらに、コンシューマー環境が一時的に利用できない状態でも Amazon SQS が連携シグナルを保持するため、メッセージの信頼性も向上します。
本記事では、Airflow 3.0 のアセットベーススケジューリングと Amazon SQS の連携を使って、クロスアカウントのオーケストレーションパターンを設計・デプロイする方法を説明します。Asset Watcher の仕組み、プロデューサー DAG からアセットイベントを発行する方法、下流の Amazon MWAA 環境で依存ワークフローをトリガーする方法を取り上げ、複数アカウントにまたがる応答性の高い疎結合なパイプラインを構築します。
AI コーディングアシスタントでインフラストラクチャを構築・デプロイしている場合、ソリューションのリポジトリには Agent Skills 標準に基づくエージェントスキルが含まれており、本記事のアーキテクチャとベストプラクティスを組み込んでいます。
ソリューション概要
本ソリューションでは、次のような複数の Amazon MWAA 環境によるオーケストレーションアーキテクチャを示します。
- プロデューサー Amazon MWAA 環境 (アカウント A) がデータ処理ワークフローを実行し、データセットの作成・更新時に Amazon SQS キューへアセットイベントを発行します。
- Amazon SQS キュー がメッセージブローカーとして働き、プロデューサー環境とコンシューマー環境を疎結合にします。
- コンシューマー Amazon MWAA 環境 (アカウント B) が Asset Watcher で Amazon SQS キューを監視し、関連するアセットイベントが届くと下流の DAG を自動的にトリガーします。
主なメリット
イベント駆動のアプローチには、従来のポーリングと比べていくつかの利点があります。
- ポーリングの負荷が不要: 継続的なセンサーのポーリングを、イベント到着時に反応する Asset Watcher に置き換えられます。
- ほぼリアルタイムの応答: スケジュールされたポーリング間隔を待つことなく、下流の DAG が数秒以内にトリガーされます。
- 環境の独立性: プロデューサーとコンシューマーの Amazon MWAA 環境に直接的な依存関係がないため、各チームは互いに影響を与えずに環境をスケールしたり更新したりできます。
- 確実なメッセージ配信: コンシューマー環境が一時的に利用できない場合でも、Amazon SQS が耐久性の高いメッセージ配信を担います。
- 明確なチームオーナーシップ: 自チームの Amazon MWAA 環境を維持しながら、複雑なクロスアカウントワークフローを連携できます。
- 実装の高速化: 要件を自然言語で記述すれば、エージェントスキルが本記事のベストプラクティスを組み込んだ、デプロイ可能なプロデューサー DAG とコンシューマー DAG を生成します。
アーキテクチャ概要
次のアーキテクチャでは、AWS アカウントをまたいで独立した Amazon MWAA 環境を接続します。一方の環境でパイプラインが完了すると、環境間の直接的な結合やポーリングの負荷なしに、もう一方の環境で依存ワークフローが自動的にトリガーされます。

図 1: Amazon SQS を使った Amazon MWAA 環境間のクロスアカウントイベント駆動オーケストレーション
アーキテクチャのコンポーネント
アーキテクチャは 4 つの主要コンポーネントで構成されます。プロデューサー DAG はアセットをアウトレットとして定義し、タスクが正常に完了すると Amazon SQS キューへイベントを発行します。Amazon SQS キューはアカウントをまたぐ耐久性の高いメッセージブローカーとして働き、AWS Identity and Access Management (IAM) ポリシーでプロデューサーにメッセージ送信の権限、コンシューマーに受信の権限を付与します。コンシューマー側では Asset Watcher がキューを監視し、メッセージが届くとアセットの状態を更新します。そのアセットをスケジュールに指定したコンシューマー DAG が自動的にトリガーされます。
前提条件
本ソリューションを実装する前に、次のものを準備してください。
- Apache Airflow 3.0 以降を実行する Amazon MWAA 環境 2 つ (同一の AWS アカウントでも、別々のアカウントでも構いません)。各環境で triggerer コンポーネントを有効にしておく必要があります。
- IAM ポリシーに関する中級レベルの知識 (クロスアカウントのロール信頼関係やリソースベースポリシーを含む)。
- Apache Airflow の DAG 作成に関する中級レベルの知識 (Python による DAG 定義やタスクオペレーターを含む)。
- 本記事のコードサンプルを読んで応用できる基本的な Python の経験 (Python 3.8 以降)。
- クロスアカウントの権限を設定した Amazon SQS 標準キュー (「クロスアカウント IAM」セクションを参照)。
- 両方の Amazon MWAA 環境と Amazon SQS キューにアクセスできる権限を持つ認証情報で設定した AWS Command Line Interface (AWS CLI)。
- 所要時間: 約 90 分 (GitHub リポジトリの手順に従った場合)。
- 推定コスト: 2 つの Amazon MWAA 環境と Amazon SQS キューの実行には AWS 料金が発生します。Amazon MWAA の料金ページと Amazon SQS の料金ページで、お使いのリージョンと使用量に応じたコストを見積もってください。継続的な課金を避けるため、完了後はリソースを削除してください。
実装
本記事で説明するソリューションをデプロイできる GitHub リポジトリを用意しています。Amazon MWAA 環境とクロスアカウントの Amazon SQS キューのセットアップから、Asset Watcher を使ったプロデューサー DAG とコンシューマー DAG のデプロイまで、実装手順を進めていきます。本記事で提供する DAG ファイル、IAM ポリシー、requirements の設定などのコードサンプルは、デモ目的のみを想定しています。本番環境にデプロイする前に、十分なテスト、セキュリティレビュー、固有の要件やコンプライアンス基準への適合確認を必ず実施してください。
考慮事項
- Asset Watcher は Airflow の scheduler ではなく triggerer 上のバックグラウンドプロセスとして動作します。イベント駆動で DAG がトリガーされる前提として、コンシューマー側の Amazon MWAA 環境で triggerer が正常に稼働していることを確認してください。triggerer が停止していると、Amazon SQS メッセージはキューに溜まるものの、triggerer が復旧するまで下流の DAG はトリガーされません。詳細は Asset Watcher のドキュメントを参照してください。
- Amazon SQS メッセージのデフォルトの保持期間は 4 日です (最大 14 日まで設定可能)。コンシューマー環境が保持期間を超えて利用できない状態が続くと、メッセージは失われます。処理に失敗したメッセージを捕捉するデッドレターキューの設定を検討し、復旧要件に合わせて
MessageRetentionPeriodを調整してください。 - クロスアカウントの Amazon SQS アクセスには、プロデューサーの実行ロールに付与する IAM アイデンティティポリシーと、Amazon SQS キューのリソースベースポリシーの両方が必要です。どちらかが欠けていたり、設定が誤っていたりすると、メッセージ配信はエラーを出さずに失敗します。クロスアカウントアクセスのパターンについては、Four ways to grant cross-account access on AWS を参照してください。
- Amazon SQS の
VisibilityTimeoutは、Asset Watcher がメッセージを処理するのに要する想定時間より長く設定してください。タイムアウトが短すぎると、メッセージが再配信されて DAG が重複実行される可能性があります。この値を調整する際は、Amazon SQS の可視性タイムアウトのドキュメントを確認してください。 - Amazon MWAA 環境には、DAG 数、triggerer 数、DAG の同時実行数に上限があります。複数の Asset Watcher で異なる Amazon SQS キューを監視する構成にスケールする予定がある場合は、設計を決める前に現在の Amazon MWAA のクォータを確認してください。
- アセット URI は、Asset Watcher の定義とコンシューマー DAG の
scheduleパラメータで完全に一致させる必要があります。大文字小文字や末尾の文字が異なるだけでも、コンシューマー DAG はトリガーされません。不整合を避けるため、アセットは単一の DAG ファイルで定義してください。 - プロバイダーパッケージ
apache-airflow-providers-amazonとapache-airflow-providers-common-messagingは、Airflow と互換性のあるバージョンに固定してください。互換性のないバージョンではインポートエラーが発生し、triggerer が起動しないことがあります。依存関係の競合を避けるため、本記事で説明する constraints ファイルを使用してください。
エージェントスキル
AI コーディングアシスタントは、一般的なプログラミングパターンだけでなく、対象のアーキテクチャや制約に関するコンテキストを持っているときに最も役立ちます。Anthropic が開発し、2025 年 12 月に公開標準としてリリースされた Agent Skills は、この目的に適した可搬性の高いフォーマットです。SKILL.md ファイルに手順の知識、ベストプラクティス、ワークフローを記述しておくと、対応する AI コーディングエージェントが必要に応じて検出して適用できます。この標準は現在、Kiro、Strands Agents、Anthropic Claude Code、OpenAI Codex、Cursor、Gemini CLI などのツールでサポートされています。ここで提供するソリューションには、この標準に基づくエージェントスキル (agent-skill/) が含まれ、本記事のクロスアカウントオーケストレーションアーキテクチャと運用のベストプラクティスが記述されています。「注文パイプライン用にクロスアカウントの Amazon MWAA DAG を書いて」のように AI コーディングアシスタントに指示すると、スキルがエージェントを一連のワークフローに沿って導きます。
- Amazon SQS キュー URL の収集。
- 正しく構成されたプロデューサー DAG とコンシューマー DAG のファイル生成。
- 必要に応じて Amazon MWAA 環境へのデプロイ。
このスキルでは、AWS アカウント ID や Amazon MWAA 環境名を事前に指定する必要はありません。ローカルに設定された AWS CLI 認証情報を使って aws mwaa list-environments と aws sts get-caller-identity を実行して環境を自動検出し、どちらがプロデューサーでどちらがコンシューマーかの確認を求めます。
スキルには 2 つのモードがあります。
- サンプルモード: クロスアカウントの動作をすばやく検証するためのリファレンス実装として、プロデューサー DAG とコンシューマー DAG を生成します。入力は Amazon SQS キュー URL のみです。
- カスタムモード: DAG テンプレートを固有のビジネスロジックに合わせて調整します。たとえば、プロデューサーが AWS Glue の抽出、変換、ロード (ETL) ジョブを実行し、コンシューマーが data build tool (dbt) のモデル更新をトリガーする構成です。このモードでは、正しい Asset Watcher のパターンを保ちながら、DAG ID、タスク名、スケジュール、処理ロジックをカスタマイズします。
コード生成に加えて、スキルには自動デプロイのフローも含まれます。デプロイフローは既存の Amazon MWAA 環境を検出し、事前チェック (Amazon Virtual Private Cloud (Amazon VPC) のネットワーク、プロバイダーのバージョン、triggerer の健全性、Amazon SQS キューへのアクセス可否) を実行し、正しい Amazon Simple Storage Service (Amazon S3) バケットに DAG をアップロードして、エンドツーエンドで準備が整っているかを検証します。インフラストラクチャを変更するステップでは、いずれもユーザーの明示的な確認が必要です。使い方については GitHub リポジトリも参照してください。
ベストプラクティス
Amazon SQS を使った Airflow の Asset Watcher が常に最適とは限りません。適した場面であっても、センサーベースのポーリングとは異なる運用上の考慮事項が生じます。
本セクションでは、環境間オーケストレーションのパターンをどう選ぶか、Asset Watcher が依存するインフラストラクチャ (IAM、Amazon VPC、依存関係) をどう設定するか、本番環境で信頼性の高いプロデューサー DAG とコンシューマー DAG をどう設計するかを説明します。
クロスアカウント IAM
- プロデューサーの実行ロールには
sqs:SendMessageとsqs:GetQueueUrlが必要です。sqs:*を避け、対象キューの ARN に絞って付与します。 - Amazon SQS キューのリソースポリシーでは、プロデューサーロールに
sqs:SendMessage、コンシューマーロールにsqs:ReceiveMessage、sqs:DeleteMessage、sqs:GetQueueAttributes、sqs:GetQueueUrlを許可する必要があります。 - DAG をデプロイする前に、AWS CLI でクロスアカウントアクセスをテストします。Airflow のタスクログから AWS IAM をデバッグするのは、CLI レベルで設定ミスを見つけるよりはるかに手間と時間がかかります。
- 本番環境のキューでは Amazon SQS のサーバー側暗号化を有効にします。
triggerer の健全性
- Airflow の Asset Watcher は scheduler ではなく triggerer で動作します。コンシューマー DAG をデプロイした後、Airflow UI で triggerer の健全性を確認します。
- health API は、コンポーネントが壊れていても healthy と報告することがあります。Triggerer のロググループに Amazon CloudWatch のログストリームが存在するかを併せて確認してください。
airflow-<ENV>-Triggererの CloudWatch ログでClientError、QueueDoesNotExist、ImportErrorを監視します。- Amazon SQS の
ApproximateNumberOfMessagesVisibleと、デッドレターキュー (DLQ。受信試行の上限回数を超えても処理できなかったメッセージを捕捉するキュー) の深さに Amazon CloudWatch アラームを設定します。 - 依存関係の競合を防ぐため、constraints ファイルでプロバイダーのバージョンを固定します。
Amazon VPC のネットワーク
- プライベートサブネットは 0.0.0.0/0 を NAT ゲートウェイにルーティングする必要があります。設定していないと、ウェブサーバーは正常に見えるのに、ワーカーと triggerer がエラーを出さずに失敗します。
- 本番環境の高可用性のため、NAT ゲートウェイはアベイラビリティーゾーンごとに 1 つ、合計 2 つ使用します。
- プライベートルーティングモードでは、NAT の代わりに Amazon VPC エンドポイント (Amazon S3、Amazon SQS、Amazon CloudWatch Logs、Amazon Elastic Container Registry (Amazon ECR)) を使用します。
- Scheduler、Worker、DAGProcessing、Triggerer の Amazon CloudWatch ログストリームが存在するか確認します。ロググループが空の場合、コンテナが動作していません。
- セキュリティグループは、自己参照のインバウンドトラフィックと制限のないアウトバウンドを許可する必要があります。
依存関係の管理
- プロバイダーのバージョンは == で固定し、constraints ファイルを使用します。固定していないと、環境の更新時に動作しなくなります。
- デプロイ前に MWAA Docker イメージでローカルに依存関係をテストします。
- 更新後は
requirements_install_ipのログストリームを確認します。環境の作成時にネットワークが利用できなかった場合は、新しいrequirements-s3-object-versionで再インストールを強制します。 - バージョンの競合を避けるため、
requirements.txtに追加する前にプリインストール済みのベースパッケージを確認します。
オーケストレーションパターンの選択
環境間の依存関係すべてに Asset Watcher が必要なわけではありません。Airflow 3.0 には主に 3 つのオーケストレーションパターンがあります。Amazon SQS を使った Asset Watcher、MwaaTriggerDagRunOperator、そしてセンサーベースのポーリングで、それぞれ応答時間、結合度、リソース消費のトレードオフが異なります。実装を決める前に、次の表でユースケースに合ったパターンを選んでください。
| パターン | 仕組み | 応答時間 | 結合度 | ワーカー占有 | 適した用途 |
|---|---|---|---|---|---|
| Asset Watcher + SQS (本記事) |
コンシューマーの triggerer が SQS を監視し、メッセージ到着時に DAG をトリガーする | 数秒 | 疎結合 | なし | クロスアカウントのパイプライン、ファンアウト、独立したリリースサイクル |
| MwaaTrigger DagRunOperator |
プロデューサーが MWAA API を呼び出して別環境の DAG を開始する | 数秒 | 密結合 | あり (wait_for_completion 使用時) |
同一アカウント内の 1 対 1 のトリガー |
| センサー (ポーリング) |
コンシューマーが定期的に条件を確認する | ポーリング間隔 | 中程度 | あり (deferrable でない場合) | 状態が持続する条件、環境内の依存関係 |
- 状態が持続するトリガー (
S3KeyTriggerなど) を Asset Watcher に組み込むのは避けてください。条件が解除されないため、継続的に発火してしまいます。
DAG の作成
- モジュールレベルのコードは最小限にします。DAG ファイルはサイクルごとに再パースされ、重いインポートはパースループ全体を遅くします。
- 1 回実行しても複数回実行しても同じ結果になるようにタスクを設計します (べき等性と呼ばれる性質です)。リトライ時に Amazon SQS メッセージが重複することがあるため、重複レコードを避けるには INSERT ではなく UPSERT (挿入または更新) を選びます。
- シークレットは DAG ファイルやメッセージ本文に含めません。代わりに Airflow の Connections (
aws_conn_id) を使用します。 - S3 にアップロードする前に、
python your_dag.pyでローカルに DAG のインポートをテストします。 - S3 へのアップロード後は DAG のパースが完了するまで待つか、
dags reserializeで強制します。
プロデューサー DAG の設計
- コンシューマーがプロデューサーに問い合わせずにルーティングできるよう、Amazon SQS メッセージに
dag_id、run_id、logical_date、データセット固有のコンテキストを含めます。 - 生の
boto3パッケージではなくSqsHookを使用します。aws_conn_idの設定が反映され、Airflow のロギングとも統合されます。 - 発行の失敗はそのまま例外として上げ、Airflow のリトライ機構で再配信を処理させます。
コンシューマー DAG の設計
- メッセージはキューを直接読むのではなく
triggering_asset_events経由で参照します。Amazon SQS メッセージは Asset Watcher がすでに消費しています。 - メッセージのペイロードは防御的に検証します。プロデューサーはスキーマを変更していく可能性があります。
- 複数アセットにまたがる複雑な依存関係には、条件付きのアセットスケジューリング (& / |) を使用します。
リソースのクリーンアップ
継続的な AWS 料金を避けるため、本ソリューションで作成したリソースは完了後に削除してください。GitHub リポジトリには、Amazon SQS キュー、Amazon MWAA 環境、IAM ロールとポリシー、Amazon S3 バケットを削除する手順を段階的に用意しています。
プロビジョニングしたリソースの削除については、GitHub リポジトリのクリーンアップ手順を参照してください。
まとめ
Apache Airflow 3.0 のアセットベーススケジューリングと Asset Watcher により、ポーリングの負荷や密結合なしに Amazon MWAA 環境間でワークフローを連携させる実用的な手段が手に入ります。Amazon SQS を信頼性の高いメッセージブローカーとして使うことで、従来のポーリング機構の運用負荷を伴わずに、複数の Amazon MWAA 環境と AWS アカウントにまたがる応答性の高い疎結合なデータパイプラインを構築できます。
Asset Watcher を使うことで環境間オーケストレーションのレイテンシーは数分から数秒に短縮され、カスタムセンサーは宣言的なアセットベーススケジューリングに置き換わります。複雑なワークフローを連携させながら、チームごとに独立した Amazon MWAA 環境を維持する柔軟性も得られます。Amazon SQS の耐久性のあるメッセージ配信により、環境が一時的に停止している間でもシグナルを失うリスクが下がります。
始めるには、次の手順を進めてください。
- アーキテクチャの確認 (5 分): リポジトリのアーキテクチャ図を開き、どの Amazon MWAA 環境がプロデューサーで、どれがコンシューマーになるかを確認します。
- Amazon SQS キューのセットアップ (15 分): クロスアカウントの Amazon SQS 標準キューを作成し、「クロスアカウント IAM」セクションの IAM アイデンティティポリシーとリソースベースポリシーを適用します。次に進む前に AWS CLI でアクセスを確認します。
- DAG サンプルのデプロイと検証 (30 分): 「実装」セクションのプロデューサー DAG とコンシューマー DAG のスニペットを Amazon MWAA 環境にコピーし、プロデューサー DAG を手動でトリガーして、コンシューマー DAG が自動的に実行されることを確認します。
- 事前チェックの実行 (20 分): 「ベストプラクティス」セクションの Amazon VPC ネットワーク、プロバイダーバージョン、triggerer の健全性のチェックを進めます。環境の準備完了とする前に、Triggerer のロググループに Amazon CloudWatch のログストリームが存在することを確認してください。
- 必要に応じてエージェントスキルを使う: AI コーディングアシスタントを使っている場合は、リポジトリからスキルをインストールし、ビジネスロジックを自然言語で記述して、パイプラインに合わせたデプロイ可能な DAG を生成します。
複数のアカウントや AWS リージョンにデータ運用をスケールしていくうえで、Asset Watcher によるアセットベーススケジューリングは、AWS 上でモダンなイベント駆動データアーキテクチャを構築する基盤になります。まずは基本的なプロデューサー・コンシューマーのパターンから始め、オーケストレーションの要件が増えるにつれて複数アセットにまたがる複雑な依存関係へ段階的に発展させてください。
詳細は次の資料を参照してください。
- Apache Airflow 3 on Amazon MWAA Launch Blog Post。
- Apache Airflow のドキュメント。
- Amazon MWAA ユーザーガイド。
- Amazon SQS のドキュメント。
- AWS IAM のクロスアカウントアクセス。
- Amazon MWAA の Amazon CloudWatch モニタリング。
- Airflow の外部で人による承認ステップや複雑な分岐ロジックを必要とするオーケストレーションパターンには AWS Step Functions。
- AWS Well-Architected Framework — Data Analytics Lens。
- Amazon MWAA のクォータ。
- Amazon VPC のネットワークに関するベストプラクティス。