Flink DataStream API使用readFile在批处理模式下仅读取文件夹部分文件问题
问题根因
- 你使用的
StreamExecutionEnvironment.readFile是Flink遗留的流读文件API,在BATCH运行模式下的分片逻辑默认与全局环境并行度绑定:分片数等于当前环境的并行度,每个分片对应一个读任务。 - 本地运行Flink作业时,IDE默认会将作业并行度设为1,因此仅会生成1个读分片,最多只有1个读任务运行,出现仅读取到1个文件的现象。
- 集群运行时你设置了大于1的并行度,会生成对应数量的读分片,因此读取的文件数量与并行度一致。
- 额外注意:你当前的代码同时给
RowCsvInputFormat和readFile方法传入了路径,部分版本Flink会出现路径解析冲突,也可能导致漏读目录下的文件。
修复方案
- 方案1(推荐):替换为Flink流批一体原生的
FileSource连接器,分片逻辑不受全局并行度强制绑定,会根据文件数量和大小自动拆分,不会出现漏读问题,示例代码:
import org.apache.flink.connector.file.src.FileSource import org.apache.flink.connector.file.src.reader.CsvReaderFormat import org.apache.flink.api.common.eventtime.WatermarkStrategy val fileSource = FileSource .forBulkFileFormat( CsvReaderFormat.forRowType(MyClass.getRowTypeInformation()), new Path(salesPath) ) .build() val trans_data: DataStream[MyClass] = env .fromSource(fileSource, WatermarkStrategy.noWatermarks(), "csv文件读取源") .map(x => MyClass.convertFromRow(x))
- 方案2(兼容旧API):调整现有代码的路径传参逻辑,同时显式指定读取操作的并行度和集群配置对齐即可:
val env = StreamExecutionEnvironment.getExecutionEnvironment env.setRuntimeMode(RuntimeExecutionMode.BATCH) // 本地运行时可以手动设置全局并行度和集群对齐,验证效果 // env.setParallelism(4) val trans_data: DataStream[MyClass] = env.readFile( // InputFormat不需要提前指定路径,由readFile统一传入 RowCsvInputFormat.builder( MyClass.getRowTypeInformation() ).build(), salesPath ).setParallelism(4) // 替换为你实际需要的读并行度 .map(x => MyClass.convertFromRow(x))
内容的提问来源于stack exchange,提问作者Luka Zecevic
相关产品推荐
相关产品推荐

