返回
RSS AWS Big Data Blog AI 逐段翻译 发布 2026-09-30 00:17

为 Amazon MWAA 构建 LLM 驱动的 DAG 故障分析插件

DataHot 速览

Apache Airflow 已成为数据管道编排核心,但当管道扩展到数百个 DAG 并跨 AWS Glue、Amazon EMR、Amazon Athena、Amazon Redshift 等服务时,单个任务失败排查非常耗时。本文介绍如何构建自定义 Airflow 插件,集成 Amazon Bedrock(Anthropic Claude)自动分析 DAG 任务失败,并给出可操作的诊断洞察。插件部署在 Amazon MWAA 上,在 Airflow UI 中提供按需 AI 根因分析。完整源码已发布在 GitHub 示例仓库中。

为什么值得关注:数据从业者可了解如何用 LLM 增强数据管道可观测性与故障排查,降低 SLA 风险与数据工程运维负担。

本文目录 23 节
  1. 解决方案概述
  2. 工作原理
  3. 操作器感知的上下文收集
  4. 先决条件
  5. 插件设计
  6. 插件结构
  7. 插件注册
  8. 分析引擎
  9. 操作符脚本获取
  10. 提示工程
  11. 安全措施
  12. 部署插件
  13. 步骤 1:克隆代码仓库
  14. 步骤 2:打包并上传到 Amazon S3
  15. 步骤 3:更新 Amazon MWAA 环境
  16. 步骤 4:配置 Amazon Bedrock 连接
  17. 步骤 5:验证部署
  18. 测试解决方案
  19. 成本考虑
  20. 最佳实践
  21. 扩展解决方案
  22. 清理
  23. 结论

译文

AI 逐段翻译

Apache Airflow已成为各行业数据管道的编排骨干。但随着这些管道增长到跨服务(如AWS Glue、Amazon EMR、Amazon Athena和Amazon Redshift)的数百个有向无环图(DAG),调试单个任务失败就变成了一项重大的运维挑战。当任务失败时,数据工程师需要筛选日志、交叉比对DAG配置、分析错误消息以找到根因,这会延迟管道的服务级别协议(SLA)并影响团队生产力。

在这篇文章中,我们将向您展示如何构建一个自定义 Apache Airflow 插件,与Amazon Bedrock集成,以自动分析 DAG 任务失败并提供可操作的诊断洞察。该插件部署到Amazon Managed Workflows for Apache Airflow(Amazon MWAA),并按需提供 AI 驱动的根因分析。

此解决方案的完整源代码可在sample-aws-mwaa-llm-powered-plugin GitHub 仓库中获取。克隆该仓库,并跟随我们在本文中解释设计决策。

解决方案概述

Apache Airflow是一个被广泛采用的开源平台,用于以编程方式编写、调度和监控复杂的数据管道。团队使用 Airflow 跨行业编排提取、转换和加载(ETL)流程、机器学习工作流和数据湖管理。

Amazon MWAA 是一项托管服务,让您可以直接在 AWS 上运行 Apache Airflow,而无需承担管理底层基础设施的运维负担。使用 Amazon MWAA,您可以专注于编写工作流和业务逻辑,而 AWS 负责预置、修补、扩展和保护您的 Airflow 环境。

该解决方案使用以下 AWS 服务:

该插件直接在您的 Airflow UI 中添加了一个分析视图。概括来说,当任务失败并且您触发分析时,该插件会自动执行以下操作:

  1. 从 Airflow 元数据数据库中检索失败的任务实例元数据。
  2. 收集全面的上下文,包括任务日志、DAG 源代码和特定于操作器的脚本。
  3. 将丰富后的上下文发送到 Amazon Bedrock 进行分析。
  4. 返回结构化诊断报告,包括根因识别、分步解决方案和预防建议。

工作原理

前面的四个步骤都发生在单个 Analyze Task 操作背后。下面的图表和流程展示了高层架构以及插件如何执行这些步骤。

任务分析器插件的架构,连接 Amazon MWAA 上的 Airflow UI 与 Amazon Bedrock 和 Amazon S3

图 1:Amazon MWAA 上 LLM 驱动的任务分析器插件的高层架构

该插件遵循多步骤分析流程:

  1. 用户触发分析 – 在 Airflow UI 中,您选择失败的任务,然后选择Analyze Task。
  2. 上下文收集 – 插件从 Airflow 元数据数据库和 Amazon S3 中检索任务元数据、执行日志和 DAG 源代码。
  3. 操作器感知的丰富化 – 根据操作器类型,插件获取实际失败的代码或查询(例如,来自 AWS Glue 的 PySpark 脚本或来自 Amazon Athena 的 SQL 查询)。
  4. 基础模型分析 – 丰富后的上下文被发送到 Amazon Bedrock,后者返回结构化诊断报告。
  5. 结果展示 – 分析结果在 Airflow UI 中显示,并附有可操作的建议。

