返回
RSS AWS Big Data Blog 发布 2026-08-07 00:16 25

Amazon MWAA与Airflow 3.0实现事件驱动跨账户编排

AWS大数据博客介绍了如何在Amazon MWAA上使用Airflow 3.0构建事件驱动的跨账户管道编排。传统跨环境协调依赖轮询或传感器,存在延迟和可靠性问题。通过Airflow 3.0的基于资产的调度和Asset Watchers,结合Amazon SQS作为消息代理,可实现近实时触发,将编排延迟从分钟级降至秒级,并减少worker资源消耗。
推荐理由:数据工程师可了解如何在多账户MWAA环境中用Airflow 3.0和SQS替代轮询,提升管道响应效率与可靠性。
AWSAmazon

译文 AI 逐段翻译

在多个 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 Watchers 用事件驱动的触发器取代了基于轮询的传感器。这种方法将编排延迟从几分钟减少到几秒,并回收了先前由轮询传感器占用的 worker 资源。它还提高了消息可靠性,因为即使消费者环境暂时不可用,Amazon SQS 也会保留协调信号。

在本文中,您将学习如何使用 Airflow 3.0 与 Amazon SQS 集成,设计和部署基于资产的跨账户编排模式。您将了解 Asset Watchers、如何从生产者 DAG 发布资产事件,以及如何在下游 Amazon MWAA 环境中触发依赖工作流,从而构建响应迅速、解耦的跨多个账户的管道。

如果您使用 AI 编码助手来构建和部署基础设施,解决方案存储库包含一个基于 Agent Skills 标准的代理技能,该技能编码了本文中的架构和最佳实践。

解决方案概览

此解决方案演示了一个多 MWAA 编排架构,其中:

  1. 生产者 Amazon MWAA 环境(账户 A) 运行数据处理工作流,在数据集创建或更新时将资产事件发布到 Amazon SQS 队列。
  2. Amazon SQS 队列 充当消息代理,将生产者和消费者环境解耦。
  3. 消费者 Amazon MWAA 环境(账户 B) 使用 Asset Watchers 监控 Amazon SQS 队列,并在相关资产事件到达时自动触发下游 DAG。

主要优势

这种事件驱动的方法相比传统轮询具有几个优势:

  • 不再有轮询开销: 您用事件驱动的 Asset Watchers 取代连续的传感器轮询,它们在事件到达时做出响应。
  • 近乎实时的响应: 下游 DAG 在几秒内触发,而不是等待计划中的轮询间隔。
  • 独立环境: 生产者和消费者 Amazon MWAA 环境没有直接依赖,因此每个团队可以独立扩展和更新其环境,而不影响另一方。
  • 可靠的消息传递: Amazon SQS 提供持久的消息传递,即使消费者环境暂时不可用。
  • 明确的团队所有权: 您和您的团队维护自己的 Amazon MWAA 环境,同时协调复杂的跨账户工作流。
  • 更快的实施: 用自然语言描述需求,代理技能即可生成部署就绪的生产者和消费者 DAG,并内置本文中的最佳实践。

架构概述

以下架构展示了如何在 AWS 账户之间连接独立的 Amazon MWAA 环境,以便一个环境中完成的管道自动触发另一个环境中的依赖工作流,而无需直接的环境耦合或轮询开销。

生产者 Amazon MWAA 环境向 Amazon SQS 队列发布资产事件,消费者环境使用 Asset Watcher 监控该队列以触发下游 DAG

图 1:使用 Amazon SQS 在 Amazon MWAA 环境之间进行跨账户事件驱动编排

架构组件

该架构有四个主要组件。生产者 DAG 将资产定义为输出,并在任务成功完成时将事件发布到 Amazon SQS 队列。Amazon SQS 队列充当账户之间的持久消息代理,AWS Identity and Access Management (IAM) 策略授予生产者发送消息的权限和消费者接收消息的权限。在消费者方面,Asset Watcher 监控队列并在消息到达时更新资产状态,这会自动触发调度在该资产上的消费者 DAG。

前提条件

