Occasional corruption of parquet files , parquet writer might not be calling ParquetFileWriter->end()
- 主要语言
- Java
- 星标
- 3.1k
- 派生
- 1.6k
- 平均合并
- 3 天 12 小时
- 30 天内合并 PR
- 33
描述
We have a high volume streaming service which works most of the time . But off late we have been observing that some of the parquet files written out by write flow are getting corrupted. This is manifested in our reading flow with the following exception
Writer version - 1.6.0 , Reader version - 1.7.0
Caused by: java.lang.RuntimeException: hdfs://Ingest/ingest/jobs/2017-11-30/00-05/part4139 is not a Parquet file. expected magic number at tail [80, 65, 82, 49] but found [-28, -126, 1, 1]
at org.apache.parquet.hadoop.ParquetFileReader.readFooter(ParquetFileReader.java:422)
at org.apache.parquet.hadoop.ParquetFileReader.readFooter(ParquetFileReader.java:385)
at org.apache.parquet.hadoop.ParquetRecordReader.initializeInternalReader(ParquetRecordReader.java:157)
at org.apache.parquet.hadoop.ParquetRecordReader.initialize(ParquetRecordReader.java:140)
at org.apache.spark.rdd.SqlNewHadoopRDD$$anon$1.(SqlNewHadoopRDD.scala:180)
at org.apache.spark.rdd.SqlNewHadoopRDD.compute(SqlNewHadoopRDD.scala:126)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:306)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:270)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:306)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:270)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:306)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:270)
at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:73)
at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41)
at org.apache.spark.scheduler.Task.run(Task.scala:89)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:227)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)```
After looking at the code , i can see that one of the possible causes is/are
1] footer not being serialized in the writer due to end not being called
but we are not seeing any exceptions on the writer.
2] data size - does data size has impact ? There will be cases when row group sizes will be huge as it is activity data of a user
We are using default parquet block size and hdfs block size . Other than upgrading to the latest version and re-test , what are the options we have to debug a issue like this
**Reporter**: [venkata yerubandi](https://issues.apache.org/jira/secure/ViewProfile.jspa?name=raoyvn)
**Note**: *This issue was originally created as [PARQUET-1176](https://issues.apache.org/jira/browse/PARQUET-1176). Please see the [migration documentation](https://issues.apache.org/jira/browse/PARQUET-2502) for further details.*
贡献指南
这个仓库没有索引到贡献指南
调研方向
首先检查 ParquetFileReader.readFooter 的堆栈跟踪位置,尤其是 ParquetFileReader.java:422,并检查写入路径是否调用了 ParquetFileWriter->end()。在高容量流式服务中复现损坏问题,同时调整 row-group 和 HDFS block 的大小。完成标准是确定缺失尾部 magic number 的原因,并记录经过验证的调试或修正路径。
由索引模型根据 Issue 内容生成。
评估
- 技术栈
- java, spark
- 领域
- data-engineering
- Issue 类型
- 缺陷
- 难度
- 4/5
- 预计耗时
- 3-5 天
- 活跃度
- 停滞
- 描述清晰度
- 需要澄清
- 新手友好度
- 25/100