ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink写Hudi遇半成品Parquet文件:从魔数报错到时间线排障与防复发

Flink写Hudi遇半成品Parquet文件:从魔数报错到时间线排障与防复发 凌晨两点多手机连续震了好几下我就知道生产又出事了。打开钉钉群里的监控截图一条Flink批作业失败再往下翻找到真正的异常堆栈Caused by: java.io.IOException: hdfs://nameservice1/user/hive/warehouse/ods_mall.db/user_actions/dt2024-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/dt2024-06-17/ # 把可疑文件拷到本地再检查 hdfs dfs -get /user/hive/warehouse/ods_mall.db/user_actions/dt2024-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] bPAR1 and b[-4:] bPAR1 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。我用可疑文件名里的fileId8712bc29-df3a-4c7a-9a55-xxxx-xxxx去分区目录里搜了一下hdfs dfs -ls -h /user/hive/warehouse/ods_mall.db/user_actions/dt2024-06-17/ | grep 8712bc29结果发现这个file group下同时存在一个正常的parquet文件和一个log文件而那个可疑的12KB文件用的fileId也是同一个。这就说明这个12KB文件是同一个写入批次里被分裂出来的另一个文件块正常文件写完了它没写完。到这一步排障链路已经基本清晰了Flink的写入任务在某次运行中创建了这个文件但写入过程异常中断Hudi的元数据没有正确提交它文件就残留在了HDFS上。3. 根因落地Flink写入中断留下的半成品文件是怎么混进读取链路的很多人到这里会有一个很大的疑问Hudi不是有事务机制吗为什么一个没提交的文件会进入读取链路这就要聊一聊Flink写入Hudi时的底层工作方式了。3.1 Hudi写入端的工作流从flush文件到提交CommitFlink写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/dt2024-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.intervalFlink做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 tail0字节空文件残留ls -l看到大小0直接删除或隔离file is not a parquet file, expected metadata at footerfooter长度字段异常parquet-tools meta失败从上游数据重刷expected magic number at head, but found xxx文件头魔数错误可能是非Parquet文件被改了后缀名头部前4字节非PAR1检查文件来源malformed thrift structurefooter存在但Thrift解析失败结构损坏或版本异常用parquet-tools检查重刷数据判断这类问题有一个百试百灵的顺序看大小 - 看头尾魔数 - 看footer - 对照Hudi时间线。按这个顺序走基本五分钟内能判断出坏文件的严重程度和来源。还要提醒一点有时候报错不是IOException而是Schema相关的AvroTypeException字面上没有“is not a Parquet file”但读的同样是Hudi表。那是Schema演进或文件版本混杂的问题处理方式完全不同——不要拿本文的“删文件”思路去处理Schema不匹配否则真的会删出一条错误路线。回到这次事故本身我个人最大的收获不是记住了某条命令而是建立了“物理文件与元数据对照”的排障视角。从那以后遇到Hudi的诡异报错我第一件事永远是拉文件系统快照和时间线对比检查。这个习惯建议每一个维护Flink-Hudi环境的朋友都养成。下次如果你也在凌晨被这种报错叫醒先别慌备份好坏文件然后从文件尾部的魔数开始查起大概率能少走一大段弯路。
返回列表