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

Airflow虚拟环境任务调用外部函数报NameError的问题求助

问题:Airflow virtualenv任务调用外部函数触发NameError

原始代码

from datetime import datetime
from airflow import DAG
from airflow.decorators import task

default_args = {
    'owner': 'airflow',
    'retries': 1,
}

with DAG(
    "test_dag_venv",
    default_args=default_args,
    description='Dag to test venv',
    schedule_interval="@once",
    start_date=datetime(2022, 1, 6, 10, 45),
    tags=['testing'],
    concurrency=1,
    is_paused_upon_creation=True,
    catchup=False  # dont run previous and backfill; run only latest
) as dag:
    def print_test_1():
        print('print test 1')
    @task.virtualenv(task_id="print_test", requirements=['numpy'], system_site_packages=False)
    def print_test():
        import numpy as np
        print(np.__version__)
        print(print_test_1()) 
    t1 = print_test()
    t1

错误信息

[2022-09-20 11:14:01,994] {process_utils.py:135} INFO - Executing cmd: /tmp/venvzfr52lj1/bin/python /tmp/venvzfr52lj1/script.py /tmp/venvzfr52lj1/script.in /tmp/venvzfr52lj1/script.out /tmp/venvzfr52lj1/string_args.txt
[2022-09-20 11:14:02,003] {process_utils.py:139} INFO - Output:
[2022-09-20 11:14:02,439] {process_utils.py:143} INFO - 1.19.5
[2022-09-20 11:14:02,439] {process_utils.py:143} INFO - Traceback (most recent call last):
[2022-09-20 11:14:02,440] {process_utils.py:143} INFO -   File "/tmp/venvzfr52lj1/script.py", line 33, in <module>
[2022-09-20 11:14:02,440] {process_utils.py:143} INFO -     res = print_test(*arg_dict["args"], **arg_dict["kwargs"])
[2022-09-20 11:14:02,440] {process_utils.py:143} INFO -   File "/tmp/venvzfr52lj1/script.py", line 31, in print_test
[2022-09-20 11:14:02,441] {process_utils.py:143} INFO -     print(print_test_1())
[2022-09-20 11:14:02,441] {process_utils.py:143} INFO - NameError: name 'print_test_1' is not defined
[2022-09-20 11:14:02,562] {taskinstance.py:1501} ERROR - Task failed with exception

错误原因

Airflow的@task.virtualenv装饰器会把被装饰的函数代码单独提取,在临时创建的虚拟环境中作为独立脚本执行。外部定义的print_test_1不会被打包到这个临时脚本里,因此在虚拟环境的执行上下文里找不到该函数,触发NameError。

即便给print_test_1也加@task.virtualenv装饰器,它会变成另一个独立的Airflow任务,每个任务都是隔离的进程/虚拟环境,依然无法直接在print_test函数内部调用。

解决方向

1. 若print_test_1仅为辅助函数,无需独立成任务

将print_test_1嵌套到print_test函数内部,这样它会被一起打包到虚拟环境的执行脚本中:

@task.virtualenv(task_id="print_test", requirements=['numpy'], system_site_packages=False)
def print_test():
    # 把辅助函数放到virtualenv任务内部
    def print_test_1(input_val):
        print('print test 1')
        return f"processed_{input_val}"
    
    import numpy as np
    print(np.__version__)
    # 调用内部辅助函数并传值
    result = print_test_1("test_input")
    print(result)
    # 基于返回结果执行后续代码
    final_result = f"final_output_{result}"
    print(final_result)

2. 若print_test_1需要作为独立任务执行,传递结果给print_test

使用Airflow的任务依赖和XCom机制传递数据,实现任务间的结果交互:

with DAG(
    "test_dag_venv",
    default_args=default_args,
    description='Dag to test venv',
    schedule_interval="@once",
    start_date=datetime(2022, 1, 6, 10, 45),
    tags=['testing'],
    concurrency=1,
    is_paused_upon_creation=True,
    catchup=False
) as dag:
    # 定义独立任务print_test_1
    @task
    def print_test_1(input_val):
        print('print test 1')
        return f"processed_{input_val}"

    # virtualenv任务接收上游任务的结果
    @task.virtualenv(task_id="print_test", requirements=['numpy'], system_site_packages=False)
    def print_test(prev_result):
        import numpy as np
        print(np.__version__)
        print(f"Received from print_test_1: {prev_result}")
        # 基于上游结果执行后续逻辑
        final_result = f"final_{prev_result}"
        return final_result

    # 设置任务依赖,自动传递结果
    t1 = print_test_1(input_val="initial_input")
    t2 = print_test(t1.output)
    t1 >> t2

这种方式下,print_test_1的返回值会通过XCom自动传递给print_test,实现跨任务的传值和后续逻辑执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 23:30:56