凌晨两点多,手机连续震了好几下,我就知道生产又出事了。打开钉钉群里的监控截图,一条Flink批作业失败,再往下翻,找到真正的异常堆栈:
Caused by: java.io.IOException: hdfs://nameservice1/user/hive/warehouse/ods_mall.db/user_actions/dt=2024-06-17/8712bc29-df3a-4c7a-9a55-xxxx-xxxx_20240617000123.parquet is not a Parquet file. expected magic number at tail, but found 0这个报错太典型了。做Flink-Hudi生产维护的同学,十有八九都会在某个深夜跟它碰面。表面意思是:那个后缀是.parquet的文件,并不是一个合法的Parquet文件。但真正的问题远没有这么简单——排障绕了一晚上加一个上午,最后定位到根因的时候,发现报错信息只是冰山一角。这篇文章就把完整的排查过程和底层原理写透,给正在跟Flink、Hudi、Parquet打交道的人一个可复用的排障思路。
1. 报错并不会告诉你真正的病因:先拆解“is not a Parquet file”在指什么
很多人第一次看到这种报错,第一反应是“文件损坏了”,第二反应是“是不是有人往表目录里丢了非Parquet格式的文件”。这两个猜测都有道理,但都不够准确。
1.1 报错究竟是谁抛出来的
先别急着折腾数据,先看清楚这个异常是从哪一层抛出来的。完整堆栈长这样:
Caused by: java.io.IOException: .../xxx.parquet is not a Parquet file. expected magic number at tail, but found 0 at org.apache.parquet.hadoop.ParquetFileReader.readFooter(ParquetFileReader.java:452) at org.apache.parquet.hadoop.ParquetFileReader.readFooter(ParquetFileReader.java:388) at org.apache.hudi.hadoop.HoodieParquetFileFormat.getReader(HoodieParquetFileFormat.java:...) ...看到ParquetFileReader.readFooter就有意思了。Parquet的读取流程是从文件尾部开始的,读取器要先去定位footer(文件尾部的元数据区),footer里保存了schema、row group列表、统计信息这些关键内容。如果footer读不出来或者末尾的魔术数字不对,就会抛出“expected magic number at tail”这类异常。
1.2 Parquet文件的“身份证”是什么
一个合法的Parquet文件,头部和尾部都有一段特殊的标记,叫magic number,就是四个字节的PAR1(十六进制是50 41 52 31)。文件结构大致是:
PAR1 -- 文件头魔数 [列数据块 / row groups] [footer元数据] -- 包裹的Thrift结构体 [footer长度] -- 4字节 PAR1 -- 文件尾魔数读取的时候,Parquet工具会先跳到文件末尾,检查最后四个字节是不是PAR1;是的话再往前读footer长度,然后用Thrift反序列化footer内容。所以“尾部魔数校验失败”这句话,翻译成人话就是:文件在写入过程中没有正常收尾。
1.3 文件扩展名只是标签,文件状态才是关键
这里有个反直觉的点:报错说“is not a Parquet file”,并不代表这个文件本来就不是Parquet,绝大多数情况下是它还没来得及成为一个完整Parquet文件。就像一个Word文档写到一半被强杀了,文件后缀还是.docx,但用Office打开就会提示“文件格式损坏”或者“内容不可读”。
生产实践中这种“半成品”的常见成因有三类:
- 写入进程异常中断:磁盘满了、内存溢出、task manager被杀、网络抖动导致rename失败,文件只写了一半就停了。
- 作业被强制终止:有人直接kill了Flink任务,或者运维重启集群,Hudi的rollback流程没有机会跑完,残留在文件系统里的文件就成了孤儿。
- 并发写同一个分区:两个任务同时操作同一张Hudi表,某个文件被另一个任务清理或覆盖,读到的时间窗口里文件正处于“薛定谔状态”。
所以看到这个报错,第一反应应该是:别急着删文件,先确认这个文件在Hudi元数据里是什么身份。是已经提交的正式文件,还是没人认领的孤儿文件?这一步判断错了,后面所有操作都可能把问题放大。
2. 现场取证:从一行报错到HDFS文件状态与Hudi时间线体检
我建议大家排障时养成一个习惯:把报错文本当作线索,而不是答案。真正要查的是文件系统上的实体文件和Hudi时间线之间的关系。
2.1 第一手信息:先拿全文件路径和堆栈类名
我当时的做法是先把这条异常的完整堆栈和上下文日志保存下来。除了报错这一行,还要看:
- 报错文件的完整路径,包括分区目录名和文件名。
- 是哪个组件在读取这个文件。比如通过Hive查询触发的,还是Spark读Hudi表触发的,或者Flink任务自己启动时做的checkpoint恢复。读取方不同,排查重心完全不同。
- 堆栈里出现了哪些类。如果是
HoodieParquetFileFormat,基本确定是Spark SQL读Hudi表;如果直接是ParquetFileReader,可能是Hive外部表或某些工具直接列了目录。
这一步至少要花五分钟仔细看,很多人跳过直接搜报错文本,很容易误判方向。
2.2 给文件做“体检”:大小、头部和尾部
拿到文件路径后,最快的排查手段是看文件的基本信息:
# 先看文件大小,0字节或者明显小于同分区其他文件,立刻引起警惕 hdfs dfs -ls -h /user/hive/warehouse/ods_mall.db/user_actions/dt=2024-06-17/ # 把可疑文件拷到本地再检查 hdfs dfs -get /user/hive/warehouse/ods_mall.db/user_actions/dt=2024-06-17/8712bc29-df3a-4c7a-9a55-xxxx-xxxx_20240617000123.parquet /tmp/suspect.parquet # 本地看出文件类型 file /tmp/suspect.parquet # 看文件尾部,正常结尾应该能看到 PAR1 xxd /tmp/suspect.parquet | tail -5 # 用parquet-tools做元数据解析 parquet-tools meta /tmp/suspect.parquet我当时拿到的结果非常明确:文件只有12KB,而同一分区正常文件都是200MB以上;xxd查看尾部,完全看不到PAR1魔数,最后几十个字节全是0x00。parquet-tools meta直接报MalformedThriftException,说footer解析不了。
补一个小技巧,如果想批量检查某个分区目录下有没有可疑Parquet文件,可以用一个简单的Python脚本辅助:
from pathlib import Path base = Path('/tmp/parquet_check') for p in base.glob('*.parquet'): b = p.read_bytes() ok = len(b) >= 12 and b[:4] == b'PAR1' and b[-4:] == b'PAR1' if not ok: print(f'[SUSPECT] {p.name}, size={len(b)}, tail={b[-4:]!r}')这个脚本也适合事后加进巡检流程当第一道粗筛。
2.3 用Hudi时间线给文件“验明正身”
文件格式检查只是确认了“文件确实不完整”,但更关键的问题是:为什么一个不完整的文件会出现在这张表的分区目录里?它是不是Hudi正式提交的一部分?
Hudi表和普通目录表最大的区别就是它有一套时间线机制(timeline)。每次写入、清理、压缩都会在表的.hoodie/timeline目录下生成一个instant记录。可以去Hudi表的.hoodie目录下看一眼:
hdfs dfs -ls -h /user/hive/warehouse/ods_mall.db/user_actions/.hoodie/timeline/你会看到类似这样的文件:
20240617000123.commit.requested 20240617000123.commit.inflight 20240617000123.commit 20240617000124.clean关键是要判断文件名里带的instant时间对不对。比如可疑文件名是:
8712bc29-df3a-4c7a-9a55-xxxx-xxxx_20240617000123.parquet文件名里的20240617000123理应是它对应的commit时间。然后去时间线里查一下20240617000123这个instant是否存在、状态是不是completed。
2.4 把同一个File Group里的兄弟文件也翻出来
Hudi的组织方式是:一个分区下有多个file group,每个file group可以有base file(.parquet)和若干个log file(.log后缀)。同一个file group的文件,在文件名前半段拥有相同的fileId。
我用可疑文件名里的fileId:
8712bc29-df3a-4c7a-9a55-xxxx-xxxx去分区目录里搜了一下:
hdfs dfs -ls -h /user/hive/warehouse/ods_mall.db/user_actions/dt=2024-06-17/ | grep "8712bc29"结果发现这个file group下同时存在一个正常的parquet文件和一个log文件,而那个可疑的12KB文件用的fileId也是同一个。这就说明:这个12KB文件是同一个写入批次里被分裂出来的另一个文件块,正常文件写完了,它没写完。
到这一步,排障链路已经基本清晰了:Flink的写入任务在某次运行中创建了这个文件,但写入过程异常中断,Hudi的元数据没有正确提交它,文件就残留在了HDFS上。
3. 根因落地:Flink写入中断留下的半成品文件是怎么混进读取链路的
很多人到这里会有一个很大的疑问:Hudi不是有事务机制吗?为什么一个没提交的文件会进入读取链路?这就要聊一聊Flink写入Hudi时的底层工作方式了。
3.1 Hudi写入端的工作流:从flush文件到提交Commit
Flink写Hudi的背后是HoodieFlinkWriteClient。每次写入批次大致会经过这几个阶段:
- Flink算子把数据写入内存缓冲区。
- 缓冲区满了或者到了checkpoint时机,写入端会创建一个新的base file,并持续向HDFS写出Parquet数据。
- 文件写出完成后,会在Hudi时间线上生成一个commit动作,把文件路径和元数据写进commit文件。
- commit成功后,这个文件才被认为是这张表正式数据的一部分。
问题就出在第2步和第3步之间。
如果在第2步和第3步之间,Flink作业失败或者被强制终止,HDFS上已经落地的Parquet文件就处于“孤儿”状态。正常情况下Hudi的rollback机制会清理掉这种文件,但rollback的触发是有条件的:作业能够收到失败信号并走完rollback流程。如果作业进程直接被kill、task manager崩溃,或者磁盘写入本身处于半挂起状态,rollback根本没有机会执行。
3.2 为什么读取端几乎必然踩雷
读取端的情况更有意思。以Hive查询Hudi表为例,Hive Metastore里登记的是表目录和分区路径,查询引擎(Spark或Tez)列出的其实是分区目录下的文件集合。Hudi自己的InputFormat会尝试把这些文件归属到file group里,但如果文件系统里正好残留了一个名字符合Parquet模式、却没有被commit记录引用的物理文件,某些读取路径就会直接把它当作Parquet文件去解析。
结果就是:元数据层面它不存在,物理层面它又真实躺在那里。读到它的查询任务自然就炸了。
这也是为什么这个问题在生产中会反复出现:Flink写入任务可能通过checkpoint恢复了,后续批次继续写新的数据,读取任务却被同一个残留文件反复卡死。我那次就是Flink任务重启后跑通了,但下游Spark任务一读Hudi表就失败,因为那个坏文件一直没被清理。
3.3 “不是Parquet”只是表象,元数据和物理文件不一致才是根
把这次事故的本质抽出来,其实是三个层面的问题:
- 物理文件层面:HDFS上存在一个不完整的Parquet文件。
- Hudi元数据层面:时间线里没有一个对应的有效commit记录。
- 读取逻辑层面:某些读引擎并不会仔细检查“这个文件是否在commit元数据里”,而是直接尝试解析它。
明白了这三层关系,你就不会再被“not a Parquet file”这个表象牵着走了。它只是一个很表面的症状,真正的病灶是写中断之后,文件系统和时间线之间出现了不一致。
4. 止血操作记录:恢复服务、清理残留、完成数据核验
排障到这里,剩下的就是动手解决了。我整理一下当时的完整操作,给大家一条尽量安全的路径。
4.1 先恢复服务:让读取任务绕开坏文件
止血永远比追责优先。我当时做的第一件事,是跟数据研发确认这个坏文件对应的分区数据是否可以重新生成。可以的话,最干脆的方案是把这个坏文件移动到备份目录,而不是直接删除——毕竟生产环境里“删了就后悔”的情况太多了。
hdfs dfs -mv /user/hive/warehouse/ods_mall.db/user_actions/dt=2024-06-17/8712bc29-df3a-4c7a-9a55-xxxx-xxxx_20240617000123.parquet /tmp/recovery_bucket/20240617_suspect.parquet移动完以后,让下游Spark任务重新读一次Hudi表。如果能正常通过count和抽样查询,说明阻塞读取的就是这个文件。这时候不要急着宣告恢复,还要继续往下验证数据完整性。
4.2 更规范的清理姿势:借助Hudi的元数据能力
如果你们集群上有Hudi CLI或者Spark环境,更好的做法是用Hudi自身的能力去清理,而不是手动mv文件。用Hudi CLI连接表之后,可以查看时间线状态:
commits show看看有没有处于inflight或requested状态的instant。如果确认某个instant是失败的,并且它的文件都是半成品,可以通过Hudi的元数据操作把它的残留文件纳入rollback处理。这样表的时间线和物理文件都能保持一致,后续clean、compaction跑起来也更干净。
手动mv文件是应急手段,能解决当前阻塞,但不会修正时间线。所以有条件的话,还是要把时间线状态同步处理好。
4.3 数据完整性核验:不是文件能读就大功告成
文件能读了只是第一步,还要确认这张表的数据没有丢、没有重复。我当时的核验分三步:
- 行数对比:用Spark读Hudi表,按分区count;
- 时间范围对比:查该分区数据的业务时间字段的min和max,和上游数据源核对;
- 抽样内容对比:抽查几个关键维度的去重数量,比如用户数、订单号,和上游平台的统计值做对齐。
只有这三步都通过了,才敢跟业务说数据可用。
4.4 复盘一下这次操作代价最小的路线
整个操作顺序应该是:
备份坏文件 -> 读取任务重试 -> 行数/内容核验 -> 确认无碍后清理时间线残留 -> 通知业务方恢复消费如果反过来先删文件再验证,一旦删错,就得重新从上游重刷整个分区,代价会大很多。生产排障的原则永远是先备份、再隔离、后清理。
5. 防复发设计:写入端参数、读取端校验与异常监控
解决了一次事故,如果不做任何配置调整,大概率还会在另一个凌晨再次遇到。这里写一下我在事发后针对写入链路和读取链路做的防御性改动。
5.1 写入端:把检查点和Hudi行为参数化
Flink + Hudi的写入,本质上是“Flink管理状态、Hudi管理存储”。两者之间需要一组合理的参数来降低“半成品文件”出现的概率。以下是基于当时生产环境做的一组配置,不同Hudi版本参数名可能有差异,核心是理解每个参数在干什么:
| 参数 | 作用 | 实践经验 |
|---|---|---|
execution.checkpointing.interval | Flink做Checkpoint的频率 | 设太短会导致文件数量膨胀,设太长会导致恢复成本高。我们最终调整为3到5分钟 |
write.task.max.size | 控制写入端算子数据攒批阈值 | 这决定了单个Parquet文件在flush前累积多少数据,间接影响小文件数量 |
hoodie.parquet.small.file.limit | 小于该阈值的文件会被视为小文件并触发合并 | 调大小文件阈值前,要评估资源占用,别盲目放大 |
hoodie.compact.inline | 是否开启内联压缩 | 开启后数据文件更整齐,log文件不会无限膨胀 |
hoodie.clean.automatic | 是否自动清理过期文件 | 保持默认开启,但要注意clean不清理“未提交的孤儿文件” |
注意:不同Hudi版本(0.12、0.13、0.14、1.x)的参数命名有差异。我配置时用的是当时线上Hudi 0.12.x的写法,升级版本前要去对应版本的官方配置页核对一遍,别直接照抄。
5.2 读取端:加一道轻量级文件预检
读取端可以在任务调度前对Hudi分区目录做一次快速体检,比如用前面那个Python脚本扫描新分区目录下的文件头尾魔数,一旦发现异常文件立刻报警并终止下游任务启动。这个预检不用全表扫,只扫最近新增的分区,成本很低,但能避免读任务被一个坏文件拖死。
5.3 监控告警:把“孤儿文件”当成独立指标
我后来在监控系统里加了一个指标:Hudi表分区目录下的Parquet文件数量,和Hudi时间线里commit文件引用的Parquet文件数量,是否一致。这两个数字一旦对不上,就说明有孤儿文件或者“幽灵文件”。用脚本跑对比,每半个小时一次,报警阈值设为超过0就告警。
这个指标救过我第二次——有一次一个上游任务被kill,就是靠这个指标在业务方发现之前拦截下来的。
5.4 顺手排掉两个容易搞混的报错
排查过程中,也会看到很多同学把Hudi的Parquet异常跟Flink的JDBC连接器异常弄混。比如网络上常有人问“Flink的JDBC连接器报错怎么办”,那个通常是你MySQL或ClickHouse的连接稳定性问题,跟Hudi的Parquet文件无关。
还有一类是“Flink写Hudi时目录权限不足”的报错,报错里会出现Permission denied,那是集群权限策略的问题。所有这些报错的共性规律是一样的:先确认出错的组件是谁,再去看它正在访问的资源状态。别看到一个IOException就朝着“文件损坏”的方向死查,那很可能是另一个层面的问题。
6. “not a Parquet file”家族:五种相似报错的一次性识别
回头想你可能会发现:这个报错虽然文本相似,但背后的成因和处置方式完全不同。为了让大家下次不用再一步步试,我把这一类报错按特征拆成了五种变体。
| 报错特征 | 常见现场 | 判定手段 | 处理方式 |
|---|---|---|---|
expected magic number at tail, but found 0 | 文件写入中断,footer没有落盘 | 文件很小+尾部无PAR1 | 备份坏文件,重刷分区 |
empty file expected magic number at tail | 0字节空文件残留 | ls -l看到大小0 | 直接删除或隔离 |
file is not a parquet file, expected metadata at footer | footer长度字段异常 | parquet-tools meta失败 | 从上游数据重刷 |
expected magic number at head, but found xxx | 文件头魔数错误,可能是非Parquet文件被改了后缀名 | 头部前4字节非PAR1 | 检查文件来源 |
malformed thrift structure | footer存在但Thrift解析失败 | 结构损坏或版本异常 | 用parquet-tools检查,重刷数据 |
判断这类问题,有一个百试百灵的顺序:看大小 -> 看头尾魔数 -> 看footer -> 对照Hudi时间线。按这个顺序走,基本五分钟内能判断出坏文件的严重程度和来源。
还要提醒一点:有时候报错不是IOException而是Schema相关的AvroTypeException,字面上没有“is not a Parquet file”,但读的同样是Hudi表。那是Schema演进或文件版本混杂的问题,处理方式完全不同——不要拿本文的“删文件”思路去处理Schema不匹配,否则真的会删出一条错误路线。
回到这次事故本身,我个人最大的收获不是记住了某条命令,而是建立了“物理文件与元数据对照”的排障视角。从那以后,遇到Hudi的诡异报错,我第一件事永远是拉文件系统快照和时间线对比检查。这个习惯,建议每一个维护Flink-Hudi环境的朋友都养成。下次如果你也在凌晨被这种报错叫醒,先别慌,备份好坏文件,然后从文件尾部的魔数开始查起,大概率能少走一大段弯路。