Apache Airflow 集成 Atlassian Jira 通知:用 JiraNotifier 在 DAG/Task 失败时自动创建 Issue
2026/9/14 4:58:07 网站建设 项目流程

Apache Airflow 集成 Atlassian Jira 通知:用 JiraNotifier 在 DAG/Task 失败时自动创建 Issue

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

本篇技术指南围绕 Apache Airflow 的apache-airflow-providers-atlassian-jiraProvider 展开,核心讲解其通知能力:通过JiraNotifier(即send_jira_notification)在 DAG 或任务触发on_failure_callback时自动在 Jira 实例中创建 Issue。读完本文,你将掌握该 Notifier 的全部配置参数、DAG 级与 Task 级两种接入写法、Jinja 模板渲染机制、同步/异步两种执行路径,以及底层 Hook 与连接配置方式,可直接落地到自己的工作流监控场景。

一、概述:Jira 通知能做什么

在 Airflow 中,DAG 与 Task 均暴露了on_*_callbacks系列回调(如on_failure_callbackon_success_callback等)。Jira Provider 提供的 JiraNotifier 正是面向这些回调设计的通知器:当任务或整个 DAG 运行失败时,它自动在 Jira 中创建一个 Issue,把运行状态转化为可追踪、可指派、可流转的工单,从而把 Airflow 的调度结果无缝接入团队已有的工单流。

该能力的官方说明位于 jira-notifier-howto-guide.rst,本文以其为骨架,结合 Provider 源码、Hook 实现与单元测试进行深度展开。

二、安装与环境要求

在使用 Jira 通知之前,需要先安装对应 Provider 包。根据 README.rst:

pip install apache-airflow-providers-atlassian-jira

其运行时依赖包括:

依赖包版本要求
apache-airflow>=2.11.0
apache-airflow-providers-common-compat>=1.10.1
apache-airflow-providers-http无显式版本
atlassian-python-api>3.41.10

其中atlassian-python-api是同步JiraHook底层封装的 Jira REST API 客户端,而异步路径则依赖apache-airflow-providers-http提供的HttpAsyncHook

三、核心组件:JiraNotifier 与 send_jira_notification

Jira 通知能力由一个类和一个便捷别名构成,均定义在 notifications/jira.py:

  • JiraNotifier:继承自airflow.providers.common.compat.sdk.BaseNotifier的通知器类,是回调真正的执行实体;
  • send_jira_notification:源码末尾一行send_jira_notification = JiraNotifier(见 jira.py)表明它只是JiraNotifier的别名,两种写法完全等价,官方文档示例中使用的是函数式别名。

JiraNotifier实现了两个执行入口:

  • notify(context)(同步):调用self.hook.get_conn().create_issue(fields),其中hook是懒加载的JiraHook(通过@cached_property缓存),最终走的是atlassian-python-apiJira客户端;
  • async_notify(context)(异步):调用await self.async_hook.create_issue(fields),其中async_hookJiraAsyncHook,基于aiohttp直接向 Jira REST API 发起POST请求(见 jira.py)。

3.1 构造参数详解

JiraNotifier.__init__的参数签名(见 jira.py)如下:

参数类型默认值说明
jira_conn_idstr"jira_default"指向 Jira 实例的 Airflow Connection ID
proxiesAnyNone调用 Jira REST API 时使用的代理,可选
api_versionstr/int"2"使用的 Jira API 版本
api_rootstr"rest/api"API 请求的根路径
descriptionstr必填Issue 正文内容
summarystr必填Issue 标题
project_idint必填创建 Issue 所属项目的 ID
issue_type_idint必填Issue 类型(类别)的 ID
labelslist[str]None应用到 Issue 上的标签,缺省时内部置为[]

值得注意的细节:

  • 模板字段template_fields = ("description", "summary", "project_id", "issue_type_id", "labels")(见 jira.py),意味着这五个字段都支持 Airflow 的 Jinja 模板渲染,可以在其中引用{{ dag.dag_id }}{{ ti.task_id }}等上下文变量;
  • 字段组装_get_fields()方法把上述参数组装成 Jira REST API 所需的 payload(见 jira.py):
{ "description": self.description, "summary": self.summary, "project": {"id": self.project_id}, "issuetype": {"id": self.issue_type_id}, "labels": self.labels, }

3.2 版本兼容处理

源码中对 Airflow 3.1+ 做了兼容分支:AIRFLOW_V_3_1_PLUS为真时把**kwargs透传给父类(3.1.0 起BaseNotifier支持接收 context 等额外参数),否则调用无参的super().__init__()(见 jira.py)。这意味着在较新版本 Airflow 上,你可以在构造 Notifier 时传入额外的上下文相关参数。

四、完整示例:DAG 级与 Task 级失败通知

官方 how-to 指南(jira-notifier-howto-guide.rst)给出了同时覆盖 DAG 级与 Task 级回调的完整代码。下面在保留原示例全部要素的基础上补充了注释:

