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_callback、on_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-api的Jira客户端;async_notify(context)(异步):调用await self.async_hook.create_issue(fields),其中async_hook是JiraAsyncHook,基于aiohttp直接向 Jira REST API 发起POST请求(见 jira.py)。
3.1 构造参数详解
JiraNotifier.__init__的参数签名(见 jira.py)如下:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
jira_conn_id | str | "jira_default" | 指向 Jira 实例的 Airflow Connection ID |
proxies | Any | None | 调用 Jira REST API 时使用的代理,可选 |
api_version | str/int | "2" | 使用的 Jira API 版本 |
api_root | str | "rest/api" | API 请求的根路径 |
description | str | 必填 | Issue 正文内容 |
summary | str | 必填 | Issue 标题 |
project_id | int | 必填 | 创建 Issue 所属项目的 ID |
issue_type_id | int | 必填 | Issue 类型(类别)的 ID |
labels | list[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-api的Jira客户端的封装:
- 通过
self.get_connection(jira_conn_id)读取 Airflow Connection; - 从 Connection 的
extra中解析verify(SSL 校验开关,默认True); - 用
url=conn.host、username=conn.login、password=conn.password、verify_ssl=verify、proxies、api_version、api_root构造Jira客户端; - 调用
client.create_issue(fields)完成建单。
另外,JiraHook通过get_connection_form_widgets()在 Airflow UI 的 Connection 表单中增加了一个Verify SSL布尔控件,并通过get_ui_field_behaviour()隐藏了schema与extra字段(见 jira.py)。
5.2 异步路径:JiraAsyncHook 与 aiohttp
异步async_notify()走 JiraAsyncHook,它继承自HttpAsyncHook:
- 请求方法固定为
POST,默认请求头为Content-Type: application/json与Accept: application/json; get_resource_url()将api_root、api_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 ID:
jira_default(与JiraNotifier的jira_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便于后续按标签筛选; - 充分利用模板渲染:
description、summary、project_id、issue_type_id、labels均可模板化,可把运行上下文(如 execution date、失败原因等)写入 Issue 正文,让工单自解释; - 同步与异步的选择:默认
notify()为同步调用;若你的执行环境强调非阻塞 IO(如 Triggerer 场景),可使用async_notify()异步路径,它不依赖atlassian-python-api的同步客户端,而是直接以aiohttp发 REST 请求; - 先确认 ID 再上线:
project_id、issue_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),仅供参考