使用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
相关产品推荐
相关产品推荐