from datetime import datetime from airflow import DAG from airflow.providers.standard.operators.bash import BashOperator from airflow.providers.atlassian.jira.notifications.jira import send_jira_notification with DAG( "test-dag", start_date=datetime(2023, 11, 3), # DAG 级失败回调:整个 DAG 运行失败时触发 on_failure_callback=[ send_jira_notification( jira_conn_id="my-jira-conn", description="Failure in the Dag {{ dag.dag_id }}", summary="Airflow Dag Issue", project_id=10000, issue_type_id=10003, labels=["airflow-dag-failure"], ) ], ): BashOperator( task_id="mytask", # Task 级失败回调:该任务失败时触发 on_failure_callback=[ send_jira_notification( jira_conn_id="my-jira-conn", description="The task {{ ti.task_id }} failed", summary="Airflow Task Issue", project_id=10000, issue_type_id=10003, labels=["airflow-task-failure"], ) ], bash_command="fail", retries=0, # 关闭重试,让任务立即失败以验证通知 )

要点解读:

  • 回调传入方式on_failure_callback接收一个可调用对象列表,因此可以同时挂多个通知器;
  • 模板变量description中直接使用了{{ dag.dag_id }}{{ ti.task_id }},这得益于前面提到的template_fields
  • retries=0:示例故意关闭重试,保证任务失败时立即触发回调,便于验证效果;
  • project_id/issue_type_id为整数:它们对应 Jira 实例中的项目与 Issue 类型的数字 ID,而非名称。

五、底层执行链路:从回调到 Jira Issue

5.1 同步路径:JiraHook 与 atlassian-python-api

同步notify()依赖 JiraHook。它是一个对atlassian-python-apiJira客户端的封装:

  1. 通过self.get_connection(jira_conn_id)读取 Airflow Connection;
  2. 从 Connection 的extra中解析verify(SSL 校验开关,默认True);
  3. url=conn.hostusername=conn.loginpassword=conn.passwordverify_ssl=verifyproxiesapi_versionapi_root构造Jira客户端;
  4. 调用client.create_issue(fields)完成建单。

另外,JiraHook通过get_connection_form_widgets()在 Airflow UI 的 Connection 表单中增加了一个Verify SSL布尔控件,并通过get_ui_field_behaviour()隐藏了schemaextra字段(见 jira.py)。

5.2 异步路径:JiraAsyncHook 与 aiohttp

异步async_notify()走 JiraAsyncHook,它继承自HttpAsyncHook

  • 请求方法固定为POST,默认请求头为Content-Type: application/jsonAccept: application/json
  • get_resource_url()api_rootapi_version、资源名拼接为 URL,例如默认情况下rest/api/2/issue
  • create_issue(fields)使用aiohttp.ClientSession(支持代理proxy),将{"fields": fields}序列化为 JSON 后 POST 到该端点。

也就是说,异步路径本质上是直接调用 Jira 的标准POST /rest/api/{version}/issueREST 接口,与同步路径殊途同归。

六、单元测试验证

Provider 的单元测试位于 tests/unit/atlassian/jira/notifications/test_jira.py,它们印证了上述全部行为:

  • test_jira_notifier:验证通过send_jira_notification创建的通知器调用JiraHook.get_conn().create_issue,且 payload 与_get_fields()组装结果一致;
  • test_jira_notifier_with_notifier_class:验证JiraNotifier类直用与别名写法行为等价;
  • test_jira_notifier_templated:验证模板字段渲染——传入description="Test operator failed for dag: {{ dag.dag_id }}."后,实际调用时的 payload 中该字段被渲染为具体的 dag_id 值;
  • test_jira_notifier_get_fields:直接断言_get_fields()的输出与预期 payload 完全一致;
  • test_jira_async_notifier:以pytest.mark.asyncio验证异步路径,确认async_notify()会调用JiraAsyncHook.create_issue并传入相同 payload。

测试中的标准 payload 形如:

{ "description": "Test operator failed", "summary": "Test Jira issue", "project": {"id": 10000}, "issuetype": {"id": 10003}, "labels": ["airflow-dag-failure"], }

七、配置 Jira Connection

通知器依赖一个已配置好的 Jira Connection。官方连接说明见 connections.rst:

  • 默认 Connection IDjira_default(与JiraNotifierjira_conn_id默认值一致,见 jira.py);
  • Host:Jira 主机地址,必须带协议 scheme(如https://your-jira.example.com);
  • Port:连接 Jira 使用的端口(按需填写);
  • Login:用于调用 Jira API 认证的用户名;
  • Password:上述用户的密码;
  • Verify SSL:连接 Jira API 时是否校验 SSL,默认True,可在 Connection 表单的额外配置中调整。

八、进阶使用建议

  • 区分 DAG 级与 Task 级回调:DAG 级on_failure_callback在 DAG 整体运行失败时触发,适合标记"整条链路异常";Task 级回调则精确到单个任务,可在description/labels中携带ti.task_id便于后续按标签筛选;
  • 充分利用模板渲染descriptionsummaryproject_idissue_type_idlabels均可模板化,可把运行上下文(如 execution date、失败原因等)写入 Issue 正文,让工单自解释;
  • 同步与异步的选择:默认notify()为同步调用;若你的执行环境强调非阻塞 IO(如 Triggerer 场景),可使用async_notify()异步路径,它不依赖atlassian-python-api的同步客户端,而是直接以aiohttp发 REST 请求;
  • 先确认 ID 再上线project_idissue_type_id是数字 ID 而非名称,建议先在 Jira 中确认目标项目与 Issue 类型的实际 ID,再写入 DAG;
  • 与重试策略配合:示例中retries=0是为了快速验证;生产环境可结合retries与失败回调,让"最终失败"才产生工单,避免噪音。

至此,你已经掌握了 JiraNotifier 从参数配置、回调挂载、模板渲染到底层 Hook 与 REST 调用链路的完整脉络,可以把它作为 Airflow 工作流失败监控的标准化组件接入团队工单体系。

【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询