Spark中调用repartition函数时能否限制每个分区的行数?
回答
原生repartition方法不支持直接指定单分区行数上限的参数,但你完全可以通过简单的计算配合repartition实现需求,整个过程不需要落盘,重分区后的数据集可以直接衔接后续的处理逻辑。
具体实现逻辑
- 首先统计待处理数据集的总记录数
- 按照单分区最多100行的要求,向上取整计算需要的总分区数,计算公式为:
所需分区数 = 向上取整(总记录数 / 100) - 将计算得到的分区数传入
repartition()方法完成重分区即可,Spark会将数据随机均匀打散到对应数量的分区中,每个分区的行数不会超过100行
代码示例(PySpark)
# 源数据集 # df = ... max_rows_per_part = 100 total_count = df.count() # 整数向上取整计算分区数 target_part_num = (total_count + max_rows_per_part - 1) // max_rows_per_part # 完成重分区,返回的DataFrame可直接用于后续处理 df_repartitioned = df.repartition(target_part_num) # 后续可直接接任意转换操作,无需落盘 # df_final = df_repartitioned.withColumn(...)
补充说明
- 这种方式得到的分区行数是近似均匀的,因为repartition采用随机哈希分发的逻辑,各分区行数会在100上下小幅波动,不会出现分区数据量严重倾斜的情况。如果需要严格校验分区行数,可以重分区后通过
spark_partition_id()统计每个分区的实际行数做校验,若有少量偏差可以微调分区数。 - 如果你的数据集体量极大,不想因为
count()触发一次全量扫描,也可以通过采样的方式估算总记录数来计算分区数,代价是会有少量分区的行数略超过100行,可根据业务容忍度选择方案。 - 如果你后续还有聚合、join等会触发shuffle的操作,注意这类操作会自动重新分区,如果你需要保持单分区100行的限制,需要在这类shuffle操作之后重新执行上述分区逻辑。
内容的提问来源于stack exchange,提问作者Sonali Bisht
相关产品推荐
相关产品推荐

