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

Pyarrow使用S3文件系统写入分区Parquet数据集时出现数据覆盖问题

PyArrow 写入分区Parquet到S3的数据覆盖问题说明

本地分区Parquet写入逻辑

当在本地向数据集写入两个Parquet文件时,Arrow可以正确地向分区追加数据。例如,如果使用Arrow按指定列对两个文件进行分区,首次写入带分区规则的Parquet文件时,Arrow会生成分区列每个唯一值对应的子文件夹目录结构。写入第二个文件时,Arrow可以智能地将数据写入对应分区,因此如果两个文件的分区列存在共同取值,该共同取值对应的子文件夹下会出现两个独立文件。

本地写入代码示例:

df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_07.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(table, str(base +  "parquet_dataset_partition_combined"), 
                    partition_cols=['PartitionPoint'])

df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_08.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(table, str(base +  "parquet_dataset_partition_combined"), 
                    partition_cols=['PartitionPoint'])

运行结果:
结果示意图1
结果示意图2

上述示例中分区列的基数为2(取值A和B),因此会生成2个分区文件夹,且PartitionPart=A子文件夹下存在2个文件,因为actual_07和actual_08两个文件都有归属PartitionPart=A分区的数据。

S3写入覆盖问题复现

完全相同的代码使用S3作为文件系统时无法实现上述追加效果,复现代码如下:

from pyarrow import fs

s3  = fs.S3FileSystem(region="us-east-2")

df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_07.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(table, "parquet-storage", 
                    partition_cols=['PartitionPoint'],
                    filesystem=s3)

df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_08.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(table, "parquet-storage", 
                    partition_cols=['PartitionPoint'],
                   filesystem=s3)

实际运行时第二条写入语句会覆盖S3中的已有数据,任意时刻每个PartitionPart=A文件夹内始终只存在1个文件。

问题原因与解决方案

这不是S3文件系统的已知限制,是PyArrow写入逻辑的参数默认配置导致的问题:

  • 本地文件系统写入时,PyArrow会自动检测已有文件,自动生成不重复的输出文件名,因此不会出现覆盖
  • S3作为对象存储没有原生目录结构,旧版本PyArrow处理S3写入时默认使用固定的文件名模板,两次写入相同分区的文件名完全一致,而S3对相同路径(对象Key)的写入默认执行覆盖操作,因此后写入的数据会覆盖之前的文件

修复方式非常简单,写入时显式指定不重复的文件名模板,同时开启追加模式即可,修改后的代码如下:

from pyarrow import fs
import uuid

s3  = fs.S3FileSystem(region="us-east-2")
# 写入第一个文件
df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_07.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(
    table, 
    "parquet-storage", 
    partition_cols=['PartitionPoint'],
    filesystem=s3,
    append=True, # 开启追加模式
    # 每次写入生成带唯一UUID的文件名,避免重名覆盖
    basename_template=f"part-{uuid.uuid4().hex}-{{i}}.parquet"
)
# 写入第二个文件
df = pd.read_parquet('~/Desktop/rough/parquet_experiment/actual_08.parquet')
table = pa.Table.from_pandas(df)
pq.write_to_dataset(
    table, 
    "parquet-storage", 
    partition_cols=['PartitionPoint'],
    filesystem=s3,
    append=True,
    basename_template=f"part-{uuid.uuid4().hex}-{{i}}.parquet"
)

如果你使用的是PyArrow 15.0以上版本,也可以用existing_data_behavior='overwrite_or_ignore'参数替代append=True,兼容性更好。


内容的提问来源于stack exchange,提问作者Abhishek Malik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 16:15:03