在实施此解决方案之前,您需要:

  • 两个运行 Apache Airflow 3.0 或更高版本的 Amazon MWAA 环境,可以在相同或不同的 AWS 账户中。每个环境必须启用 triggerer 组件。
  • IAM 策略的中级知识,包括跨账户角色信任关系和基于资源的策略。
  • Apache Airflow DAG 编写的中级知识,包括基于 Python 的 DAG 定义和任务运算符。
  • 基本的 Python 经验(Python 3.8 或更高版本)以阅读和调整提供的代码示例。
  • 一个配置了跨账户权限的 Amazon SQS 标准队列(参见跨账户 IAM 部分)。
  • AWS Command Line Interface (AWS CLI) 配置有权限访问两个 Amazon MWAA 环境和 Amazon SQS 队列的凭证。
  • 完成时间: 约 90 分钟(按照 GitHub 存储库中的说明)。
  • 预估成本: 运行两个 Amazon MWAA 环境和 Amazon SQS 队列将产生 AWS 费用。请参阅 Amazon MWAA 定价页面Amazon SQS 定价 页面来估算您所在区域和使用的成本。完成后请记得删除资源以避免持续费用。

实施

这篇文章包含一个GitHub 仓库,你可以用它部署本文所述的解决方案。你将遵循从设置 Amazon MWAA 环境和跨账户 Amazon SQS 队列到部署带有 Asset Watchers 的生产者和消费者 DAG 的实施步骤。本文提供代码示例,包括 DAG 文件、IAM 策略和依赖配置,仅用于演示目的。在生产环境部署前,请确保你根据具体要求和合规标准进行充分测试、安全审查和验证。

注意事项

  • Asset Watchers 作为 Airflow triggerer 中的后台进程运行,而不是在调度器中运行。在期望事件驱动的 DAG 触发之前,请验证 triggerer 在消费者 Amazon MWAA 环境中健康且正在运行。如果 triggerer 宕机,Amazon SQS 消息将在队列中累积,但不会触发下游 DAG,直到 triggerer 恢复。有关更多信息,请阅读Asset Watchers 文档
  • Amazon SQS 消息默认保留期为 4 天(可配置至 14 天)。如果消费者环境不可用时间超过保留期,消息将丢失。考虑配置死信队列以捕获处理失败的消息,并根据恢复要求调整MessageRetentionPeriod
  • 跨账户 Amazon SQS 访问需要生产者的执行角色上的 IAM 身份策略和 Amazon SQS 队列上的基于资源的策略。如果任一策略缺失或配置错误,消息传递将静默失败。有关跨账户访问模式的指导,请参阅在 AWS 上授予跨账户访问权限的四种方法
  • 将 Amazon SQS 的VisibilityTimeout设置为高于 Asset Watcher 处理消息的预期时间。如果超时过短,消息可能会被重新传递并触发重复的 DAG 运行。在调整此值时,请参阅Amazon SQS 可见性超时文档
  • 每个 Amazon MWAA 环境对 DAG 数量、triggerer 数量和并发 DAG 运行数量有限制。如果你计划扩展到多个 Asset Watchers 监听不同的 Amazon SQS 队列,请在做出设计决策前检查当前的 Amazon MWAA 配额。
  • Asset URI 必须在 Asset Watcher 定义和消费者 DAG 的schedule参数之间完全匹配。任何不匹配,即使只是大小写或尾随字符的不同,都会阻止消费者 DAG 被触发。在单个 DAG 文件中定义资产以避免不一致。
  • 将提供商包apache-airflow-providers-amazonapache-airflow-providers-common-messaging固定到与 Airflow 兼容的版本。不兼容的版本可能导致导入错误,从而阻止 triggerer 启动。使用如本文所述的约束文件以避免依赖冲突。

代理技能

AI 编码助手在拥有你的特定架构和约束的上下文时最为有用,而不仅仅是通用编程模式。代理技能(Agent Skills)最初由 Anthropic 开发,并于 2025 年 12 月作为公共标准发布,为这一需求提供了可移植的格式。SKILL.md文件编码程序性知识、最佳实践和工作流程,使兼容的 AI 编码代理能够按需发现和应用它们。该标准现在受 Kiro、Strands Agents、Anthropic Claude Code、OpenAI Codex、Cursor、Gemini CLI 和其他工具支持。此处提供的解决方案包含一个基于此标准构建的代理技能(agent-skill/),它编码了本文中的跨账户编排架构和操作最佳实践。当你告诉 AI 编码助手类似“为我的订单管道编写跨账户 Amazon MWAA DAG”时,该技能会指导代理完成完整的工作流程:

  • 收集 Amazon SQS 队列 URL。
  • 生成结构正确的生产者和消费者 DAG 文件。
  • 可选地将其部署到 Amazon MWAA 环境。

