Airflow测试中如何重写任务装饰器以适配API Mock?
问题背景
我使用@task.external_python装饰器指向包含额外依赖的自定义venv,希望编写集成测试,用requests-mock模拟特定API,让DAG端到端运行(pytest与DAG需在同一进程)。曾设想实现自定义任务装饰器,根据环境在external_python和python之间切换,但该方案属于猴子补丁,会侵入生产代码,因此不认可。目前已完成大量单任务测试,需要针对特定数据集测试DAG。
问题解答
1. 是否可以在测试中动态重新定义任务Operator?
可以,且不会侵入生产代码。在pytest用例或fixture中,直接遍历DAG任务,将ExternalPythonOperator替换为PythonOperator即可,示例代码:
from airflow.providers.cncf.kubernetes.operators.python import ExternalPythonOperator from airflow.operators.python import PythonOperator def test_dag_end_to_end(dag): # 遍历DAG中的所有任务 for task in dag.tasks: if isinstance(task, ExternalPythonOperator): # 替换Operator类型 task.__class__ = PythonOperator # 删除ExternalPythonOperator特有的参数,避免兼容问题 if hasattr(task, "external_python"): del task.external_python
替换后任务会在测试进程内运行,可直接结合requests-mock进行API模拟。
2. 如果不行,如何在不构建Provider的情况下编写自定义任务装饰器?
无需构建Provider,直接基于Airflow原生装饰器封装环境分支逻辑即可,示例代码:
from airflow.decorators import task import os def conditional_task(**kwargs): def decorator(func): # 通过环境变量标记测试环境 is_test_env = os.getenv("AIRFLOW_TEST_MODE") == "true" if is_test_env: # 测试环境使用python装饰器,在当前进程运行 return task.python(**kwargs)(func) else: # 生产环境使用external_python装饰器 return task.external_python(**kwargs)(func) return decorator
生产代码中使用@conditional_task(external_python="/path/to/venv/bin/python")装饰任务,测试前设置环境变量export AIRFLOW_TEST_MODE=true即可切换到测试模式。该方案仅做环境分支判断,不属于猴子补丁,不会侵入生产逻辑。
3. 或者是否有更优方案为任务中的API提供Mock?
推荐两种更简洁的方案:
方案一:结合Operator替换+requests-mock注入
先按问题1的方法将任务切换为PythonOperator,再利用requests-mock的fixture直接模拟API,示例代码:
import pytest import requests_mock from airflow.models import DagBag @pytest.fixture def target_dag(): # 加载目标DAG dagbag = DagBag(dag_folder="/your/dag/path", include_examples=False) return dagbag.get_dag("your_dag_id") def test_dag_with_mock(target_dag, requests_mock): # 替换ExternalPythonOperator为PythonOperator for task in target_dag.tasks: if isinstance(task, ExternalPythonOperator): task.__class__ = PythonOperator del task.external_python # 模拟目标API requests_mock.get("https://api.example.com/target", json={"result": "test_data"}) # 端到端运行DAG target_dag.run()
方案二:让任务代码支持Mock注入
修改任务函数,增加可选的session参数,测试时传入Mock的session,示例代码:
import requests from airflow.decorators import task @task.external_python(external_python="/path/to/venv/bin/python") def fetch_api_data(session=None): # 默认使用真实session,测试时可传入Mock session session = session or requests.Session() response = session.get("https://api.example.com/target") return response.json() # 测试用例 def test_fetch_data(requests_mock): # 创建带Session支持的Mock mock_session = requests_mock.Mocker(session=True) mock_session.get("https://api.example.com/target", json={"result": "test"}) with mock_session: # 执行任务并传入Mock session result = fetch_api_data(session=mock_session).execute() assert result == {"result": "test"}
该方案无需修改Operator,仅增强任务的可测试性,结合Operator替换方案即可实现端到端的Mock测试。
内容的提问来源于stack exchange,提问作者Vlad Miller

