EMR-Spark生成随机损坏的Snappy压缩Parquet文件,读取报FAILED_TO_UNCOMPRESS
偶发Snappy压缩Parquet文件损坏问题排查及解决方案
问题描述
有多款作业除Spark依赖库外无任何公共代码,均会输出Snappy压缩的Parquet格式文件,近期随机出现文件损坏问题:写入作业运行无报错、无异常退出,绝大多数场景下输出文件正常,仅在后续读取文件时才会发现损坏情况。该问题偶发随机,难以稳定复现,写入侧日志中无任何和文件损坏相关的错误记录。
该类问题自11月5日起首次出现,此前全年运行稳定,问题出现后未重新部署过EMR集群,也未修改过作业代码。
读取损坏文件报错汇总
- Parquet Tools读取报错:
OSError: IOError: Corrupt snappy compressed data. - Spark读取报错:
Caused by: java.io.IOException: could not read page Page [bytes.size=1048612, valueCount=58058, uncompressedSize=1048612] in col [client_id] optional binary client_id (UTF8) at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readPageV1(VectorizedColumnReader.java:618) at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.access$000(VectorizedColumnReader.java:49) at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader$1.visit(VectorizedColumnReader.java:547) ... 44 more Caused by: java.io.IOException: FAILED_TO_UNCOMPRESS(5)
- Athena读取报错:
GENERIC_INTERNAL_ERROR: Malformed input: offset=259867
已知技术栈版本
- emr 5.31.0
- Spark 2.4.6
- openjdk 1.8.0_302
- EMRFS consistent view: Disabled
排查思路与解决方案
1. 优先排查底层存储静默数据损坏
因为EMRFS一致性视图关闭,且无代码/集群变更,首先检查底层对象存储是否在11月5日前后有静默写入异常:
- 对比损坏文件的写入时间点是否集中在特定可用区/存储节点故障时段
- 校验损坏文件的etag与写入作业日志中记录的输出文件md5是否匹配,确认是否为存储端写入时的数据截断/篡改
2. 排查Snappy压缩相关的隐性Bug
Spark 2.4.6配套的Snappy-java版本存在已知偶发压缩数据损坏问题,在特定JDK内存对齐场景下会触发:
- 临时验证方案:将作业压缩格式切换为Gzip,运行3-7天观察是否还有损坏文件出现,排除压缩库问题
- 永久修复方案:若确认为Snappy库问题,可升级snappy-java版本到1.1.8.4及以上,在Spark作业启动参数添加
--conf spark.driver.extraClassPath=./snappy-java-1.1.8.4.jar --conf spark.executor.extraClassPath=./snappy-java-1.1.8.4.jar,优先加载高版本Snappy库覆盖EMR内置版本
3. 排查Parquet写入时的内存溢出隐性问题
偶发的Executor堆外内存溢出但未触发进程退出的场景,会导致Parquet页写入时内存数据错乱:
- 检查作业在11月5日之后的数据量变化,是否有单条数据大小突增、并行度不足导致Executor内存占用升高的情况
- 增加Executor堆外内存配置:
--conf spark.executor.memoryOverhead=2048,观察问题是否复现
4. 临时规避方案
在写入作业后增加校验步骤:作业写完所有Parquet文件后,新增一个小任务遍历所有输出文件,调用spark.read.parquet加载全量文件并触发count操作,主动识别损坏文件,触发重跑逻辑
内容的提问来源于stack exchange,提问作者tdebroc
相关产品推荐
相关产品推荐

