Spark on Yarn集群运行疑问:WordCount与SparkSQL相关问题
Hey there, let’s clear up why your WordCount job is spitting out two part files (part-00000 and part-00001)—this is totally normal behavior tied to how Spark handles parallel processing and data partitioning!
核心逻辑:输出文件数 = 结果数据集的分区数
Spark writes output files one per partition of your final RDD or DataFrame. When your WordCount job finishes, the number of part files directly matches the number of partitions in your result dataset. This is how Spark maintains parallelism even during the final write step.
结合你的集群配置来分析
Your cluster has 4 worker nodes (each with 16 cores and 12GB RAM), but Spark doesn’t automatically set the partition count to match the number of workers. The final partition count depends on a few factors:
- The parallelism settings for shuffle stages (like
spark.default.parallelismfor RDDs orspark.sql.shuffle.partitionsfor DataFrames) - The size of your input data (Spark may reduce partitions if data volume is small to avoid tiny files)
- How YARN allocated executors for your job (maybe only 2 executors were used for the final stage)
In your case, the WordCount job ended up with 2 partitions in the result set, hence two part files.
如何控制输出文件数量?
If you want a single output file (though this isn’t recommended for large datasets, as it kills parallelism), you can use coalesce() or repartition() to adjust the partition count before writing:
// For RDD-based WordCount wordCountResult.coalesce(1).saveAsTextFile("hdfs://your/output/path") // For DataFrame-based WordCount (using SparkSQL) wordCountDF.coalesce(1).write.mode("overwrite").text("hdfs://your/output/path")
coalesce(1)is more efficient here because it avoids a full shuffle (it just merges partitions without re-distributing data)repartition(1)will trigger a shuffle, which is unnecessary unless you need to rebalance data
为什么不是4个输出文件?
It’s likely your job didn’t utilize all 4 worker nodes for the final stage. Spark optimizes resource usage based on data size—if your input wasn’t large enough to justify 4 parallel partitions, it would scale back to a smaller number to avoid overhead from managing too many tiny files.
内容的提问来源于stack exchange,提问作者Essex

