如何为Apache Airflow中的PythonOperator任务编写单元测试?
Airflow任务单元测试问题
我是Apache Airflow的完全新手,刚搭建了基础环境。我已掌握使用pytest对任务关联函数编写单元测试的方法,但想了解:
- 能否针对DAG中由PythonOperator创建的
process_user、store_user等任务本身编写单元测试,且无需运行完整DAG? - 是否可将这些任务导入单独的Python文件编写测试用例,避免原文件过于杂乱?
相关任务及关联函数
任务逻辑函数
def _process_user(ti): user = ti.xcom_pull(task_ids="extract_user") user = user['results'][0] processed_user = json_normalize({ 'firstname': user['name']['first'], 'lastname': user['name']['last'], 'country': user['location']['country'], 'username': user['login']['username'], 'password': user['login']['password'], 'email': user['email'] }) processed_user.to_csv('/tmp/processed_user.csv', index=None, header=False) def _store_user(): hook = PostgresHook(postgres_conn_id='postgres') hook.copy_expert( sql="COPY users FROM stdin WITH DELIMITER as ','", filename='/tmp/processed_user.csv' )
DAG中的任务定义
process_user = PythonOperator( task_id='process_user', python_callable=_process_user ) store_user = PythonOperator( task_id='store_user', python_callable=_store_user ) create_table >> is_api_available >> extract_user >> process_user >> store_user
问题解答
1. 可以针对PythonOperator任务编写单元测试,无需运行完整DAG
PythonOperator本质就是封装了你定义的python_callable函数,测试任务本身核心还是测试对应的函数,只要处理好Airflow的依赖(比如ti对象、Hook)就行:
- 测试
_process_user函数:需要mockti.xcom_pull返回模拟的用户数据,再验证生成的CSV内容是否符合预期
import pytest import pandas as pd from unittest.mock import Mock from your_dag_file import _process_user def test_process_user(): # 模拟XCom返回的用户数据 mock_user_data = { "results": [ { "name": {"first": "John", "last": "Doe"}, "location": {"country": "USA"}, "login": {"username": "johndoe", "password": "testpass"}, "email": "john@example.com" } ] } # Mock TaskInstance对象 mock_ti = Mock() mock_ti.xcom_pull.return_value = mock_user_data # 执行函数 _process_user(mock_ti) # 验证CSV内容 df = pd.read_csv('/tmp/processed_user.csv', names=['firstname', 'lastname', 'country', 'username', 'password', 'email']) assert df.iloc[0]['firstname'] == 'John' assert df.iloc[0]['country'] == 'USA'
- 测试
_store_user函数:需要mockPostgresHook和它的copy_expert方法,验证调用参数是否正确
from unittest.mock import patch from your_dag_file import _store_user def test_store_user(): with patch('your_dag_file.PostgresHook') as mock_hook_class: # 获取mock的hook实例 mock_hook = mock_hook_class.return_value # 执行函数 _store_user() # 验证Hook初始化参数 mock_hook_class.assert_called_once_with(postgres_conn_id='postgres') # 验证copy_expert的调用参数 mock_hook.copy_expert.assert_called_once_with( sql="COPY users FROM stdin WITH DELIMITER as ','", filename='/tmp/processed_user.csv' )
如果想直接测试PythonOperator实例而非函数,也可以调用task.execute(context)方法,传入mock的context字典:
from your_dag_file import process_user from unittest.mock import Mock def test_process_user_operator(): mock_context = {'ti': Mock()} mock_context['ti'].xcom_pull.return_value = {"results": [/* 模拟数据 */]} # 直接执行Operator的execute方法 process_user.execute(mock_context) # 后续验证逻辑和之前一致
2. 完全可以将测试用例放在单独文件中
Airflow的DAG文件就是普通Python代码,有两种方式实现:
- 抽离逻辑函数到单独模块:把
_process_user、_store_user放到比如dags/utils/user_processing.py里,然后在DAG文件中导入使用:
# dags/utils/user_processing.py def process_user(ti): # 原_process_user的逻辑 def store_user(): # 原_store_user的逻辑 # dags/your_dag.py from airflow.operators.python import PythonOperator from utils.user_processing import process_user, store_user process_user_task = PythonOperator( task_id='process_user', python_callable=process_user ) # 其他任务定义...
然后在项目根目录创建tests文件夹,编写测试文件(比如tests/test_user_processing.py),直接导入抽离后的函数测试,完全不会污染原DAG文件。
- 直接从DAG文件导入:即使不抽离函数,也可以直接从DAG文件中导入
_process_user、_store_user或者Operator实例到测试文件中编写用例,只要项目结构符合Python模块导入规则就行。
内容的提问来源于stack exchange,提问作者Aatr Sin
相关产品推荐
相关产品推荐

