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

使用pytest测试Airflow任务时遇Task Instance State失败问题

解决Airflow task.run测试中Task Instance State失败报错

我在对比Airflow中task.run与task.execute的差异时,编写了pytest测试代码验证prev_ds、ds等Jinja变量能否通过run方法自动渲染。测试代码断言通过,但出现了Task Instance State失败的报错日志:

tests/test_instance_context.py [2025-04-26T12:51:18.289+0000] {taskinstance.py:2604} INFO - Dependencies not met for <TaskInstance: test_dag.test manual__2025-04-05T00:00:00+00:00 [failed]>, dependency 'Task Instance State' FAILED: Task is in the 'failed' state.
[2025-04-26T12:51:18.303+0000] {taskinstance.py:2604} INFO - Dependencies not met for <TaskInstance: test_dag.test manual__2025-04-06T00:00:00+00:00 [failed]>, dependency 'Task Instance State' FAILED: Task is in the 'failed' state.   

相关测试代码

tests/conftest.py

import datetime
import pytest
from airflow.models import DAG

@pytest.fixture
def test_dag():
    return DAG(
        "test_dag",
        default_args={
            "owner": "airflow",
            "start_date": datetime.datetime(2025, 4, 5),
            "end_date": datetime.datetime(2025, 4, 6)
        },
        schedule=datetime.timedelta(days=1)
    )

tests/test_instance_context.py

import datetime

from airflow.models import BaseOperator
from airflow.models.dag import DAG
from airflow.utils import timezone

class SampleDAG(BaseOperator):
    template_fields = ("_start_date", "_end_date")

    def __init__(self, start_date, end_date, **kwargs):
        super().__init__(**kwargs)
        self._start_date = start_date
        self._end_date = end_date
    
    def execute(self, context):
        context["ti"].xcom_push(key="start_date", value=self.start_date)
        context["ti"].xcom_push(key="end_date", value=self.end_date)
        return context


def test_execute(test_dag: DAG):
    task = SampleDAG(
        task_id="test",
        start_date="{{ prev_ds }}",
        end_date="{{ ds }}",
        dag=test_dag
    )

    task.run(
        start_date=test_dag.default_args["start_date"], 
        end_date=test_dag.default_args["end_date"]
    )

    expected_start_date = datetime.datetime(2025, 4, 5, tzinfo=timezone.utc)
    expected_end_date = datetime.datetime(2025, 4, 6, tzinfo=timezone.utc)
    assert task.start_date == expected_start_date
    assert task.end_date == expected_end_date

报错原因与解决方法

原因

task.run()会为每个调度周期创建真实的Task Instance并写入元数据存储,测试结束后这些实例会残留failed状态(因为测试没有完整模拟Airflow的调度流程),后续依赖检查时就会抛出该报错。

解决方法

1. 用task.test()替代task.run()(推荐)

task.test()是Airflow专为测试设计的方法,它会模拟Task Instance执行,不会留下持久化的状态残留,同时能正确渲染模板变量。修改测试代码中的执行逻辑:

task.test(
    start_date=test_dag.default_args["start_date"], 
    end_date=test_dag.default_args["end_date"]
)

2. 修复Operator属性访问问题

原SampleDAG类中使用self.start_date但未定义对应属性,需添加属性访问器确保模板渲染后能正确取值:

class SampleDAG(BaseOperator):
    template_fields = ("_start_date", "_end_date")

    def __init__(self, start_date, end_date, **kwargs):
        super().__init__(**kwargs)
        self._start_date = start_date
        self._end_date = end_date
    
    @property
    def start_date(self):
        return self._start_date
    
    @property
    def end_date(self):
        return self._end_date
    
    def execute(self, context):
        context["ti"].xcom_push(key="start_date", value=self.start_date)
        context["ti"].xcom_push(key="end_date", value=self.end_date)
        return context

3. 强制跳过依赖检查(不推荐)

如果坚持使用task.run(),可以添加参数跳过状态依赖检查,但可能引入其他测试隐患:

task.run(
    start_date=test_dag.default_args["start_date"], 
    end_date=test_dag.default_args["end_date"],
    ignore_depends_on_past=True,
    force=True
)

内容的提问来源于stack exchange,提问作者hhk

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 08:28:14