该技能不需要你预先提供 AWS 账户 ID 或 Amazon MWAA 环境名称。相反,它会通过运行aws mwaa list-environmentsaws sts get-caller-identity来自动发现你的环境,使用本地配置的 AWS CLI 凭据,然后要求你确认哪个环境是生产者,哪个是消费者。

该技能支持两种模式:

  • 示例模式:生成参考生产者和消费者 DAG,用于快速跨账户验证,仅需 Amazon SQS 队列 URL 作为输入。
  • 自定义模式:使 DAG 模板适应特定的业务逻辑。例如,生产者运行 AWS Glue 提取、转换和加载(ETL)作业,消费者触发数据构建工具(dbt)模型刷新。此模式自定义 DAG ID、任务名称、计划和处理逻辑,同时保留正确的 Asset Watcher 模式。

除了代码生成,该技能还包括自动部署流程。此流程发现现有的 Amazon MWAA 环境,运行预检(Amazon Virtual Private Cloud(Amazon VPC)网络、提供商版本、triggerer 健康和 Amazon SQS 队列可访问性),将 DAG 上传到正确的 Amazon Simple Storage Service(Amazon S3)存储桶,并验证端到端就绪性。每个修改基础设施的步骤都需要明确的用户确认。另请参阅GitHub 仓库以获取使用说明。

最佳实践

Airflow Asset Watchers 与 Amazon SQS 并非总是合适的选择。当它们合适时,它们会引入与基于传感器的轮询方法不同的操作考虑。

本节涵盖如何选择正确的跨环境编排模式、如何配置 Asset Watchers 依赖的基础设施(IAM、Amazon VPC、依赖项),以及如何设计在生产环境中可靠的生产者和消费者 DAG。

跨账户 IAM

  • 生产者执行角色需要sqs:SendMessagesqs:GetQueueUrl权限,并限定到特定队列 ARN,以避免使用sqs:*
  • Amazon SQS 队列资源策略必须允许生产者角色执行sqs:SendMessage,并允许消费者角色执行sqs:ReceiveMessagesqs:DeleteMessagesqs:GetQueueAttributessqs:GetQueueUrl
  • 在部署DAG之前,使用AWS CLI测试跨账户访问。通过Airflow任务日志调试AWS IAM比在CLI级别捕获配置错误要困难得多,也慢得多。
  • 为生产队列启用Amazon SQS服务器端加密。

触发器健康状况

  • Airflow Asset Watchers运行在triggerer中,而不是scheduler中。部署消费者DAG后,在Airflow UI中验证triggerer健康状况。
  • 即使组件出现故障,健康API也可能报告健康。通过验证Triggerer日志组是否存在Amazon CloudWatch日志流来进行交叉检查。
  • 监控airflow-<ENV>-Triggerer的CloudWatch日志,查找ClientErrorQueueDoesNotExistImportError
  • 在Amazon SQS的ApproximateNumberOfMessagesVisible和死信队列(DLQ)深度上设置Amazon CloudWatch警报,DLQ会捕获在达到最大接收尝试次数后仍处理失败的消息。
  • 使用约束文件固定提供者版本,以防止依赖冲突。

Amazon VPC网络

  • 私有子网必须将0.0.0.0/0路由到NAT网关。否则,worker和triggerer会静默失败,而web服务器看起来正常。
  • 为生产高可用性使用两个NAT网关(每个可用区一个)。
  • 对于私有路由模式,使用Amazon VPC终端节点(Amazon S3、Amazon SQS、Amazon CloudWatch Logs和Amazon Elastic Container Registry(Amazon ECR))而不是NAT。
  • 确认Scheduler、Worker、DAGProcessing和Triggerer存在Amazon CloudWatch日志流。空日志组意味着容器没有运行。
  • 安全组必须允许自引用入站流量和无限制出站。

依赖管理

  • 使用==固定提供者版本,并使用约束文件。未固定的版本在环境更新时会中断。
  • 在部署前,使用MWAA Docker镜像在本地测试依赖项。
  • 更新后检查requirements_install_ip日志流。如果创建时网络不可用,使用新的requirements-s3-object-version强制重新安装。
  • 在将包添加到requirements.txt之前,检查预安装的基础包以避免版本冲突。

选择编排模式

并非每个跨环境依赖都需要Asset Watcher。Airflow 3.0提供三种主要的编排模式:带Amazon SQS的Asset Watchers、MwaaTriggerDagRunOperator以及基于传感器的轮询,每种模式在响应时间、耦合和资源消耗方面都有不同的权衡。在实施之前,使用下表将您的用例与正确的模式匹配。

