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'])
运行结果:

上述示例中分区列的基数为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
相关产品推荐
相关产品推荐

