Airflow @task.external_python中传递Pandas DataFrame失败排查
Airflow 2.4.1(Docker部署)外部虚拟环境任务间传递Pandas DataFrame失败
目标
- 使用Docker部署的Airflow 2.4.1版本
- 每个任务配置独立外部Python虚拟环境
- 需要在任务间传递Pandas DataFrame
- 此前传递Python列表可成功,但切换为DataFrame后流程崩溃
- 第一个任务可正常执行
问题代码
from airflow.decorators import dag, task from pendulum import datetime from datetime import timedelta my_default_args = { 'owner': 'Anonymus', # 'email': ['random@random.com'], # 'email_on_failure': True, # 'email_on_retry': False, #only allow if it was allowed in the scheduler # 'retries': 1, #only allow if it was allowed in the scheduler # 'retry_delay': timedelta(minutes=1), # 'depends_on_past': False, } @dag( dag_id='test_global_variable', schedule='12 11 * * *', start_date=datetime(2023,2,1,tz="UTC"), catchup=False, default_args=my_default_args, tags=['sample_tag', 'sample_tag2'], ) def write_var(): @task.external_python(task_id="task_1", python='/opt/airflow/v1/bin/python3') def add_to_list(my_list): print(my_list) import pandas as pd df = pd.DataFrame(my_list) return df @task.external_python(task_id="task_2", python='/opt/airflow/v1/bin/python3') def add_to_list_2(df): print(df) df = df.append([19]) return df add_to_list_2(add_to_list([23, 5, 8])) write_var()
错误日志(核心部分)
[2023-02-07, 14:06:24 GMT] {process_utils.py:187} INFO - Traceback (most recent call last): [2023-02-07, 14:06:24 GMT] {process_utils.py:187} INFO - File "/tmp/tmddiox599m/script.py", line 17, in <module> [2023-02-07, 14:06:24 GMT] {process_utils.py:187} INFO - arg_dict = pickle.load(file) [2023-02-07, 14:06:24 GMT] {process_utils.py:187} INFO - AttributeError: 无法在模块 'pandas._libs.internals'(路径:/opt/airflow/venv1/lib/python3.8/site-packages/pandas/_libs/internals.cpython-38-x86_64-linux-gnu.so)上找到属性 '_unpickle_block'
问题原因
该错误由Pandas版本不兼容导致:Airflow主环境与外部虚拟环境的Pandas版本不一致,用pickle序列化的DataFrame在反序列化时,找不到对应版本的内部方法_unpickle_block。即使两个任务使用同一个虚拟环境,只要主环境和虚拟环境的Pandas版本不同,就会触发这个问题。
解决办法
方法1:统一所有环境的Pandas版本
- 查看外部虚拟环境的Pandas版本:
/opt/airflow/v1/bin/python3 -c "import pandas; print(pandas.__version__)"
- 在Airflow主环境安装完全相同的版本:
pip install pandas==x.x.x # 替换x.x.x为虚拟环境的版本号
方法2:改用文件格式中转(推荐生产环境使用)
不直接传递DataFrame对象,而是将其序列化为CSV/Parquet等通用格式,任务间传递文件路径:
from airflow.decorators import dag, task from pendulum import datetime import pandas as pd import os my_default_args = { 'owner': 'Anonymus', } @dag( dag_id='test_global_variable', schedule='12 11 * * *', start_date=datetime(2023,2,1,tz="UTC"), catchup=False, default_args=my_default_args, tags=['sample_tag', 'sample_tag2'], ) def write_var(): @task.external_python(task_id="task_1", python='/opt/airflow/v1/bin/python3') def add_to_list(my_list): print(my_list) df = pd.DataFrame(my_list) # 保存为临时CSV(生产环境建议用共享存储或对象存储) temp_path = '/tmp/temp_df.csv' df.to_csv(temp_path, index=False) return temp_path @task.external_python(task_id="task_2", python='/opt/airflow/v1/bin/python3') def add_to_list_2(file_path): df = pd.read_csv(file_path) print(df) # 注意:pd.append已弃用,改用pd.concat new_row = pd.DataFrame([19]) df = pd.concat([df, new_row], ignore_index=True) return df add_to_list_2(add_to_list([23, 5, 8])) write_var()
方法3:配置Airflow XCom自定义序列化器
在airflow.cfg中添加以下配置,让Airflow专门处理Pandas对象的序列化:
[xcom] enable_xcom_pickling = True custom_serializers = pandas.DataFrame custom_deserializers = pandas.DataFrame
注意:此方法仍需保证所有环境的Pandas版本一致,否则仍可能出现序列化失败。
内容的提问来源于stack exchange,提问作者sogu
相关产品推荐
相关产品推荐

