如何在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()
关键说明
异常捕获与区分:
send_slack函数接收Airflow自动传入的context参数,从中提取任务实例的异常信息。- 通过判断异常字符串中的关键词(如Docker服务器地址、私有仓库地址、特定错误提示),区分两类目标错误,生成针对性的通知内容。
通知逻辑替换:
- 代码中用
print模拟通知发送,实际使用时需替换为Slack的API调用逻辑(比如使用slack-sdk库发送消息到指定频道)。
- 代码中用
覆盖范围:
- 除了目标两类错误,还添加了未知错误的兜底处理,确保所有失败场景都能收到通知。
内容的提问来源于stack exchange,提问作者S.Hashiba
相关产品推荐
相关产品推荐

