AWS EMR Step运行Spark程序时HDFS合并命令执行异常求助
你遇到的这个问题我太熟悉了——在EMR Step模式下用shell管道合并HDFS文件,结果数据全打到日志里,任务无限跑。本质原因是EMR Step会捕获任务的所有stdout输出,而你用hadoop fs -cat | hadoop fs -put -的管道逻辑,在Step环境里没法正确把cat的输出传给put的stdin,反而把所有数据都输出到日志流里,导致任务永远在打印数据,没法完成。
下面给你几个可行的替代方案,按推荐程度排序:
1. 使用Hadoop官方的getmerge命令(最推荐)
Hadoop自带的getmerge工具就是专门用来合并HDFS文件的,它会把指定目录下的文件合并到本地文件,再由你上传回HDFS,全程不会把数据打到stdout。
修改你的Scala代码为:
// 构造getmerge命令:先合并到本地临时目录,再上传回HDFS,最后清理本地文件 val mergeCmd = s""" hadoop fs -getmerge ${HDFSOutPath}/part* /tmp/${fileName}.csv && hadoop fs -put /tmp/${fileName}.csv ${HDFSOutPath}/ && rm -f /tmp/${fileName}.csv """.trim.replaceAll("\\s+", " ") // 执行命令 mergeCmd.!
这个方案的好处是:
- 依赖Hadoop原生工具,稳定性高
- 不会产生大量stdout输出,避免日志爆炸
- 性能比
coalesce(1)好很多,因为不需要shuffle数据
2. 使用Spark 3.3+的saveAsSingleFile API
如果你的Spark版本是3.3及以上,可以用官方新增的saveAsSingleFile方法,它是专门优化过的单文件输出,比coalesce(1)性能好太多(不会把所有数据强行拉到一个节点)。
示例代码:
import org.apache.spark.sql.SaveMode // 直接保存为单个文件(会生成part-00000格式的文件) df.write .mode(SaveMode.Overwrite) .option("header", "true") // 如果需要输出表头可以加上 .saveAsSingleFile(HDFSOutPath + "/single_temp") // 用HDFS命令重命名为目标文件名,并清理临时目录 val renameCmd = s"hadoop fs -mv ${HDFSOutPath}/single_temp/part-00000* ${HDFSOutPath}/${fileName}.csv && hadoop fs -rm -r ${HDFSOutPath}/single_temp" renameCmd.!
注意:saveAsSingleFile会在指定目录下生成一个单独的part文件,你需要额外做一步重命名操作。
3. 用Scala直接操作HDFS API(最灵活)
如果不想依赖shell命令,可以直接用Hadoop的FileSystem API在Scala代码里完成文件合并,全程在JVM内执行,完全避免shell管道的问题。
示例代码:
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.hadoop.conf.Configuration import java.io.BufferedInputStream val conf = new Configuration() val fs = FileSystem.get(conf) val targetPath = new Path(s"${HDFSOutPath}/${fileName}.csv") val outputStream = fs.create(targetPath) // 遍历所有part文件(过滤掉_SUCCESS等无关文件) val partFiles = fs.listStatus(new Path(HDFSOutPath)) .filter(status => status.isFile && status.getPath.getName.startsWith("part-")) .map(_.getPath) // 逐个读取并写入目标文件 partFiles.foreach { partPath => val inputStream = new BufferedInputStream(fs.open(partPath)) val buffer = new Array[Byte](4096) var bytesRead = inputStream.read(buffer) while (bytesRead != -1) { outputStream.write(buffer, 0, bytesRead) bytesRead = inputStream.read(buffer) } inputStream.close() } outputStream.close()
这个方案的好处是完全可控,不需要调用外部shell命令,适合复杂的合并场景(比如需要过滤某些行、添加额外内容等)。
为什么原来的管道命令不行?
在EMR Step模式下,任务的stdout和stderr会被EMR的日志系统(CloudWatch Logs)实时捕获。你原来的hadoop fs -cat ... | hadoop fs -put -逻辑,在SSH登录集群执行时是正常的——因为shell会把cat的输出直接传给put的stdin。但在Step环境中,任务的stdin可能没有正确关联到管道的输出,导致cat的所有数据都被输出到stdout,进而被日志系统捕获,任务就会一直打印数据,永远无法完成。
内容的提问来源于stack exchange,提问作者matterwaves

