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

