You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Airflow DockerOperator创建镜像失败时发送通知

Airflow DockerOperator 执行失败时发送针对性通知

要实现DockerOperator遇到指定两类错误时发送通知,只需改造on_failure_callback函数,让它能捕获并区分异常类型,再发送对应内容的通知。以下是完整实现代码及说明:

改造后的代码

from airflow import DAG
from datetime import datetime, timedelta
from airflow.providers.docker.operators.docker import DockerOperator

def send_slack(context):
    # 获取任务实例的异常信息
    exception = str(context['task_instance'].exception)
    task_id = context['task_instance'].task_id
    dag_id = context['task_instance'].dag_id
    
    # 根据异常关键词判断错误类型
    if "Cannot connect to the Docker daemon" in exception or "20.21.22.23:2375" in exception:
        error_msg = f"DAG [{dag_id}] 任务 [{task_id}] 执行失败:Docker执行服务器(20.21.22.23:2375)未运行或不可达"
    elif "10.11.12.13" in exception or "pull access denied" in exception or "connection refused" in exception:
        error_msg = f"DAG [{dag_id}] 任务 [{task_id}] 执行失败:私有Docker仓库(10.11.12.13)未运行或镜像拉取失败"
    else:
        error_msg = f"DAG [{dag_id}] 任务 [{task_id}] 执行失败:未知错误\n{exception}"
    
    # 这里替换为实际的Slack发送逻辑,比如调用Slack API
    print(error_msg)
    # 示例:slack_client.chat_postMessage(channel='#airflow-alerts', text=error_msg)

default_args = {
    'on_failure_callback': send_slack,
}

with DAG(
    dag_id='test_dag',
    default_args=default_args,
    schedule_interval='45 * * * *',
    start_date=datetime(2021, 1, 1),
    catchup=False,
    dagrun_timeout=timedelta(minutes=420),
    concurrency=1,
    tags=['test']
) as dag:

    t = DockerOperator(
        task_id="test_operator",
        container_name="test_container",
        image=f"10.11.12.13/myapp:latest",
        force_pull=False,
        auto_remove=True,
        command="python my_test.py",
        docker_url="tcp://20.21.22.23:2375",
        cpus=1,
        mem_limit="1g",
        mount_tmp_dir=False
    )

    t

if __name__ == "__main__":
    dag.cli()

关键说明

  1. 异常捕获与区分:

    • send_slack函数接收Airflow自动传入的context参数,从中提取任务实例的异常信息。
    • 通过判断异常字符串中的关键词(如Docker服务器地址、私有仓库地址、特定错误提示),区分两类目标错误,生成针对性的通知内容。
  2. 通知逻辑替换:

    • 代码中用print模拟通知发送,实际使用时需替换为Slack的API调用逻辑(比如使用slack-sdk库发送消息到指定频道)。
  3. 覆盖范围:

    • 除了目标两类错误,还添加了未知错误的兜底处理,确保所有失败场景都能收到通知。

内容的提问来源于stack exchange,提问作者S.Hashiba

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 04:05:20