凌晨十二点四十,手机上的告警把我从椅子上拽起来:实时数仓的 Flink 作业连续重启三次,全部失败。点开日志一看,堆栈里躺着一个不太常见的异常——DuplicateFileIdException。单独看这个名字,像是写文件时撞了“重号”;可再往后翻,它直接让 Hudi Sink 的 commit 阶段中断,整个 Checkpoint 都没法完成。当时第一反应是偶发问题,大不了从 Checkpoint 恢复,结果第二天同一时间附近它又来了。
这条异常不是简单的网络抖动或者 OOM,它跟 Hudi 的文件管理机制强相关。如果你维护的 Flink 任务也是写 Hudi 表,尤其是用了 Bucket Index、开了异步 Compaction 或 Clustering,那这篇内容值得看完。我会把这个问题的完整排查过程、根因分析、以及最终落地的修复参数全部拆开讲。
1. 故障现场:异常堆栈、任务拓扑和第一时间的取舍
1.1 异常堆栈到底长什么样
我们线上的报错不是每次都在同一个算子位置抛出来的,但关键的异常对象是一致的,大致是这样的内容:
org.apache.hudi.exception.DuplicateFileIdException: duplicate fileId at org.apache.hudi.common.fs.FSUtils.getFileIdFromFileName(FSUtils.java:...) at org.apache.hudi.common.table.timeline.HoodieTimeline.lambda... at org.apache.hudi.sink.commit....(HoodieSink)... Caused by: ...不同 Hudi 版本里的包路径会有点差异,但异常信息里通常会带上具体的fileId,有的版本还会带create_file_id和delete_file_id两个字段。这两个字段非常关键,后半部分我会讲它们分别代表什么。
我们的作业拓扑也不复杂:上游是 MySQL Binlog,通过 Flink CDC 同步到 Kafka,再经过一个 Flink 任务做轻量清洗后以 UPSERT 方式写入 Hudi 表。Hudi 表用的是 MERGE_ON_READ,索引类型为 BUCKET,主键是业务订单号,Sink 并行度是 48,bucket.index.num.buckets配的是 1024。
1.2 为什么第一眼会误判成普通写入失败
看到异常的那一刻,我们最自然的做法是去查是不是磁盘满了、HDFS 是否抖动、上游 JDBC 连接器有没有断连。因为这些因素都会引发写失败,但一般不会导致DuplicateFileIdException这种带有明显业务语义的异常。它指的不只是“写不进去”,而是“有两个写入上下文都认为自己拥有同一个文件组”。
误判的代价是浪费了快一个小时的排查时间,中间还做了几次无效的 Checkpoint 恢复。后来把日志级别临时调到 DEBUG,才看到两个 subtask 在同一批次里向同一个fileId发起了写入。这个信息直接改变了排查方向。
所以如果你们也遇到类似异常,首先不要急着恢复作业、不要急着改并行度。先做三个动作:
- 把完整堆栈、
create_file_id、delete_file_id保存下来; - 从 Flink UI 上把作业的并行度、Checkpoint 配置、算子链截图存档;
- 去 Hudi 表目录下看一眼
.hoodie时间线,确认是否有多个 active instant。
2. 先理解 Hudi 的 fileId 到底是什么:从文件名到文件组分配逻辑
2.1 一个 Hudi 数据文件里最重要的部分是前缀
Hudi 表目录下的数据文件看起来是一串很长的名字,比如:
a1b2c3d4e5f6..._0-1-0_20250101120000.parquet很多人只关注文件格式和时间戳,容易忽略一个关键点:整个文件名最前面的那一串fileId。它相当于这条数据所属“文件组”的门牌号。同一个文件组所有数据文件,包括基础文件.parquet和增量日志文件.log,都共享同一个fileId前缀。
Hudi 把数据组织成“表 -> 分区 -> 文件组 -> 文件”的层级。文件组是并发控制和数据合并的最小单位。两个不同的写入者写两个不同的文件组,互不干扰;但如果两个写入者同时认为同一个文件组归自己管理,Hudi 在提交时就会发现目录里有重复归属关系,直接抛出DuplicateFileIdException。
2.2 Bucket Index 让 fileId 从随机数变成了固定房号
Hudi 老版本里 fileId 更接近随机生成的 UUID,每次写一条新数据都可能生成新的文件组,对大批量流式写入不太友好。后来的 Bucket Index 把思路改成:先用分区字段和主键字段做哈希,求出数据属于哪个 bucket,再让 bucket 对应固定的 fileId。
打个比方:分区是楼层,bucket 是房间。一条订单数据到达后,Flink Sink 会根据哈希值把它分配到一个固定房间;你手里没有房卡,只凭主键就能算出它会去哪个房间。Hudi 读取时也能用同样算法直接定位到房间,不需要全局索引扫描。
好处是写入和读取都能快速定位,坏处也很明显:如果两个写入者同时往同一个房间放东西,系统必须保证它们不会把两套家具互相覆盖。Hudi 有锁机制来协调,但锁只能保证“提交阶段”的一致性,无法保证“写入阶段”所有写入上下文不会重复分配同一个 fileId。
2.3 Flink 写入时的 writeToken 和两层状态
Flink 的 Hudi Sink 不是直接把每条 record 写一个文件,而是每个 subtask 缓存一批数据,再写对应的 bucket。每个 subtask 内部持有writeContext,里面包含文件路径、fileId、writeToken 信息。writeToken是用来区分同一个 fileId 下多个临时文件标记的,正常情况应该全局唯一递增。
Flink 任务做 Checkpoint 时,这些 writeContext 状态会序列化保存。作业重启恢复时,Coordinator 会把这些状态重新分发给新的 subtask。如果恢复时机不对、或者旧状态之间存在重叠,两个 subtask 可能同时持有同一个 fileId 的写入上下文,接着就会在同一批次写出同一个 fileId 的数据文件,提交时自然触发重复异常。
3. 四步定位:从 Flink 状态、Hudi 时间线到数据目录交叉验证
3.1 第一步:先把“有多少个作业写同一张表”查明白
这是最便宜、最优先要确认的。很多时候DuplicateFileIdException之所以偶发,是因为有两个 Flink 作业在同一时间段写同一张 Hudi 表,平时数据分布不重叠,某一天某批 key 同时落在同一个 bucket,就撞上了。
我当时的做法是把所有 Flink 作业列表拉出来,逐个过滤日志里的 table path:
curl -s "http://flink-rest:8081/jobs/overview" | jq -r '.jobs[] | "\(.id) \(.name)"'再结合日志搜索关键字hudi_table_path,结果发现线上除了主实时任务,还有一个数据补数任务也在写同一张表。补数任务的并行度是 16,数据是从另一个 Kafka Topic 来的。按理说业务主键一致,分区也一致,但两个任务在 Hudi 层面完全是两个客户端,它们没有共享写锁和写上下文。
3.2 第二步:看.hoodie时间线,找同一个 instant 里的重复归属
Hudi 的每次提交都对应 timeline 里的一个 instant。如果两个客户端同时提交,它们各自生成自己的 instant,时间一前一后,表面上看起来是连续的,但在文件系统层面可能已经留下“同一个 fileId 被两个临时文件引用”的影子。
到 Hudi 表目录下执行:
hdfs dfs -ls /warehouse/hudi/table_name/.hoodie/commits、deltacommits、inflight开头的文件会比较直观。我们找到一条 inflight 的 delta commit 和另一条 completed 的 commit 时间戳非常接近,它们都记录到了同一个 fileId。这就是“看起来都成功,合到一起就重复”的铁证。
3.3 第三步:按 fileId 反查是哪个算子写出来的
拿到异常信息里的 fileId 后,不要只看堆栈,要去 Flink Web UI 的 TaskManager 日志里搜索它,找到是什么时间、哪个 subtask 写出的文件名,并和任务拓扑中的算子编号对应。
我们定位到一个很有意思的现状:两个作业各自只创建了自己的 Flink sink,但 Hudi 表 Bucket Index 的 bucket 数量是 1024,而真实数据量并不大,大多数 bucket 在某一整段小时内没有数据。所以两个作业虽然各有各的并行度,但最终在 commit 阶段,都去抢那几个“热门 bucket”对应的 fileId。某个补数作业恰好和实时任务写到了同一批 key 时,异常就会稳定出现。
3.4 第四步:检查 Compaction / Clustering 等异步服务是否在抢同一个文件组
如果你确认没有多个作业写同一张表,那第二步的重点要放在异步服务上。很多 Hudi 表配置了compaction.async.enabled = true或者clustering.async.enabled = true。Compaction 会把一个文件组的 base 文件和 log 文件合并成新的 base 文件,文件组还是同一个,fileId 不变;Clustering 则会把小文件重新布局,也可能涉及 fileId 的退役和新建。
问题出在这类异步服务同样是独立的客户端,它们不一定和实时写任务的锁上下文完全串行。极端情况下,Compaction 正在创建新 base 文件时,实时任务也在向同一个文件组的 log 文件写入,提交时出现“既有 create_file_id,又有 delete_file_id,两者指向同一个 fileId”的冲突。
create_file_id: b7c4f9a1... delete_file_id: b7c4f9a1...如果看到两个字段一样,那基本就是 Table Service 和主写入管线的时间窗口重叠。如果 create 出现了两个不同的 fileId,则是多个写入客户端各自声明了新文件组,协调器判断归属重复。
4. 根因整理:DuplicateFileIdException 背后最常见的三种成因
4.1 多作业并发写同一张表,且热门 key 的哈希分布过于集中
这是最常见也最容易忽略的。单独看每个 Flink 作业都是健康的,但它们都是独立的 Hudi WriteClient。为了缓解这个问题,经常会加大hoodie.bucket.index.num.buckets,但这只增加了房间数量,并没有解决“多个房东收同一套房租”的冲突。
如果业务上不允许拆表,那唯一的稳妥方案就是把多路写入合并到同一个 Flink 作业的同一个 Sink 算子里。这个不只是在架构图上“看起来单条链路”,而是要让 Flink 算子内部共用同一套 writeContext 状态,才能从根本上避免两个客户端各自维护一份文件归属表。
4.2 Checkpoint 恢复时重复分配写入状态
作业因为其他原因做一次正常的 Checkpoint 恢复,理论上不会出问题,但如果上游 Flink 版本和 Hudi connector 版本存在已知 bug,恢复时会把旧的 writeContext 状态重复广播到多个 subtask,导致同一个 fileId 被两个新的 subtask 同时持有。
这类问题尤其容易出现在“从 Savepoint 恢复之后修改了并行度”的场景里。并行度调整后,文件状态重新 hash 分配,原来的边界被打散,某些 fileId 就会同时出现在两个 subtask 中。所以如果没有明确原因,生产任务尽量不要随便调整 Hudi Sink 的并行度。
4.3 异步 Table Service 与主写入链路的时间窗口重叠
Compaction 和 Clustering 这类服务本身是后台增强能力,不是核心写入链路。Hudi 社区已经在做并发控制,但到了生产环境,只要有多个活跃客户端同时操作同一张表,就有可能在时间窗口重叠时触发 fileId 归属冲突。
我建议把这类服务的自动化程度先降下来。不是要我们彻底放弃它们,而是把触发时机从“任意时刻自动跑”改成“固定窗口统一跑一次”。
5. 修复方案与防御性参数:我们最终采用的组合拳
5.1 主线修复:把两路写入合并成单作业
我们的核心修复,是把补数作业并入实时任务。补数数据先写入 Kafka 同一个 Topic,主任务在 Source 端做UNION ALL,经过同一个转换逻辑后由同一个 Hudi Sink 写表。这样整个表只有一个活跃的 FLink WriteClient,fileId 的状态由同一个 Operator Coordinator 分配,冲突自然消失。
如果你暂时不能合并,那至少要保证同一时刻只有一个作业处于“写提交”状态。这可以通过外部调度锁实现,比如在提交 Flink 作业之前先占用一个 ZooKeeper 节点,拿到锁才允许提交。
5.2 参数层收口:锁定 bucket、关闭异步、固定并行度
经过这次事故,我在 Hudi 表的 Flink WITH 参数里加了这么一组配置:
WITH ( 'connector' = 'hudi', 'table.type' = 'MERGE_ON_READ', 'hoodie.index.type' = 'BUCKET', 'hoodie.bucket.index.num.buckets' = '4096', 'hoodie.bucket.index.hash.field' = 'order_id', 'write.operation' = 'upsert', 'compaction.async.enabled' = 'false', 'clustering.async.enabled' = 'false', 'hoodie.clean.async' = 'false', 'hoodie.write.lock.provider' = 'org.apache.hudi.client.transaction.lock.FileSystemBasedLockProvider' )重点解释下每个参数:
hoodie.bucket.index.num.buckets:我把它从 1024 调到 4096,从概率上分散热门 bucket 的压力。但注意,bucket 数量不是一个可以随时改的参数。表一旦写入,bucket 映射关系就已经固化了,临时改大会导致新旧索引不一致,必须重刷全量数据。所以最好在建表前按峰值数据量规划好。compaction.async.enabled和clustering.async.enabled:先关闭自动执行,用定时任务每天低峰跑一次。线上表如果 log 文件增长很快,可以选择开 schedule,但执行动作必须单独拉起。hoodie.write.lock.provider:让多客户端共享文件锁,至少能在更早阶段暴露冲突,而不是等到 commit 才爆炸。
5.3 版本升级带来的变化与坑
我们的 Flink 版本是 1.14,Hudi 最初用的是 0.12.3。排障时翻了社区 issue,发现 0.12 到 0.13 之间对 Flink Sink 的writeToken生成逻辑做过修复,专门针对恢复状态时可能出现重复 token 的情况。后来我们升级到 0.13.1,再往后到 0.14.1,同类异常再也没有出现在主链路。
版本升级不是银弹,但如果你还在比较旧的版本上反复遇挫,先看一眼 Hudi 的 release notes,搜索DuplicateFileIdException或者writeToken关键字。有些修复确实只在特定版本里生效。
升级前要注意的地方:Hudi connector 要和 Flink 的 shaded-hadoop 版本匹配,最好在测试环境完整跑一轮写入、Compaction、查询再上线。不要把生产表直接拿去做版本升级实验。
5.4 可观测性:让“下一次”提前暴露出来
这个异常最怕的不是难修,而是出现频率低、恢复一次就以为没事。所以要给 Hudi 表加监控,重点看两个指标:
- 最近一小时失败的 commit 数量;
.hoodie目录下的 inflight instant 数量。
Flink 作业本身的 Checkpoint 监控也要分组看。如果某个 TaskManager 上的 Sink 子任务经常出现“Backpressure”或者“Rescale”相关日志,就要警惕状态在重新分配时把相同的 fileId 交到了多个 subtask。
另外,我习惯把每个 Hudi 表的参数变更提交到 git 仓库存档,出现问题时可以直接对照“这个表是什么时候调了并行度、什么时候改了 bucket 数量”。很多诡异报错,其实就是某个参数在某次变更中被悄悄改了。
6. 真·经验总结:如果再遇到,我会按什么顺序排查
我踩过一次,现在反而淡定了。以后再看到DuplicateFileIdException,排查顺序固定下来:
第一,看当前有几个作业在写同一个表路径;第二,看异常信息里create_file_id和delete_file_id是否相同;第三,看.hoodie下的 instant 是否有重叠窗口;第四,再看作业最近一次恢复是否调整过并行度。
最后再说一个容易被忽略的细节:尽量别用“删除 Hudi 表部分分区文件”的方式来清理冲突数据。文件系统层删了文件,但时间线上还留着那些 instant 记录,后面提交时反而可能引发其他未知异常。更稳妥的做法是保留 Hudi 的元数据完整性,把任务停住,确认问题来自哪个写入者,然后从源头关闭那个写入者,再用最近一个干净 Checkpoint 恢复。
这套组合拳打完以后,我们的实时链路已经稳定运行了半年多。如果你们也在被这个异常反复折腾,希望这篇复盘能帮你少走几个小时的弯路。