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

解决Airflow跨Task传递DataFrame时的Parquet空结构体写入错误

问题:Airflow Task间传递DataFrame时序列化失败

问题背景与代码示例

原本使用pandas 2.0时,Task间传递Dict无异常;升级至pandas 2.2.0并安装pyarrow后,将DataFrame转Dict的逻辑移至upsert方法,需要直接传递DataFrame,出现序列化错误。

DAG代码示例:

def ozon_items():

    @task.external_python(python='/usr/local/airflow/venv_main/bin/python')
    def get_ozon_items():
        from src.transform import OzonTables
        # 请求API并转换,返回DataFrame
        return OzonTables().items()

    @task.external_python(python='/usr/local/airflow/venv_main/bin/python')
    def update_ozon_items(data):
        from src.databases.DatabaseWorker import DatabaseWorker
        from src.databases.markets.ozon_schema import OzonItems

        db = DatabaseWorker(host='host.docker.internal', port=5431)

        with db.session() as session:
            session.upsert(
                OzonItems,
                data,
                on_conflict='do_update',
                deletable=True
            )
    update_ozon_items(get_ozon_items())

报错信息

关键错误:

pyarrow.lib.ArrowNotImplementedError: Cannot write struct type 'optional_description_elements' with no child field to Parquet. Consider adding a dummy child field.

完整报错堆栈:

{taskinstance.py:2699} ERROR - Task failed with exception
Traceback (most recent call last):
File "/usr/local/lib/python3.11/site-packages/airflow/models/taskinstance.py", line 440, in _execute_task
 task_instance.xcom_push(key=XCOM_RETURN_KEY, value=xcom_value, session=session)
File "/usr/local/lib/python3.11/site-packages/airflow/utils/session.py", line 76, in wrapper
 return func(*args, **kwargs)
            ^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/models/taskinstance.py", line 2981, in xcom_push
 XCom.set(
File "/usr/local/lib/python3.11/site-packages/airflow/utils/session.py", line 76, in wrapper
 return func(*args, **kwargs)
            ^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/models/xcom.py", line 247, in set
 value = cls.serialize_value(
             ^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/models/xcom.py", line 662, in serialize_value
 return json.dumps(value, cls=XComEncoder).encode("UTF-8")
            ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/json/__init__.py", line 238, in dumps
 **kw).encode(obj)
           ^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/utils/json.py", line 104, in encode
 return super().encode(o)
            ^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/json/encoder.py", line 200, in encode
 chunks = self.iterencode(o, _one_shot=True)
              ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/json/encoder.py", line 258, in iterencode
 return _iterencode(o, 0)
            ^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/utils/json.py", line 91, in default
 return serialize(o)
            ^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/serialization/serde.py", line 145, in serialize
 data, classname, version, is_serialized = _serializers[qn].serialize(o)
                                               ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/airflow/serialization/serializers/pandas.py", line 51, in serialize
 pq.write_table(table, buf, compression="snappy")
File "/usr/local/lib/python3.11/site-packages/pyarrow/parquet/core.py", line 3071, in write_table
 with ParquetWriter(
          ^^^^^^^^^^^^^^
File "/usr/local/lib/python3.11/site-packages/pyarrow/parquet/core.py", line 990, in __init__
 self.writer = _parquet.ParquetWriter(
                   ^^^^^^^^^^^^^^^^^^^^^^^
File "pyarrow/_parquet.pyx", line 1753, in pyarrow._parquet.ParquetWriter.__cinit__
File "pyarrow/error.pxi", line 144, in pyarrow.lib.pyarrow_internal_check_status
File "pyarrow/error.pxi", line 121, in pyarrow.lib.check_status
pyarrow.lib.ArrowNotImplementedError: Cannot write struct type 'optional_description_elements' with no child field to Parquet. Consider adding a dummy child field.

解决方案

1. 问题根源

Airflow通过PyArrow将DataFrame序列化为Parquet格式存储到XCom时,DataFrame中存在无子女段的Struct类型列optional_description_elements,PyArrow不支持将空Struct写入Parquet,导致序列化失败。

2. 处理空Struct列

针对optional_description_elements列,选择以下两种方式之一处理:

  • 删除空Struct列:如果该列数据无用,在返回DataFrame前执行删除操作:
    def get_ozon_items():
        from src.transform import OzonTables
        df = OzonTables().items()
        # 删除空Struct列
        if 'optional_description_elements' in df.columns:
            df = df.drop(columns=['optional_description_elements'])
        return df
    
  • 给空Struct添加dummy字段:如果需要保留该列,将空Struct转为包含默认字段的结构:
    import pandas as pd
    
    def get_ozon_items():
        from src.transform import OzonTables
        df = OzonTables().items()
        # 给空Struct添加dummy字段
        def fill_empty_struct(x):
            if pd.isna(x) or (isinstance(x, dict) and len(x) == 0):
                return {'dummy': None}
            return x
        
        if 'optional_description_elements' in df.columns:
            df['optional_description_elements'] = df['optional_description_elements'].apply(fill_empty_struct)
        return df
    

3. 验证数据结构

确保DataFrame中所有Struct类型列都有明确的字段定义,避免出现空Struct(无子女段的结构)。

4. 可选回退方案

如果暂时无法修改数据处理逻辑,可回退至pandas 2.0版本,此时Airflow会用旧的序列化方式处理Dict类型数据,但不推荐长期使用,建议优先修复数据结构问题。


内容的提问来源于stack exchange,提问作者Sergey Bakaev Rettley

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 04:27:05