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

