Apache Flink读取HDFS文件失败及writeAsText输出异常求助
解决Flink读取HDFS文件并写入本地的两个问题
一、任务失败的原因与解决
你最初遇到的任务失败问题,本质是集群可用的Task Slots数量不足以支撑作业的并行度。Flink默认的作业并行度会根据集群配置自动设定,如果你的集群初始Task Slots设置过少(比如默认仅1个),而作业需要多个并行任务来处理数据,就会出现任务无法分配、最终失败的情况。
你通过增加Task Slots解决问题的思路完全正确——这给作业提供了足够的资源来分配并行任务,让作业能正常执行。
二、writeAsText生成目录而非单个文件的原因
至于你提到的writeAsText没有生成预期的单个文件,反而创建了同名目录和多个分片文件,这是Flink分布式计算的正常行为:
- Flink的DataSet API是并行执行的,每个并行子任务都会独立输出自己处理的数据分片。当你调用
writeAsText("/tmp/hdfs_file.txt")时,如果作业的并行度大于1,Flink会把这个路径当作目录名,然后在里面生成多个以part-开头的文件(对应每个并行任务的输出)。 - 你看到的8个文件,正好对应你作业的并行度为8——其中6个为空,说明这6个并行任务没有分到数据,另外两个分到了数据所以有内容。
如果想要生成单个文件,你可以通过以下两种方式调整:
- 设置输出的并行度为1:
在writeAsText之前调用setParallelism(1),强制让输出阶段只用一个并行任务:lines.writeAsText("/tmp/hdfs_file.txt").setParallelism(1); - 先将数据合并为一个分区:
对DataSet调用repartition(1)(适合大数据量场景),或者用collect()(会把数据拉到客户端,仅适合小数据量):// 大数据量场景的合并输出方式 lines.repartition(1).writeAsText("/tmp/hdfs_file.txt");
需要注意的是,把并行度设为1会牺牲输出阶段的并行性能,所以如果是处理大数据量,不建议这么做——分布式输出多个分片才是Flink的最优实践,后续分析也可以直接读取这些分片文件,完全不影响业务逻辑。
内容的提问来源于stack exchange,提问作者Enrique García Rodríguez
相关产品推荐
相关产品推荐

