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

