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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:52:18