如何读写含字典列表列的DataFrame到Parquet并还原原对象
解决Parquet读写后字典列表结构失真问题
当将包含字典列表的DataFrame写入Parquet再读回时,出现字典字段重复、缺失字段被填充为None的问题,无法还原原始结构。复现代码及现象如下:
复现代码:
import pyarrow as pa from pyarrow import parquet import pandas as pd COLUMN1_SCHEMA = pa.list_(pa.struct([('Id', pa.string()), ('Age', pa.string())])) SCHEMA = pa.schema([pa.field("column1", COLUMN1_SCHEMA), ('column2', pa.int32())]) df = pd.DataFrame({ "column1": [[{"Id": "1"}, {"Age": "16"}], [{"Id": "2"},{"Age": "17"}]], "column2": [1, 2], }) table = pa.Table.from_pandas(df, schema=SCHEMA) parquet.write_table(table, "f.parquet") df = pa.parquet.read_table("f.parquet", schema=SCHEMA).to_pandas()
读回后异常结果:
column1 column2 0 [{'Id': '1','Age': None}, {'Id': None, 'Age': '16'}] 1 1 [{'Id': '2','Age': None}, {'Id': None, 'Age': '17'}] 2
期望结果:
column1 column2 0 [{'Id': '1'}, {'Age': '16'}] 1 1 [{'Id': '2'}, {'Age': '17'}] 2
问题根源
你定义的COLUMN1_SCHEMA是包含Id和Age字段的struct的列表,要求列表中每个元素必须同时拥有这两个字段。但原始数据中每个列表元素仅含单个字段,PyArrow在强制匹配schema时,会自动补全缺失字段为None,写入Parquet后读回时保留了这些补全的None值,导致结构偏离原始数据。
解决方案
方案1:使用Union类型定义schema(严格匹配原始结构)
通过PyArrow的稀疏Union类型定义schema,让每个列表元素仅存储存在的字段,无需补全None。读写后再转换回原始字典格式:
import pyarrow as pa from pyarrow import parquet import pandas as pd # 定义稀疏Union类型,每个元素为Id或Age字段 union_type = pa.union( [pa.field('Id', pa.string()), pa.field('Age', pa.string())], mode='sparse' ) COLUMN1_SCHEMA = pa.list_(union_type) SCHEMA = pa.schema([pa.field("column1", COLUMN1_SCHEMA), ('column2', pa.int32())]) df = pd.DataFrame({ "column1": [[{"Id": "1"}, {"Age": "16"}], [{"Id": "2"},{"Age": "17"}]], "column2": [1, 2], }) # 将原始数据转换为符合Union类型的结构 def convert_to_union_list(lst): pa_items = [] for d in lst: if 'Id' in d: pa_items.append(pa.scalar(d['Id'], type=pa.field('Id', pa.string()))) elif 'Age' in d: pa_items.append(pa.scalar(d['Age'], type=pa.field('Age', pa.string()))) return pa.array(pa_items, type=union_type) # 构建Arrow Table column1_data = pa.array([convert_to_union_list(lst) for lst in df['column1']], type=COLUMN1_SCHEMA) table = pa.Table.from_arrays( [column1_data, pa.array(df['column2'], type=pa.int32())], schema=SCHEMA ) # 写入并读回 parquet.write_table(table, "f.parquet") read_table = parquet.read_table("f.parquet", schema=SCHEMA) df_read = read_table.to_pandas() # 将Union类型转换回原始字典格式 def union_scalar_to_dict(scalar): field_name = scalar.type.type_names[scalar.tag] return {field_name: scalar.as_py()} def restore_column1(col): return [ [union_scalar_to_dict(item) for item in lst] for lst in col ] df_read['column1'] = restore_column1(df_read['column1']) print(df_read)
方案2:不强制指定schema(简单快捷)
去掉自定义schema,让PyArrow自动推断数据结构,读写后直接还原原始结构:
import pyarrow as pa from pyarrow import parquet import pandas as pd df = pd.DataFrame({ "column1": [[{"Id": "1"}, {"Age": "16"}], [{"Id": "2"},{"Age": "17"}]], "column2": [1, 2], }) # 自动推断schema写入 table = pa.Table.from_pandas(df) parquet.write_table(table, "f.parquet") # 读回 df_read = parquet.read_table("f.parquet").to_pandas() print(df_read)
此方法无需额外处理,但无法预先约束schema,适合对schema要求不严格的场景。
方案3:调整原始数据匹配schema(严格遵循schema约束)
如果必须使用最初定义的struct列表schema,可预先补全字典的缺失字段为None,确保数据符合schema要求:
import pyarrow as pa from pyarrow import parquet import pandas as pd COLUMN1_SCHEMA = pa.list_(pa.struct([('Id', pa.string()), ('Age', pa.string())])) SCHEMA = pa.schema([pa.field("column1", COLUMN1_SCHEMA), ('column2', pa.int32())]) # 补全缺失字段为None,匹配schema df = pd.DataFrame({ "column1": [[{"Id": "1", "Age": None}, {"Id": None, "Age": "16"}], [{"Id": "2", "Age": None}, {"Id": None, "Age": "17"}]], "column2": [1, 2], }) table = pa.Table.from_pandas(df, schema=SCHEMA) parquet.write_table(table, "f.parquet") df_read = parquet.read_table("f.parquet", schema=SCHEMA).to_pandas() print(df_read)
此方法保证schema一致性,但结构上会保留None值,需业务逻辑接受该格式。
内容的提问来源于stack exchange,提问作者Miguel
相关产品推荐
相关产品推荐

