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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 04:05:35