Spark中coalesce方法生成指定数量输出文件失败求助
我来帮你拆解这个问题~你遇到的情况其实是coalesce方法的特性导致的:
为什么coalesce(233, false)只生成4个文件?
coalesce(numPartitions, shuffle=false)的核心逻辑是合并现有分区,它只能减少分区的数量,绝对无法增加分区数——因为它不会触发数据洗牌,只是把多个小分区合并成大分区,或者保持现有分区不变(当目标分区数大于当前分区数时)。
你的原始JavaRDD<String>应该只有4个分区,所以当你调用coalesce(233, false)时,Spark发现目标分区数比当前多,直接忽略了你的请求,还是用原来的4个分区保存文件,自然就只生成4个文件了。
而repartition(233)能正常工作,是因为它本质上等同于coalesce(233, true)——会强制触发shuffle,重新将数据分配到233个分区中,所以能生成指定数量的文件,但代价是数据洗牌的开销。
无shuffle实现生成233个文件的方案
如果不想触发shuffle,你需要从源头控制RDD的分区数,而不是在已有RDD上尝试用coalesce扩容:
如果你的RDD是从文件读取的
在读取文件时直接指定分区数为233,这样Spark会直接将数据分成233个分区,后续不需要任何分区调整操作:JavaRDD<String> trainDataFeatures = sc.textFile("你的输入路径", 233); trainDataFeatures.saveAsTextFile(outputPathTrainFeatures);这样每个分区的数据量大概是
16310/233≈70条,正好符合你的需求,而且完全不会触发shuffle。如果你的RDD是通过其他转换操作生成的
检查生成这个RDD的上游操作,看看能不能在某个步骤就设置足够的分区数。比如如果是groupByKey这类支持指定分区数的操作,直接在那里设置为233,这样后续RDD就会拥有目标数量的分区,无需再调整。
总结
coalesce无shuffle模式只能缩减分区,无法扩容分区;- 无shuffle生成指定数量文件的关键是让RDD从一开始就拥有目标数量的分区;
- 如果实在没办法从源头调整,那只能接受用
repartition(233)(即带shuffle的coalesce)来实现,这是唯一能在已有RDD上扩容分区的方式。
内容的提问来源于stack exchange,提问作者s1nned

