Airflow中context导入错误:无法运行跳过任务告警类
1. 错误的context导入
你代码里写的from airflow.utils import context完全无效——Airflow根本不存在airflow.utils.context这个可导入的模块或对象,这是直接触发ImportError的原因。Airflow的任务上下文是在任务运行时自动传递的,不需要从这个路径导入。
2. 自定义类的使用逻辑错误
你把SkippedOperator当成Airflow的任务Operator来用,但它只是一个普通Python类,不能直接放到DAG的任务依赖链里(比如start_task >> [task2(), task3()] >> Skipped_Task)。Airflow的任务必须是继承自BaseOperator的类实例,或是用@task装饰的函数,普通类无法被Airflow识别为任务节点。
3. 上下文传递的时机错误
你在DAG定义阶段就实例化了Skipped_Task = SkippedOperator(..., context=context),但Airflow的任务上下文只有在任务运行时才会生成,DAG定义阶段根本不存在context这个变量,这种写法完全不符合Airflow的运行逻辑。
4. 类的参数接收逻辑错误
SkippedOperator的__init__用**context接收参数,但这不是Airflow传递上下文的正确方式——正规的Operator是通过execute方法来接收运行时上下文的,而不是在初始化阶段。
1. 将类改为Airflow自定义Operator
让类继承Airflow的BaseOperator,重写execute方法处理核心逻辑:
from airflow.models.baseoperator import BaseOperator from airflow.utils.session import provide_session from airflow.models import TaskInstance from airflow.utils.state import State from airflow.exceptions import AirflowException from airflow.providers.slack.operators.slack import SlackAPIPostOperator class SkippedAlertOperator(BaseOperator): def __init__(self, channels, mode, slack_conn_id, env, **kwargs): super().__init__(**kwargs) self.channels = channels self.mode = mode self.slack_conn_id = slack_conn_id self.env = env @provide_session def execute(self, context, session=None): dag_id = context["dag"].dag_id logical_date = context["logical_date"] # 查询当前DAG运行中被跳过的任务数量 skipped_count = session.query(TaskInstance).filter( TaskInstance.dag_id == dag_id, TaskInstance.execution_date == logical_date, TaskInstance.state == State.SKIPPED ).count() if skipped_count > 0: self.log.info(f"检测到{skipped_count}个被跳过的任务") # 构造Slack告警消息 circle_emoji = "red_circle" if self.mode == "failed" else "white_circle" task_mode = "任务失败" if self.mode == "failed" else "DAG成功" slack_msg = f""" :{circle_emoji}: {self.env} {task_mode}. *任务*: {self.task_id} *DAG*: {dag_id} *执行时间*: {logical_date} *跳过任务数*: {skipped_count} """ # 非生产环境发送Slack告警 if self.env != "prod": for channel in self.channels: slack_op = SlackAPIPostOperator( task_id=f"send_slack_alert_{channel}", channel=channel, slack_conn_id=self.slack_conn_id, text=slack_msg, username="Airflow告警" ) slack_op.execute(context=context) # 如果是失败模式,主动抛出异常标记DAG失败 if self.mode == "failed": raise AirflowException(f"DAG运行中有{skipped_count}个任务被跳过,触发失败告警")
2. 正确在DAG中使用自定义Operator
在DAG的任务链里实例化这个Operator,不需要手动传递context(Airflow会自动在运行时传入):
from datetime import datetime from airflow.decorators import dag from airflow.operators.dummy import DummyOperator from airflow.decorators import task from airflow.exceptions import AirflowSkipException @dag( dag_id="my_test_dag", start_date=datetime(2023, 1, 1), schedule_interval=None, catchup=False ) def generate_dag(): start_task = DummyOperator(task_id="task1") @task() def task2(): raise AirflowSkipException @task() def task3(): return # 实例化自定义告警Operator skipped_alert = SkippedAlertOperator( task_id="skipped_task_alert", channels=["my_channel"], mode="success", slack_conn_id=SLACK_CONN, env=ENV ) start_task >> [task2(), task3()] >> skipped_alert dag = generate_dag()
3. 清理无效代码
删掉所有from airflow.utils import context的语句,这个导入完全不符合Airflow的API规范。
内容的提问来源于stack exchange,提问作者jonhatan_schilino

