Apache Airflow + pushplus:DAG 失败推到微信
效果:任务或整条 DAG 失败时,微信收到 dag_id、task_id 和执行日期,不用一直盯着 Grid 视图。
Airflow 没有内置 pushplus 通知类型。优先用官方 Apprise Provider(URL 为 pushplus://token);也可用失败回调里的 curl,或走已有邮件通道。
前置条件
- 一个 pushplus token(官网扫码获取)
- 能改 DAG / Connection 的 Airflow 2.x+
- 选用 Apprise 时:已安装
apache-airflow-providers-apprise与apprise
配置步骤
方式一:Apprise(推荐)
- Admin → Connections → +,Connection Type 选 Apprise,URI 填:
pushplus://你的token
- 在 DAG 的
default_args或 Task 上挂on_failure_callback,用AppriseNotifier(具体类名以你安装的 Provider 版本文档为准),标题和正文用context里的dag.dag_id、task.task_id、ds。 - 应用产生的标题/正文会映射为 pushplus 的
title/content。
方式二:失败回调里 curl
把 token 放进 Admin → Variables(或 Secret Backend),不要写进 DAG 文件:
def notify_pushplus(context):
import os, json, urllib.request
body = {
"token": os.environ["PUSH_PLUS_TOKEN"],
"title": f"Airflow 失败 {context['dag'].dag_id}",
"content": f"task={context['task_instance'].task_id}\nstate={context['task_instance'].state}\nds={context['ds']}",
"template": "markdown",
}
req = urllib.request.Request(
"https://www.pushplus.plus/send",
data=json.dumps(body).encode(),
headers={"Content-Type": "application/json"},
method="POST",
)
urllib.request.urlopen(req)
default_args = {"on_failure_callback": notify_pushplus}
调度器 / Worker 必须能访问 www.pushplus.plus。也可用 BashOperator 或 HttpOperator(Admin → Connections 里建 HTTP 连接)在结束任务里 POST 同一 JSON。
已有 SMTP 时,可用 EmailOperator 发信,再在邮件侧转发;直接推微信仍建议 Apprise 或上面的 HTTP。
完整字段见 消息接口。
验证
curl -X POST "https://www.pushplus.plus/send" \
-H "Content-Type: application/json" \
-d '{"token":"你的token","title":"Airflow 配置测试","content":"token 可用","template":"markdown"}'
命令行也可:apprise -t "测试" -b "测试" "pushplus://你的token"。
常见问题
回调执行了但没消息? 看 Scheduler/Worker 日志里的 HTTP 返回,code 应为 200。再查额度与网络,见 收不到消息、接口限制。
每成功一次都推? 只挂 on_failure_callback,不要挂在 on_success_callback。