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

在GCP Airflow中使用PythonVirtualenvOperator遇_pickle.PicklingError求助

问题:Airflow中使用@task.virtualenv与PythonOperator结合引发PicklingError

问题场景

将业务逻辑与DAG定义分离到两个文件:Main.py(带@task.virtualenv装饰的任务函数)和task_dag.py(用PythonOperator调用该函数),运行时触发pickle序列化错误。

Main.py 代码

from airflow.decorators import task

@task.virtualenv(
    task_id="virtualenv_python", requirements=["colorama==0.4.0"], system_site_packages=True
)
def Main_func():
    from time import sleep
    from colorama import Back, Fore, Style

    print(Fore.RED + "some red text")
    print(Back.GREEN + "and with a green background")
    print(Style.DIM + "and in dim text")
    print(Style.RESET_ALL)
    for _ in range(4):
        print(Style.DIM + "Please wait...", flush=True)
        sleep(1)
    print("Finished")

task_dag.py 代码

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
from Main import Main_func

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'max_active_runs': 1,
}

Main_Mod_dag = DAG(
    'Main_Mod_Run',
    catchup=False,
    default_args=default_args,
    description='Main Module Run',
    schedule_interval='48 11 * * 3',
    start_date=datetime(2022, 12, 11),
    tags=['Main_Mod'],
)

Main_Mod_Func = PythonOperator(task_id='Main_Mod', python_callable=Main_func, dag=Main_Mod_dag)

Main_Mod_Func

报错信息

[2023-05-10, 10:47:13 UTC] {python.py:177} INFO - Done. Returned value was: {{ task_instance.xcom_pull(task_ids='virtualenv_python', dag_id='adhoc_airflow', key='return_value') }} 
[2023-05-10, 10:47:13 UTC] {taskinstance.py:1853} ERROR - Task failed with exception 
Traceback (most recent call last): 
File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/session.py", line 72, in wrapper 
return func(*args, **kwargs) 
File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 2381, in xcom_push 
XCom.set( 
File "/opt/python3.8/lib/python3.8/site-packages/airflow/utils/session.py", line 72, in wrapper 
return func(*args, **kwargs) 
File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/xcom.py", line 206, in set 
value = cls.serialize_value( 
File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/xcom.py", line 595, in serialize_value 
return pickle.dumps(value) 
_pickle.PicklingError: Can't pickle <function Main_func at 0x7fe5c18dba60>: it's not the same object as Main.Main_func

错误原因

@task.virtualenv装饰器并非返回普通Python函数,而是将Main_func包装成了Airflow的Task对象。PythonOperator要求python_callable参数是可直接调用的函数,当尝试序列化这个包装后的Task对象时,pickle无法处理跨模块的对象引用,从而抛出错误。

解决方案

方案1:直接使用装饰后的Task对象(推荐)

无需通过PythonOperator调用,直接将装饰后的任务对象加入DAG即可:

修改后的task_dag.py:

from airflow import DAG
from datetime import datetime, timedelta
from Main import Main_func

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'max_active_runs': 1,
}

with DAG(
    'Main_Mod_Run',
    catchup=False,
    default_args=default_args,
    description='Main Module Run',
    schedule_interval='48 11 * * 3',
    start_date=datetime(2022, 12, 11),
    tags=['Main_Mod'],
) as Main_Mod_dag:
    # 直接实例化装饰后的任务对象
    Main_Mod_Func = Main_func()

方案2:显式使用PythonVirtualenvOperator

如果需要保留PythonOperator的使用方式,移除@task.virtualenv装饰器,改用PythonVirtualenvOperator定义任务:

修改后的Main.py:

def Main_func():
    from time import sleep
    from colorama import Back, Fore, Style

    print(Fore.RED + "some red text")
    print(Back.GREEN + "and with a green background")
    print(Style.DIM + "and in dim text")
    print(Style.RESET_ALL)
    for _ in range(4):
        print(Style.DIM + "Please wait...", flush=True)
        sleep(1)
    print("Finished")

修改后的task_dag.py:

from airflow import DAG
from airflow.operators.python import PythonVirtualenvOperator
from datetime import datetime, timedelta
from Main import Main_func

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'max_active_runs': 1,
}

Main_Mod_dag = DAG(
    'Main_Mod_Run',
    catchup=False,
    default_args=default_args,
    description='Main Module Run',
    schedule_interval='48 11 * * 3',
    start_date=datetime(2022, 12, 11),
    tags=['Main_Mod'],
)

Main_Mod_Func = PythonVirtualenvOperator(
    task_id='Main_Mod',
    python_callable=Main_func,
    requirements=["colorama==0.4.0"],
    system_site_packages=True,
    dag=Main_Mod_dag
)

Main_Mod_Func

说明

  • 方案1利用Airflow的装饰器语法简化任务定义,是官方推荐的写法,代码更简洁。
  • 方案2适合需要对虚拟环境任务进行更细粒度配置的场景,显式控制PythonVirtualenvOperator的参数。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 13:05:05