Airflow PythonVirtualenvOperator报错:找不到unusual_prefix_***_dag模块
Airflow PythonVirtualenvOperator 触发时报错 ModuleNotFoundError: No module named 'unusual_prefix_xxx_dag'
环境与代码
使用Airflow 2.5.3 + Kubernetes Executor + Python 3.7,编写包含PythonVirtualenvOperator的DAG,尝试传递{{ ts }}和{{ dag }}上下文变量,代码如下:
from datetime import timedelta from pathlib import Path import airflow from airflow import DAG from airflow.operators.python import PythonOperator, PythonVirtualenvOperator import pendulum dag = DAG( default_args={ 'retries': 2, 'retry_delay': timedelta(minutes=10), }, dag_id='fs_rb_cashflow_test5', schedule_interval='0 5 * * 1', start_date=pendulum.datetime(2020, 1, 1, tz='UTC'), catchup=False, tags=['Feature Store', 'RB', 'u_m1ahn'], render_template_as_native_obj=True, ) context = {"ts": "{{ ts }}", "dag": "{{ dag }}"} op_args = [context, Path(__file__).parent.absolute()] def make_foo(*args, **kwargs): print("--> making foo!") print("make foo(...): args") print(args) print("make foo(...): kwargs") print(kwargs) make_foo_task = PythonVirtualenvOperator( task_id='make_foo', python_callable=make_foo, provide_context=True, use_dill=True, system_site_packages=False, op_args=op_args, op_kwargs={ "execution_date_str": '{{ execution_date }}', }, requirements=["dill", "pytz", f"apache-airflow=={airflow.__version__}", "psycopg2-binary >= 2.9, < 3"], dag=dag)
报错信息
触发DAG时出现以下错误:
[2023-10-23, 13:30:40] {process_utils.py:187} INFO - Traceback (most recent call last): [2023-10-23, 13:30:40] {process_utils.py:187} INFO - File "/tmp/venv5ifve2a5/script.py", line 17, in <module> [2023-10-23, 13:30:40] {process_utils.py:187} INFO - arg_dict = dill.load(file) [2023-10-23, 13:30:40] {process_utils.py:187} INFO - File "/tmp/venv5ifve2a5/lib/python3.7/site-packages/dill/_dill.py", line 287, in load [2023-10-23, 13:30:40] {process_utils.py:187} INFO - return Unpickler(file, ignore=ignore, **kwds).load() [2023-10-23, 13:30:40] {process_utils.py:187} INFO - File "/tmp/venv5ifve2a5/lib/python3.7/site-packages/dill/_dill.py", line 442, in load [2023-10-23, 13:30:40] {process_utils.py:187} INFO - obj = StockUnpickler.load(self) [2023-10-23, 13:30:40] {process_utils.py:187} INFO - File "/tmp/venv5ifve2a5/lib/python3.7/site-packages/dill/_dill.py", line 432, in find_class [2023-10-23, 13:30:40] {process_utils.py:187} INFO - return StockUnpickler.find_class(self, module, name) [2023-10-23, 13:30:40] {process_utils.py:187} INFO - ModuleNotFoundError: No module named 'unusual_prefix_4c3a45107010a4223aa054ffc5f7bffc78cce4e7_dag'
错误原因
- DAG实例序列化失效:你在
op_args中直接传递了{{ dag }}(DAG实例对象),Airflow序列化DAG时会为其生成带临时前缀的模块名(如unusual_prefix_xxx_dag),但该临时模块仅存在于调度器上下文。虚拟环境独立于调度器环境,反序列化时无法找到该模块,引发报错。 - Path对象无意义:
Path(__file__).parent.absolute()传递的是原DAG文件的本地路径,虚拟环境使用临时目录,该路径在虚拟环境中无效,且Path对象的序列化也存在兼容性问题。 - 模板变量传递错误:
op_args支持模板渲染,但复杂对象(如DAG实例)无法被正确序列化到虚拟环境,仅能传递字符串、数字等可序列化简单类型。
解决方案
1. 替换DAG实例为可序列化标识
不要传递整个DAG实例,改用dag_id字符串。如果需要在callable中操作DAG,可通过dag_id从元数据库查询,或直接从kwargs中获取(因设置了provide_context=True)。
2. 修正路径传递逻辑
避免传递Path对象,改用相对路径字符串,或在callable内部通过Airflow配置获取所需路径,不依赖原DAG文件路径。
3. 优化上下文变量获取
利用provide_context=True自动注入上下文变量,无需手动在op_args中定义ts等,直接从kwargs中读取即可。
修改后的代码示例:
from datetime import timedelta import airflow from airflow import DAG from airflow.operators.python import PythonVirtualenvOperator import pendulum dag = DAG( default_args={ 'retries': 2, 'retry_delay': timedelta(minutes=10), }, dag_id='fs_rb_cashflow_test5', schedule_interval='0 5 * * 1', start_date=pendulum.datetime(2020, 1, 1, tz='UTC'), catchup=False, tags=['Feature Store', 'RB', 'u_m1ahn'], render_template_as_native_obj=True, ) # 仅传递必要的字符串类型参数 op_args = ["{{ ts }}"] def make_foo(*args, **kwargs): print("--> making foo!") print("make foo(...): args") print(args) print("make foo(...): kwargs") print(kwargs) # 从上下文直接获取所需变量 execution_date = kwargs.get('execution_date') dag_id = kwargs.get('dag').dag_id print(f"Execution Date: {execution_date}, DAG ID: {dag_id}") make_foo_task = PythonVirtualenvOperator( task_id='make_foo', python_callable=make_foo, provide_context=True, use_dill=True, system_site_packages=False, op_args=op_args, op_kwargs={ "execution_date_str": '{{ execution_date }}', }, requirements=["dill", "pytz", f"apache-airflow=={airflow.__version__}", "psycopg2-binary >= 2.9, < 3"], dag=dag)
额外说明
- 若需在虚拟环境中操作DAG对象,不要传递实例,而是通过
dag_id,在callable中使用airflow.models.DagModel.get_dag(dag_id)查询(需确保虚拟环境Airflow配置能连接元数据库)。 - 始终在
op_args/op_kwargs中传递JSON可序列化的简单类型,减少序列化反序列化风险。
内容的提问来源于stack exchange,提问作者Felix
相关产品推荐
相关产品推荐