模式工作原理响应时间耦合度占用worker?适用场景
1Asset Watchers + SQS(本文)消费者的triggerer监听SQS,消息到达时触发DAG秒级松散跨账户管道。扇出。独立发布周期
2MwaaTriggerDagRunOperator生产者调用MWAA API在另一个环境中启动DAG秒级紧密是(使用wait_for_completion同账户一对一触发
3传感器(轮询)消费者定期检查条件轮询间隔中等是(除非可延迟)持久状态条件。环境内依赖
  • 避免将持久状态触发器(例如,S3KeyTrigger)接入Asset Watchers。由于条件不会清除,它们会持续触发。

DAG编写

  • 尽量减少模块级代码。DAG文件每个周期都会重新解析,繁重的导入会拖慢整个解析循环。
  • 设计任务时,使其无论运行一次还是多次都产生相同结果(称为幂等性)。重试时可能出现重复的Amazon SQS消息,因此优先使用UPSERT(插入或更新)而不是INSERT,以避免重复记录。
  • 将机密信息排除在DAG文件和消息体之外。使用Airflow Connections(aws_conn_id)代替。
  • 在上传到S3之前,使用python your_dag.py在本地测试DAG导入。
  • S3上传后留出时间进行DAG解析,或使用dags reserialize强制解析。

生产者DAG设计

  • 在Amazon SQS消息中包含dag_idrun_idlogical_date以及数据集特定上下文,以便消费者无需回调即可路由。
  • 使用SqsHook而不是原始的boto3包。它遵循aws_conn_id并与Airflow日志集成。
  • 让发布失败抛出异常,以便Airflow重试机制处理重新投递。

消费者DAG设计

  • 通过triggering_asset_events访问消息,而不是直接读取队列。Asset Watcher已经消费了Amazon SQS消息。
  • 防御性地验证消息负载。生产者可能会随时间发展其模式。
  • 对于复杂的多资产依赖,使用条件资产调度(& / |)。

清理资源

为避免持续产生AWS费用,完成后删除您为此解决方案创建的资源。GitHub仓库包含逐步清理说明,用于删除Amazon SQS队列、Amazon MWAA环境、IAM角色和策略以及Amazon S3存储桶。

参考GitHub仓库中的清理说明以移除已配置的资源。

结论

Apache Airflow 3.0中基于资产的调度,配合Asset Watchers,为您提供了一种实用方法,可在多个Amazon MWAA环境之间协调工作流,而无需轮询开销或紧密耦合。通过使用Amazon SQS作为可靠的消息代理,您可以构建响应迅速、解耦的数据管道,跨越多个Amazon MWAA环境和AWS账户,而无需传统轮询机制的运维开销。

这种方法将跨环境编排延迟从分钟级降低到秒级,用声明式基于资产的调度取代自定义传感器,并为您和您的团队提供灵活性,在协调复杂工作流的同时维护独立的Amazon MWAA环境。Amazon SQS的持久消息传递降低了信号丢失的风险,即使在临时环境中断期间也是如此。

要开始使用:

  1. 审查架构(5分钟): 打开仓库中的架构图,确认哪些Amazon MWAA环境将是生产者,哪些将是消费者。
  2. 设置 Amazon SQS 队列(15 分钟):创建跨账户 Amazon SQS 标准队列,并应用“跨账户 IAM”部分中的 IAM 身份和基于资源的策略。在继续之前,使用 AWS CLI 验证访问权限。
  3. 部署并验证 DAG 示例(30 分钟):将“实现”部分中的生产者和消费者 DAG 片段复制到 Amazon MWAA 环境中,手动触发生产者 DAG,并确认消费者 DAG 自动运行。
  4. 运行预检(20 分钟):按照“最佳实践”部分中的 Amazon VPC 网络、提供程序版本和触发器健康检查进行操作。在声明环境就绪之前,确认 Amazon CloudWatch 日志流存在于 Triggerer 日志组中。
  5. (可选)使用代理技能:如果使用 AI 编码助手,请从存储库中安装该技能,并用自然语言描述业务逻辑,以生成适合你管道的、可用于部署的 DAG。

随着您在多个账户和 AWS 区域中扩展数据操作,使用 Asset Watchers 的基于资产的调度为在 AWS 上构建现代、事件驱动的数据架构奠定了基础。从基本的生产者-消费者模式开始,随着编排需求的增长,逐步演变为复杂的多资产依赖关系。

有关更多信息,请参阅

关于作者

分享这条资讯
分享海报
保存图片
iOS 也可以长按图片保存