所有 AWS API 调用(Amazon Bedrock、Amazon S3 和 AWS Glue)都通过aws_default Airflow 连接进行身份验证。在 Amazon MWAA 上,默认情况下此连接没有静态凭证,因此boto3会回退到环境的执行角色。这意味着无需管理或轮换密钥。如果您需要使用不同的身份调用 Amazon Bedrock 或获取脚本,可以在aws_default连接中提供这些凭证。这可以是专用的 IAM 角色或跨账户主体,用于替代执行角色。

操作器感知的上下文收集

此解决方案的一个关键差异化优势是它能够理解不同的 Airflow 操作器类型,并自动获取相关的代码或查询。与通用日志分析器不同,该插件检索的是实际失败的代码,而不仅仅是错误消息。

下表总结了该插件针对每种操作器类型获取的内容:

操作器类型插件获取的内容来源
GlueJobOperatorPySpark 或 Python 脚本Amazon S3(来自 AWS Glue 作业定义)
EmrAddStepsOperatorSpark 或 Python 脚本Amazon S3(来自步骤参数)
EmrServerlessStartJobOperatorSpark 脚本Amazon S3(来自作业驱动程序)
AthenaOperatorSQL 查询内联(来自操作器参数)
RedshiftDataOperatorSQL 查询内联(来自操作器参数)
BashOperatorBash 命令内联(来自操作器参数)
PythonOperatorPython 函数DAG 源代码

这种方法意味着基础模型可以分析实际失败的逻辑,将错误消息与代码中的特定行关联起来,以实现精确的根因识别。

先决条件

在开始之前,请确保您具备以下条件:

  • 一个运行Apache Airflow 3.x 的 Amazon MWAA 环境(本演练使用 Airflow 3.2)。该插件通过 Airflow 3.x 中引入的基于 FastAPI 的插件接口(fastapi_apps)注册其 UI。有关设置说明,请参阅Amazon MWAA 入门。
  • 在您的 AWS 区域中启用 Anthropic Claude 模型系列后对 Amazon Bedrock 的访问权限。本演练使用 Anthropic Claude,但您可以通过修改prompts.py中的提示负载格式,使插件适配 Amazon Nova 或其他基础模型。请参阅模型访问。
  • 一个AWS Identity and Access Management(IAM)执行角色,用于 Amazon MWAA,具有bedrock:InvokeModel和s3:GetObject权限。
  • 一个为您的 Amazon MWAA 环境提供支持的 Amazon S3 存储桶,并已启用存储桶版本控制。请参阅为 Amazon MWAA 创建 Amazon S3 存储桶。
  • 本地安装 Python 3.10 或更高版本。
  • 已配置适当权限的AWS 命令行界面(AWS CLI)。

注意:在大多数区域,您通过推理配置文件 ID(例如 us.anthropic.claude-sonnet-4-5-20250929-v1:0)而不是裸的按需模型 ID 来调用 Claude。运行aws bedrock list-inference-profiles以在配置模型之前确认其处于 ACTIVE 状态。

插件设计

在本节中,我们解释插件设计及其关键组件。下一节将引导您将其部署到您的 Amazon MWAA 环境。

插件结构

该插件遵循标准的Apache Airflow 插件架构。代码仓库的组织结构如下:

plugins/
├── task_analyzer_plugin.py    # Main plugin: FastAPI app, endpoints, registration
└── task_analyzer/
    ├── __init__.py
    ├── prompts.py             # Bedrock model configuration and prompt templates
    ├── script_utils.py        # Operator-specific script fetching logic
    ├── templates/
    │   └── index.html
    └── static/
        ├── css/
        │   └── styles.css
        └── js/
            ├── app.jsx
            ├── components.jsx
            ├── config.js
            ├── template.jsx
            └── utils.jsx

该代码仓库还包括示例 DAG,用于模拟不同操作符类型下的各种失败场景。

插件注册

在 Apache Airflow 3.x 中,插件的 Web 组件作为FastAPI应用程序通过fastapi_apps属性注册。在task_analyzer_plugin.py中,TaskAnalyzerPlugin类将 FastAPI 应用注册在/task-analyzer下,并向任务实例页面添加一个视图:

class TaskAnalyzerPlugin(AirflowPlugin):
    name = "task_analyzer_plugin"

    fastapi_apps = [
        {
            "app": app,
            "url_prefix": "/task-analyzer",
            "name": "Task Analyzer",
        }
    ]

    external_views = [
        {
            "name": "Analyze Task",
            "href": "/task-analyzer/",
            "url_route": "task_analyzer_view",
            "destination": "task_instance",
        }
    ]

Airflow 会自动发现 plugins 文件夹中的任何AirflowPlugin子类。无需注册调用或配置更改。在 Amazon MWAA 上,该文件通过plugins.zip提供,并解压到/usr/local/airflow/plugins/。

