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

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版本

  1. 查看外部虚拟环境的Pandas版本:
/opt/airflow/v1/bin/python3 -c "import pandas; print(pandas.__version__)"
  1. 在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 22:00:59