Arrow包write_parquet不支持追加,如何生成单个Parquet文件?
解决分批写入单个Parquet文件的方案
方案1:使用PyArrow的ParquetWriter类(官方支持追加)
PyArrow的write_parquet确实不支持追加,但底层的ParquetWriter类支持分批写入,能严格保留日期、浮点等所有数据类型。
步骤:
- 用第一批数据生成schema,初始化
ParquetWriter - 循环处理每批数据,调用
write_table写入 - 完成后关闭writer
代码示例:
import pyarrow as pa import pyarrow.parquet as pq import pandas as pd # 模拟分批生成数据 def generate_batch(batch_num): return pd.DataFrame({ 'date_col': pd.date_range(start='2023-01-01', periods=100), 'float_col': [float(i + batch_num*100) for i in range(100)], 'int_col': range(100) }) # 获取第一批数据的schema first_batch = generate_batch(0) table = pa.Table.from_pandas(first_batch) schema = table.schema # 初始化ParquetWriter writer = pq.ParquetWriter('single_output.parquet', schema) # 分批写入 for batch_num in range(5): batch_df = generate_batch(batch_num) batch_table = pa.Table.from_pandas(batch_df, schema=schema) writer.write_table(batch_table) # 关闭writer writer.close()
注意:所有批次的schema必须与初始化时一致,否则会报错。若后续批次有字段变更,需提前定义兼容的通用schema。
方案2:使用fastparquet包
fastparquet是另一个轻量型Parquet处理库,原生支持追加模式,对日期、浮点类型的处理稳定,内存占用较低。
先安装依赖:pip install fastparquet
代码示例:
import fastparquet as fp import pandas as pd # 模拟分批数据 def generate_batch(batch_num): return pd.DataFrame({ 'date_col': pd.date_range(start='2023-01-01', periods=100), 'float_col': [float(i + batch_num*100) for i in range(100)], 'int_col': range(100) }) # 第一批数据写入(创建文件) first_batch = generate_batch(0) fp.write('single_output_fast.parquet', first_batch) # 后续批次追加写入 for batch_num in range(1,5): batch_df = generate_batch(batch_num) fp.write('single_output_fast.parquet', batch_df, append=True)
注意:追加时需保证各批次的列名、数据类型完全一致,避免类型兼容问题。
方案3:先写临时文件再合并(内存紧张场景首选)
若前两种方案的内存占用仍超出限制,可先将每批数据写入单独的临时Parquet文件,最后用PyArrow读取所有临时文件并合并为单个文件,全程不会破坏数据类型。
代码示例:
import pyarrow as pa import pyarrow.parquet as pq import pandas as pd import os import shutil # 创建临时目录 temp_dir = 'temp_parquets' os.makedirs(temp_dir, exist_ok=True) # 分批写入临时文件 for batch_num in range(5): batch_df = generate_batch(batch_num) batch_table = pa.Table.from_pandas(batch_df) pq.write_table(batch_table, os.path.join(temp_dir, f'batch_{batch_num}.parquet')) # 读取所有临时文件并合并为单个表 dataset = pq.ParquetDataset(temp_dir) combined_table = dataset.read() # 写入最终单个Parquet文件 pq.write_table(combined_table, 'combined_single.parquet') # 清理临时文件 shutil.rmtree(temp_dir)
该方案的优势是单批写入时仅需处理单批数据的内存,合并阶段PyArrow会高效读取多文件,且严格保留原始数据类型。
内容的提问来源于stack exchange,提问作者Gopala
相关产品推荐
相关产品推荐

