如何使用Fugue合并Parquet文件中的DataFrame并重新分区?
合并Parquet文件并重新分区的实现方案
需求说明
现有两个Parquet文件,需将其中的DataFrame行级合并,再按指定分区数输出为新的Parquet文件。
方法一:Pandas(适合小数据集)
Pandas适合处理小规模数据,合并后可通过拆分数据实现多分区输出:
import pandas as pd import numpy as np # 读取两个Parquet文件 df1 = pd.read_parquet("/tmp/1.parquet") df2 = pd.read_parquet("/tmp/2.parquet") # 行级合并并重置索引 merged_df = pd.concat([df1, df2], ignore_index=True) # 指定分区数并输出 partition_count = 2 # 自定义分区数 for idx, part_df in enumerate(np.array_split(merged_df, partition_count)): part_df.to_parquet(f"/tmp/merged_part_{idx}.parquet")
注:Pandas默认生成单个文件,上述代码通过np.array_split将数据拆分后分别保存,实现多分区效果。
方法二:PySpark(适合大数据集)
PySpark原生支持分布式数据处理,可直接指定分区数写入,效率更高:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("MergeParquet").getOrCreate() # 读取Parquet文件 df1 = spark.read.parquet("/tmp/1.parquet") df2 = spark.read.parquet("/tmp/2.parquet") # 行级合并(union与unionAll效果一致,后者为旧API) merged_df = df1.union(df2) # 指定分区数并写入新Parquet文件 partition_count = 2 # 设置目标分区数 merged_df.repartition(partition_count).write.parquet("/tmp/merged_parquet")
注:
repartition会重新洗牌数据,确保分区数据均匀,适合调整任意分区数;- 若仅需减少分区数,可使用
coalesce(不洗牌,性能更优)。
内容的提问来源于stack exchange,提问作者Hakan Baba
相关产品推荐
相关产品推荐

