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

如何为@task装饰的Airflow任务编写单元测试

Airflow TaskFlow 任务异常分支单元测试方案

问题根因

你之前的测试不生效,核心原因有两点:

  • 被@task()装饰的函数会被包装为Airflow算子实例,直接加括号调用只会完成算子初始化,不会执行函数内部的业务逻辑,因此你配置的mock异常、调用断言都不会触发。
  • 通过创建DagRun拉取任务执行的方案属于DAG集成测试范畴,依赖Airflow元数据库、执行器等完整运行环境,执行重、环境干扰多,不适合做单函数级别的分支逻辑验证。

可直接运行的测试实现

两种轻量方案都可以跳过Airflow调度逻辑,直接验证业务分支:

方案1:调用算子绑定的原始业务函数

TaskFlow装饰器会将未包装的原始函数存在算子实例的python_callable属性上,直接调用该属性即可执行业务逻辑,不需要依赖Airflow运行时:

from unittest import TestCase, mock

class TestAccountLinkUpdatesProcess(TestCase):
    @mock.patch("dags.delta_load.updates.log")
    @mock.patch("dags.delta_load.updates.get_current_context")
    @mock.patch("dags.delta_load.updates.utils.download_file_from_s3_bucket")
    def test_file_not_found_error(self, mock_download, mock_get_ctx, mock_log):
        # 配置mock抛出目标异常
        mock_download.side_effect = FileNotFoundError("test file missing")
        # 实例化任务对象
        task = updates_process({"updates_file": "path/to/file.csv"})
        # 直接执行原始业务函数,不要直接调用task实例
        execute_result = task.python_callable({"updates_file": "path/to/file.csv"})

        # 验证逻辑符合预期
        mock_get_ctx.assert_called_once()
        mock_download.assert_called_once_with("path/to/file.csv")
        mock_log.error.assert_called_once()
        # 异常分支触发提前return,返回值应为None
        self.assertIsNone(execute_result)

方案2:测试阶段剥离TaskFlow装饰器

如果不想依赖算子的内部属性,可以在导入被测试函数前,先mock掉@task装饰器,让它直接返回原始函数,后续调用和普通Python函数完全一致:

from unittest import TestCase, mock

# 导入被测试模块前先替换task装饰器,跳过算子包装逻辑
with mock.patch("airflow.decorators.task", lambda *args, **kwargs: lambda func: func):
    from dags.delta_load.updates import updates_process

class TestAccountLinkUpdatesProcess(TestCase):
    @mock.patch("dags.delta_load.updates.log")
    @mock.patch("dags.delta_load.updates.get_current_context")
    @mock.patch("dags.delta_load.updates.utils.download_file_from_s3_bucket")
    def test_file_not_found_error(self, mock_download, mock_get_ctx, mock_log):
        mock_download.side_effect = FileNotFoundError("test file missing")
        # 此时updates_process是普通Python函数,直接调用即可
        execute_result = updates_process({"updates_file": "path/to/file.csv"})

        mock_get_ctx.assert_called_once()
        mock_download.assert_called_once_with("path/to/file.csv")
        mock_log.error.assert_called_once()
        self.assertIsNone(execute_result)

测试注意事项

  • 单元测试阶段不要使用DagRun执行链路,这套逻辑仅适合验证DAG结构、任务依赖关系,不适合单任务内部逻辑验证,执行效率低且容易被环境因素干扰。
  • mock路径必须填写被测试文件内引用对象的路径,不要mock工具类源文件中的对象,否则会出现mock不生效的问题。
  • 如果业务逻辑中用到了上下文变量,可以给mock_get_ctx配置return_value,传入自定义的假上下文字典即可,不需要真实的Airflow运行上下文。

内容的提问来源于stack exchange,提问作者Sadan A.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 14:15:41