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

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'

错误原因

  1. DAG实例序列化失效:你在op_args中直接传递了{{ dag }}(DAG实例对象),Airflow序列化DAG时会为其生成带临时前缀的模块名(如unusual_prefix_xxx_dag),但该临时模块仅存在于调度器上下文。虚拟环境独立于调度器环境,反序列化时无法找到该模块,引发报错。
  2. Path对象无意义:Path(__file__).parent.absolute()传递的是原DAG文件的本地路径,虚拟环境使用临时目录,该路径在虚拟环境中无效,且Path对象的序列化也存在兼容性问题。
  3. 模板变量传递错误: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 18:13:18