分析引擎

分析引擎是POST /api/analyze-task端点,位于task_analyzer_plugin.py中。当您触发分析时,该端点执行以下步骤:

  1. 从aws_default Airflow 连接检索 AWS 凭证。要覆盖此设置,请在 Airflow UI(aws_default连接,位于管理 > 连接)。
  2. 从请求(任务元数据、日志、DAG 源代码)组装上下文词典。
  3. 通过fetch_and_add_operator_script使用特定于操作符的脚本丰富上下文。
  4. 使用prompts.py中的模板构建提示。
  5. 调用 Amazon Bedrock 并返回结构化分析。

操作符脚本获取

位于process_operator_script中的script_utils.py函数根据操作符类型路由脚本检索:

  • 外部脚本(AWS Glue、Amazon EMR) – 插件调用AWS Glue API查找作业定义,然后从 Amazon S3 读取 PySpark 脚本。Amazon EMR 处理程序遵循相同的模式,从步骤配置或作业驱动程序中提取脚本路径。
  • 内联脚本(Amazon Athena、Amazon Redshift、BashOperator、PythonOperator、DBTOperator) – 插件直接从任务的渲染模板字段读取查询或命令,无需外部 API 调用。

该插件实现了智能获取:对于外部脚本,仅当错误消息包含代码相关模式(如 SyntaxError、TypeError 或数据类型不匹配)时,才进行 Amazon S3 API 调用。像超时这类基础设施错误会完全跳过脚本获取,从而最大限度地减少不必要的 API 调用。

提示工程

位于prompts.py中的提示模板为基础模型提供:

  • 任务元数据(DAG ID、任务 ID、运行 ID、状态)。
  • 错误消息和执行日志。
  • DAG 源代码。
  • 特定于操作符的脚本(当可用时)。

该模型生成结构化诊断报告,包含根本原因识别、分步解决方案和预防建议。模型 ID 可通过 Airflow Variables 配置,因此您可以在 Claude Sonnet 和 Claude Opus 之间切换,而无需重新部署插件。

安全措施

在将内容发送到 Amazon Bedrock 之前,该插件应用以下保护措施:

  • 凭证脱敏 – sanitize_script函数从脚本和日志中移除敏感模式(密码、令牌、访问密钥)。
  • 内容截断 – truncate_script函数限制内容大小,以保持在模型上下文窗口内。
  • 路径遍历防护 – read_allowlisted_file函数在读取任何文件之前解析规范路径并验证它们位于允许的基础目录内。

有关完整实现,请参阅script_utils.py。

可选:PII 检测和脱敏。内置的sanitize_script函数针对凭证模式。如果您的日志或脚本可能包含个人身份信息(PII),请考虑在调用 Amazon Bedrock 之前使用Amazon Comprehend添加检测步骤。DetectPiiEntities API 返回实体类型(例如姓名、电子邮件地址或账号)及其字符偏移量。您可以使用这些偏移量在上下文离开您的环境之前对跨度进行掩码或混淆。这会在每次分析中增加一次 API 调用和成本,因此请在您的合规要求需要的地方添加它。有关指导,请参阅检测 PII 实体。

部署插件

按照以下步骤将插件部署到您的 Amazon MWAA 环境。

步骤 1:克隆代码仓库

git clone https://github.com/aws-samples/sample-aws-mwaa-llm-powered-plugin.git
cd sample-aws-mwaa-llm-powered-plugin

步骤 2:打包并上传到 Amazon S3

从plugins.zip目录创建plugins/归档文件并将其上传到您的 Amazon MWAA S3 存储桶:

cd plugins
zip -r ../plugins.zip .
cd ..

aws s3 cp plugins.zip s3://<amzn-s3-demo-bucket>/plugins.zip

aws s3api head-object \
  --bucket <amzn-s3-demo-bucket> \
  --key plugins.zip \
  --query VersionId --output text

记下返回的 VersionId。下一步需要用到它。

注意:此插件仅需要fastapi和 Boto3,两者都已在 Amazon MWAA for Airflow 3.x 上预装。您不需要requirements.txt文件。跳过 requirements 文件可避免包解析冲突,这是 Amazon MWAA 环境更新失败的常见原因。

步骤 3:更新 Amazon MWAA 环境

更新您的环境以使用新的插件归档:

aws mwaa update-environment \
  --name <your-environment-name> \
  --plugins-s3-path plugins.zip \
  --plugins-s3-object-version <version-id-from-step-2>

环境会自动重启。此过程通常需要 10–30 分钟。使用以下命令监控状态:

aws mwaa get-environment \
  --name <your-environment-name> \
  --query "Environment.{Status:Status,Plugins:PluginsS3Path}" --output json

步骤 4:配置 Amazon Bedrock 连接

