
前一阵子有个需求把我逼着把文件读取这块彻底重写了一遍系统日志既要补算过去几天落在目录里的历史文件又要在新文件进来后立刻进入实时链路。以前的做法是拆成两个Flink任务一个用离线读目录一个用StreamingFileSource盯目录两套作业的解析逻辑还得各自维护改个字段就要同步改两边实在受够了。后来我换成 Flink 1.15 之后的 FileSource用 TextLineInputFormat 定义文本读取格式通过一个 Source 就同时覆盖批式读存量、流式盯增量目录持续监控也由枚举器自动完成。这篇就把这套方案的选型思路、TextLineInputFormat 的切行原理、FileSource 批流一体的具体写法以及我在生产环境里实测踩过的坑完整梳理一遍。1. 需求冲击下的选型为什么放弃老的StreamingFileSource1.1 一个任务想要同时做“补数”和“实时”当时场景很典型数据链路的上游每天把脱敏日志写入一组带日期的目录比如/data/logs/incoming/2025-04-01/*.log。业务方要求跑批的时候把最近三天的全部日志重新解析一遍生成离线指标跑实时的时候只要新文件落到目录里就立刻读取并进入后续计算不再等调度。同一个文件格式、同一套字段规范如果拆成两个任务解析逻辑必然要复制一份。哪怕一开始封装成公共函数时间长了也会出现“离线改了三行实时没同步”这类事故。所以我的目标是同一份代码通过一个开关切换批流读取同一个目录的同一套文件尽量不复制逻辑。这就是 FileSource 最吸引我的地方。它把“文件怎么切分”“行怎么读”“进度怎么记”全部收拢到统一的 Source 接口里而上层只要换RuntimeExecutionMode.BATCH还是STREAMING就能得到完全不同但符合直觉的行为。1.2 老API到底别扭在哪在 FileSource 成为主力之前流式读文件主要靠StreamingFileSource它配合StreamingFileInputFormat再由ContinuousFileMonitoringFunction去轮询目录。这套组合的问题在于职责分散监控函数负责发现文件reader operator 负责读协调逻辑散落在多个算子之间排障的时候要同时看几个组件状态语义乱断点续读的 offset 管理分散在ContinuousFileReaderOperator内部和 Source API 的新模型不统一批流割裂它本质是“流式环境下模拟读一批文件”拿到 BATCH 模式下做一次性读取并不顺手官方维护态度明确StreamingFileSource很早就被标记弃用新 feature 都往 FileSource 上走。我编译老代码时经常看到 deprecation 警告这种“能跑但官方不推荐”的状态最难受。后来干脆基于 FileSource 统一重写。老方案和新方案我整理过一个对比放在表格里更直观对比项StreamingFileSource 老方案FileSource 新方案接口归属流式专用核心散落多个算子统一 Source 接口批流共用文件发现ContinuousFileMonitoringFunction 轮询SplitEnumerator 周期发现断点状态分散在 reader operatorFileSourceSplit 统一记录读取格式StreamingFileInputFormatFileRecordFormat / TextLineInputFormat官方状态废弃新开发主力1.3 批流一体的本质是执行模式不是Source偷偷切逻辑很多人以为“批流一体”是 Source 自己判断“现在按批读现在按流读”其实不是。FileSource 在构建时会根据你调用的 builder 分支决定自己是一个有界BOUNDED还是无界CONTINUOUS_UNBOUNDED的源然后 JobGraph、调度器、网络缓冲、checkpoint 语义全部围绕这个Boundedness工作。RuntimeExecutionMode.BATCH 一次性 builder源读完当前所有 split 后 job 正常结束RuntimeExecutionMode.STREAMINGmonitorContinuously源会持续保持执行SplitEnumerator 不断发现新文件任务不会退出。底层都是FileSourceSplit在流转有界模式下 split 是一次性分配完的无界模式下 split 是不断追加进待分配队列的。理解了这一层后面所有调参都有方向感。2. TextLineInputFormat的切行逻辑从FileSplit到行缓冲2.1 文件是怎么被切成split的我最早用 FileInputFormat 那一套老接口时就被“split”这个概念绕了一下。通俗点讲一个文件物理上是一堆字节Flink 为了让多个并行子任务同时读就把字节区间切成一段一段每段就是一个FileSourceSplit里面有文件路径、起始 offset、长度、还包括一个 split ID。TextLineInputFormat 的作用是把每个 split 对应的字节区间转换成一行一行的字符串。它内部用LineReader按字节读取遇到\n就认为一行结束同时兼容\r\n。这中间最关键的点是split 的边界不会那么凑巧落在行边界上。比如一行内容跨越了两个 split 分界点如果每个 split 简单读到自己的 length 就停行就会被拦腰截断。Flink 对这类情况的处理是当前 split 读到物理边界后如果发现行还没结束会继续往后读直到读到换行符为止。也就是说单个 split 的实际读取区间可能比声明的 offset length 略大一点。提示这个“跨 split 续读”是 TextLineInputFormat 已经实现好的不需要你额外处理。前提是你别自己在外面套一层自定义字节流那样很容易打破 split 语义。2.2 最后一行没有换行符会丢这是我在实测里踩过最典型的坑。场景是用脚本写了一批不带末尾换行的文本文件然后让 FileSource 去读。结果 Sink 收到的条数比文件真实行数少一行最后那行就像凭空消失了一样。原因在于 LineReader 在读到文件末尾时如果最后一个字节不是换行符它会认为当前没有完整行可返回于是直接结束最后一个半截行就被丢弃了。这不是 Flink 的 bug而是“按行读取必须依赖行结束符”的自然结果。解决办法也不复杂约定所有落盘文本文件都以\n结尾。尤其是上游用echo -n或者程序write时别省最后那个换行。注意不同 Flink 小版本对“EOF 前最后一个字符”的处理略有差异但我的建议是永远不要依赖“最后一行没有换行也能被读出来”这个行为。生产环境里文件格式的约定越明确踩坑越少。2.3 字符编码与BOM的连带问题TextLineInputFormat 默认按 UTF-8 解码字节流这对大多数日志场景够了。但它不会帮你处理 UTF-8 BOM。如果上游用带 BOM 的方式写文件BOM 会以不可见字符出现在第一行第一个字段里。比如 CSV 解析后第一列值变成了\ufeffcustomer_id后续equals(customer_id)永远为 false。排查起来非常隐蔽因为打印日志时 BOM 不显示只有看字段长度或者用十六进制看字节才发现。我现在的做法是约定写入侧统一不带 BOM同时在下游解析时额外做一个“首个字段去掉\ufeff前缀”的兜底双保险。不建议在 TextLineInputFormat 这一层纠结因为它也没有专门处理 BOM 的参数。另外如果文件是 GBK 这类非 UTF-8 编码也不建议硬塞给 TextLineInputFormat。同一份代码里保持单一编码比在 Source 层做编码转换要省心得多。跨编码转换一旦遇到半角全角混排很容易出现乱码且不易收敛。3. FileSource的批流一体实现从SourceBuilder到运行模式3.1 一次性读和持续监控对应两个builder分支FileSource 构建入口非常简单静态方法FileSource.forRecordStreamFormat(recordFormat, path)会返回一个 builder。之后有两种收尾方式.build()一次性读取当前路径下已存在的文件读完即完成。适合补数、批计算。.monitorContinuously(Duration.ofSeconds(...)).build()按固定周期扫描目录新文件会自动生成新的 split 进入消费队列。适合实时盯增量目录。我把这段封装成了一个工具方法调用方只用传一个布尔开关import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.connector.file.src.FileSource; import org.apache.flink.connector.file.src.reader.TextLineInputFormat; import org.apache.flink.core.fs.Path; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.environment.RuntimeExecutionMode; import java.time.Duration; public class TextFileSourceFactory { public static FileSourceString buildFileSource(String path, boolean streaming) { TextLineInputFormatString format new TextLineInputFormat(); FileSource.FileSourceBuilderString builder FileSource.forRecordStreamFormat(format, new Path(path)); if (streaming) { return builder.monitorContinuously(Duration.ofSeconds(10)).build(); } return builder.build(); } public static void main(String[] args) throws Exception { boolean streaming args.length 0 --streaming.equals(args[0]); String path /data/logs/incoming; FileSourceString fileSource buildFileSource(path, streaming); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(streaming ? RuntimeExecutionMode.STREAMING : RuntimeExecutionMode.BATCH); DataStreamSourceString lines env.fromSource( fileSource, WatermarkStrategy.StringnoWatermarks().withIdleness(Duration.ofSeconds(20)), text-file-source); lines.flatMap((String line, org.apache.flink.util.CollectorString out) - { String trimmed line.trim(); if (!trimmed.isEmpty()) { out.collect(trimmed); } }).name(filter-empty-line) .print(); env.execute(text-file-source-demo); } }运行方式上我习惯用代码里env.setRuntimeMode强制指定因为这样提交脚本统一不需要每个环境都带额外参数。如果想把模式交给命令行也可以去掉这行然后用./bin/flink run -d -c TextFileSourceFactory \ -Dexecution.runtime-modebatch our-job.jar两种方式等价。3.2 split分配与并行度的实际关系FileSource 的并行度决定的是同时打开多少个 spli t 的读取通道而不是直接控制“每个文件读多快”。当目录里的 split 数量大于并行度时空闲的 reader 会不断从待分配队列里领新的 split反过来如果 split 数量少而并行度大部分子任务会一直空闲这不是故障只是资源浪费。所以并行度设定的合理参考是文件块数量。一个 128MB 的文件可能会被拆成 1 个或几个 split如果你的目录里有 50 个大文件并行度给到 20 左右通常就能比较充分地压满磁盘和网络。给的过高反而会让每个 split 的启动开销占比变大。在 WebUI 上可以观察 Source 端的numRecordsInPerSecond和每个 subtask 当前打开的 split 情况。如果发现某个 subtask 一直空闲另一个忙到反压就要考虑是不是 split 粒度太粗或者并行度设置不合适。3.3 批流共用一个Source的代价看到这里你可能会想既然一个 Source 这么方便是不是所有文件读取都统一用 FileSource 就行也不尽然。FileSource 适合“结构化良好的文件目录”比如日志按天滚动、对象按小时落盘、文件写完就立即 rename。如果是高频小文件、乱序覆盖写、单文件无限增长这类场景FileSource 的 split 枚举机制处理起来会比较吃力Kafka 或者其他消息中间件会是更合理的选择。批流共用一个 Source 的真正收益是把“同一份解析逻辑”保持住了而不是让 FileSource 在所有文件场景里都成为万能选项。4. 目录持续监控的工程细节扫描周期、临时文件与文件生命周期4.1 扫描周期不是越短越好monitorContinuously(Duration.ofSeconds(10))里的时间是 SplitEnumerator 定期扫描目录的间隔。它直接决定了“新文件多久会被发现”。我见过有人把扫描周期压到 1 秒理由是“实时性越高越好”。但实际目录扫描是去listStatus如果文件数量大每次列出目录本身就有开销太频繁反而挤占读文件的带宽。对于本地磁盘和 HDFS10 秒左右是比较稳妥的默认值如果目录在对象存储上还要考虑 List 请求的成本和延迟我一般会给到 30 秒以上。另一个要理解的点扫描间隔只是延迟的一部分。文件写入本身也需要时间上游如果持续几秒才写完整那么即使扫描发现得再快也还是要等文件处于可读状态。扫描周期、文件写入速度、下游消费速度三者需要一起权衡不能只盯着其中一个参数。4.2 先写临时文件再rename几乎是最重要的约定目录持续监控有个天然问题如果某个文件正在被上游写入而扫描恰好发生在这期间FileSource 就可能只读到文件当前已有的一部分剩下的部分不会因为你下次扫描而重新读——因为这只是一个被识别过一次的 split不会因为文件长度变化再生成新的 split。解决办法不是改 Flink 配置而是改写入侧的落盘习惯# 错误示范直接在目标目录里写 some-generator --output /data/logs/incoming/app.log # 正确做法先写临时文件完整写完再原子rename some-generator --output /data/logs/incoming/app.log.tmp mv /data/logs/incoming/app.log.tmp /data/logs/incoming/app.log只要上游遵守“先.tmp后mv”FileSource 扫描到的都是完整文件读取逻辑就变得非常简单。这个约定也帮我们挡住了绝大部分“半截文件”“读到一半文件还在涨”之类的诡异问题。4.3 子目录递归与文件名过滤FileSource 的目录列举是递归的所以你可以放心让上游按天或按小时建子目录例如/data/logs/incoming/2025-04-01/。这样历史补数只需要把路径指到上一个层级增量监控同一路径也能覆盖每天新增的子目录。如果目录里混有.tmp、.ok、.csv、.log等不同后缀就需要过滤。TextLineInputFormat 支持传入一个 PathFilter 参数TextLineInputFormatString format new TextLineInputFormat( path - path.getName().endsWith(.log) );我一般会把过滤规则放在这里而不是在下游 flatMap 里判断。尽早丢掉无关文件可以减少无效 split 的生成也避免待分配队列被没用的文件占满。4.4 “追加写”场景要慎重有人会想日志不就是追加写吗我监控一个文件文件不断增长我不断地读新追加的行这不就是 FileSource 的典型场景吗实际上 FileSource 并不是为“单个文件无限追加”设计的。它的核心模型是“文件一旦被纳入 split就按 offset 和 length 读取”如果文件在读取过程中还在增长当前 split 可能会把增长的部分也读走但这个行为依赖时机不可控。对已经消费完并记录过 checkpoint 的文件后续再追加的内容通常不会被重新读取。所以我的建议是日志链路如果是高频追加型优先考虑消息系统或者“按时间戳滚动文件”方案如果一定要用目录监控就以“新文件不断产生”为前提设计而不是“单文件无限追加”。5. 生产环境里实测踩过的坑与验证方法5.1 四个高频问题的现象与处理我整理了一个实战踩坑清单下面这几个是我在现网里真正遇到过的按频率排序场景现象原因处理方案最后一行缺失Source 输出条数比实际行数少 1最后一行无换行符LineReader 认为行不完整写入侧约定文件末尾必须有\n第一行字段带BOMCSV 首列出现\ufeff上游文件带 UTF-8 BOM写入侧去掉 BOM下游解析兜底剔除恢复时文件不存在checkpoint 恢复抛文件找不到日志目录按保留周期清理了旧文件文件保留时间大于任务保留时间批量小文件阻塞Source 端 split 堆积反压明显扫描周期太快小文件太多降低扫描频率下推过滤合理并行度5.2 水印与空闲源的问题读取文本文件时每行默认没有时间戳概念。如果你后续要做事件时间窗口就必须自己解析行内字段并指定 TimestampAssigner不能依赖 Source 自动生成水印。这里最容易被忽略的是“空闲源”问题当目录暂时没有新文件时Source 不产生任何数据水印就会一直停在上一个时间点后续事件时间窗口可能长期不触发。解决办法是用WatermarkStrategy的withIdleness参数给源一个空闲容忍时间WatermarkStrategyLogRecord strategy WatermarkStrategy .LogRecordforBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((record, ts) - record.getEventTime()) .withIdleness(Duration.ofSeconds(30));这样即使目录几分钟没有新文件水印也能依靠空闲检测继续推进下游窗口不会卡死。5.3 一个方便的本地验证方法我每次改完目录监控逻辑都会先在本地起一个小作业再用脚本模拟“持续产生完整文件”的上游验证延迟和完整性#!/usr/bin/env bash PROBE_DIR/tmp/flink-fs-probe mkdir -p $PROBE_DIR for i in $(seq 1 20); do echo line-$i $PROBE_DIR/app.log.tmp sleep 3 mv $PROBE_DIR/app.log.tmp $PROBE_DIR/app.log done注意每轮都重新生成一个app.log而不是把几十行写进同一个文件。因为 FileSource 对 rename 后的“新文件”识别非常可靠但对同一个文件不断追加就不容易把控。本地验证时可以让任务打印每个条目的来源文件名和 offset这样有没有丢行、有没有重复、延迟多久一眼就能看出来。判断有没有重复和丢行的简单办法让本地下游用一个MapStatesplitId, offset记录每条数据的来源 split 和 offset最后聚合对比文件总行数。FileSourceSplit 上报的 offset 是文件里的字节位置不是行号所以更推荐直接用“行内容去重计数”来验证。5.4 别把TextLineInputFormat的包路径搞混写代码的时候还有一个非常容易踩的坑老版org.apache.flink.api.common.io.TextLineInputFormat是批处理那套的遗留实现新版org.apache.flink.connector.file.src.reader.TextLineInputFormat才是配合 FileSource 使用的。两个类同名但构造器和接口完全不同。我刚开始迁移时IDE 自动补全给我补成了老的编译也能过但和FileSource.forRecordStreamFormat的类型完全对不上报了一堆泛型错误。后来凡是看到connector.file.src这个包路径才算真正对上了。5.5 checkpoint恢复与目录生命周期FileSource 在做 checkpoint 时会把待分配 split 和正在读取 split 的 offset 都存下来。恢复时SplitEnumerator 要重新找回这些文件。如果文件已经被清理恢复就可能失败。这一点决定了目录的保留策略清理日志文件的保留天数必须长于任务允许回溯的最长时间也要覆盖任务挂掉到恢复之间的间隔。我现在的习惯是监控目录里的文件至少保留 7 天任务重启后即使要回放一两天的数据也还有底。最后想说的是FileSource 这套方案解决的痛点是“同一份文件解析逻辑既跑补数又跑实时”。它不是万能的文件摄入方案面对单文件高频追加、超大目录低延迟这类场景还是要换更合适的工具。但只要你的业务是“一批有规律落盘的文件需要稳定地被读走”TextLineInputFormat FileSource 这套组合我实测下来是很稳的选择。