You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 16:06:05