在 Amazon MWAA 上,aws_default连接默认存在,并解析为环境的执行角色。在大多数情况下,无需任何操作。

要覆盖区域,请编辑 Airflow UI(aws_default连接,位于管理 > 连接)中的Extra字段,将其设置为:

{"region_name": "us-east-1"}

将登录名和密码留空,以便使用执行角色。

步骤 5:验证部署

环境完成更新后,导航到管理 > 插件在 Airflow UI 中。验证 task_analyzer_plugin 出现在列表中。Analyze Task 条目现在可从任何任务实例视图访问。

测试解决方案

该仓库包含 示例 DAG,用于模拟不同操作符类型的失败场景。要验证部署:

aws s3 cp dags/ s3://<amzn-s3-demo-bucket>/dags/ --recursive
  1. 将 dags/ 目录内容复制到你的 Amazon MWAA S3 存储桶的 DAGs 文件夹:
  2. 等待 Amazon MWAA 同步 DAG(通常为 1–2 分钟)。
  3. 在 Airflow UI 中,触发其中一个测试 DAG(例如,test_aws_sql_operators)并让预期的失败发生。
  4. 导航到失败的任务实例。
  5. 选择 Analyze Task 在任务实例视图中。
  6. 通过文件和行号引用识别根本原因。
  7. 带有代码示例的分步解决方案。
  8. 预防建议和监控建议。

分析通常在 5–10 秒内完成。

成本考虑

此解决方案的主要成本驱动因素是 Amazon Bedrock 推理,它按每次分析消耗的输入和输出令牌数量计费。输入令牌来自发送到模型的任务日志、DAG 源和操作符脚本。输出令牌来自模型返回的诊断报告。更大的日志和脚本会增加输入令牌,而你选择的模型会影响每令牌费率。有关当前的每模型费率,请参阅 Amazon Bedrock 定价。

为帮助控制成本,该插件包含一种缓存机制,以错误上下文的哈希为键存储结果。对同一失败模式的重复分析会返回缓存结果,而不会再次调用 Amazon Bedrock。

最佳实践

在生产环境中部署此解决方案时,请考虑以下事项:

  • IAM 最小权限 – 仅授予 bedrock:InvokeModel 为你所选的模型 ID,并将 s3:GetObject 限定到你的操作符脚本所在的特定存储桶路径。有关指导,请参阅 Amazon MWAA 执行角色。
  • 数据清理 – 该插件在将数据发送到 Amazon Bedrock 之前会编辑凭据并截断内容。将配置值存储在 AWS Secrets Manager 中,而不是将其硬编码在 DAG 源文件中。
  • 访问控制 – 该插件的端点受 Airflow 内置身份验证保护。有关大规模的 DAG 级访问管理,请参阅 Amazon MWAA 中基于标签的自动化 DAG 权限管理。
  • 运维韧性 – 在 Amazon Bedrock API 调用周围添加重试逻辑和断路器模式。使用 Amazon CloudWatch 来监控插件性能并针对失败率设置告警。

扩展解决方案

你可以通过以下方式扩展此解决方案:

  • 主动通知 – 与 Amazon Simple Notification Service(Amazon SNS)或 Slack 集成,在发生失败时自动提供分析。
  • 知识库集成 – 使用 Amazon Bedrock Knowledge Bases 构建过去分析的知识库,以实现由检索增强生成(RAG)驱动的建议,从你组织的历史失败中学习。
  • 附加操作符支持 – 为你的组织特定的自定义操作符添加处理程序,例如专有数据连接器或内部平台集成。
  • 自动化修复 – 对于充分理解的失败模式,触发自动修复,例如使用调整后的资源配置重启任务。

清理

要从你的环境中移除该插件:

aws s3 rm s3://<amzn-s3-demo-bucket>/plugins.zip
  1. 从 Amazon S3 删除插件归档文件:
  2. 更新你的 Amazon MWAA 环境以移除插件引用,然后等待环境重启。
  3. 可选地,如果不再需要,从你的执行角色中移除 Amazon Bedrock 权限。

结论

在这篇文章中,我们向你展示了如何使用 Amazon Bedrock 为 Amazon MWAA 部署一个由 LLM 驱动的 DAG 失败分析插件。操作符感知的上下文收集使此方法有别于通用日志分析器。通过从 AWS Glue、Amazon EMR 和其他服务获取实际代码,基础模型可提供精确、可操作且带有特定行引用的建议。

要开始使用,请克隆 sample-aws-mwaa-llm-powered-plugin 仓库,将其部署到开发 Amazon MWAA 环境,并使用包含的示例 DAG 进行测试。随着你的团队对分析质量建立信心,将其推广到生产环境,在那里它可作为任何管道故障的第一线调查。

这篇内容对你有用吗?

反馈只用于改善内容筛选,不等同于收藏

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