[SUPPORT] Silently corrupt Parquet files produced by Spark Insert Overwrite / Compaction / Clustering — writer succeeds, reader fails with dictionary / repetition-level index out of bounds
- Dominant language
- Java
- Stars
- 6.2k
- Forks
- 2.5k
- Avg merge
- 2d 8h
- Merged PRs (30d)
- 111
Description
### Problem
We have repeatedly hit a class of silently-corrupt Parquet files produced by Hudi's Spark write path in production. The write side reports success and no exception is thrown; the corruption only surfaces days later when downstream jobs try to read the affected file. This has already happened several times on our fleet, and each incident forces us to either roll back the offending instant or manually rewrite the whole partition, which is operationally expensive.
The corruption is low-probability — only a handful of files have been identified across a very large volume of Spark writes — but the impact of each incident is high, since the file passes every writer-side check and only fails once downstream reads hit it.
The failing write paths we have observed so far:
- INSERT OVERWRITE (Spark SQL)
- Compaction
- Clustering
### Symptoms
On the writer side:
1. The Spark job completes successfully. No task attempt fails, no exception is logged, and Hudi's commit / replacecommit metadata is written normally.
2. Markers, .hoodie timeline, and file listings are all consistent — the file passes every existing writer-side validation.
On the reader side (Spark / Hive / Flink / Parquet-cli readers using the standard Parquet reader):
1. ArrayIndexOutOfBoundsException inside the Parquet dictionary decoder (dictionary index out of range).
2. ArrayIndexOutOfBoundsException inside the repetition-level / definition-level decoder.
3. Parquet footer is intact and readable
4. Parquet page CRC checks pass
5. Partial pages within the same column chunk are readable
6. The exception type and location are not stable across incidents, for example:
```
Caused by: org.apache.parquet.io.ParquetDecodingException: Can not read value at 432297 in block 0 in file hdfs://ns22027/xx/dt=2025-06-01/2b4bae09-9918-4b51-85d1-d5a89f34616a-0_19297-1-867_20260524205029382.parquet
at org.apache.parquet.hadoop.InternalParquetRecordReader.nextKeyValue(InternalParquetRecordReader.java:264)
at org.apache.parquet.hadoop.ParquetRecordReader.nextKeyValue(ParquetRecordReader.java:207)
at org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:39)
at org.apache.spark.sql.execution.datasources.RecordReaderIterator$$anon$1.hasNext(RecordReaderIterator.scala:61)
at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:125)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:513)
at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage9.sort_addToSorter_0$(Unknown Source)
at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage9.processNext(Unknown Source)
at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:760)
at org.apache.spark.sql.execution.datasources.FileFormatWriter$.executeTask(FileFormatWriter.scala:393)
at org.apache.spark.sql.execution.datasources.WriteFilesExec.$anonfun$doExecuteWrite$1(WriteFiles.scala:100)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:897)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:897)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:337)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:92)
at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:161)
at org.apache.spark.scheduler.Task.run(Task.scala:140)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:562)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1555)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:565)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.ArrayIndexOutOfBoundsException: 53
at org.apache.spark.sql.execution.datasources.parquet.ParquetRowConverter$ParquetStringConverter.addValueFromDictionary(ParquetRowConverter.scala:530)
at org.apache.parquet.column.impl.ColumnReaderBase$1.writeValue(ColumnReaderBase.java:186)
at org.apache.parquet.column.impl.ColumnReaderBase.writeCurrentValueToConverter(ColumnReaderBase.java:440)
at org.apache.parquet.column.impl.ColumnReaderImpl.writeCurrentValueToConverter(ColumnReaderImpl.java:30)
at org.apache.parquet.io.RecordReaderImplementation.read(RecordReaderImplementation.java:406)
at org.apache.parquet.hadoop.InternalParquetRecordReader.nextKeyValue(InternalParquetRecordReader.java:234)
... 44 more
```
```
java.lang.RuntimeException: Failed on record 282883 in 6b1dd48d-d7d9-43dd-91e9-883669333eb4-1_23816-1-5364_20260625203938302.parquet
at org.apache.parquet.cli.commands.ScanCommand.run(ScanCommand.java:75)
at org.apache.parquet.cli.Main.run(Main.java:163)
at org.apache.hadoop.util.ToolRunner.run(ToolRunner.java:82)
at org.apache.parquet.cli.Main.main(Main.java:191)
Caused by: org.apache.parquet.io.ParquetDecodingException: Can not read value at 282885 in block 0 in file file: 6b1dd48d-d7d9-43dd-91e9-883669333eb4-1_23816-1-5364_20260625203938302.parquet
at org.apache.parquet.hadoop.InternalParquetRecordReader.nextKeyValue(InternalParquetRecordReader.java:280)
at org.apache.parquet.hadoop.ParquetReader.read(ParquetReader.java:136)
at org.apache.parquet.cli.BaseCommand$1$1.advance(BaseCommand.java:356)
at org.apache.parquet.cli.BaseCommand$1$1.next(BaseCommand.java:350)
at org.apache.parquet.cli.commands.ScanCommand.run(ScanCommand.java:70)
... 3 more
Caused by: java.lang.ArrayIndexOutOfBoundsException: Index 254 out of bounds for length 2
at org.apache.parquet.io.RecordReaderImplementation.read(RecordReaderImplementation.java:431)
at org.apache.parquet.hadoop.InternalParquetRecordReader.nextKeyValue(InternalParquetRecordReader.java:249)
... 7 more
```
### To Reproduce
The corruption is low-probability — only a handful of files have been identified across a very large volume of Spark writes
### Environment Description
* Hudi version: 0.13.1
* Spark version: 3.4.3
* Flink version:
* Hive version: 2.0.0
* Hadoop version:
* Storage (HDFS/S3/GCS..):
* Running on Docker? (yes/no)
Contributor guide
No contributing guide indexed for this repository
Research direction
Start with the named Spark write paths—INSERT OVERWRITE, compaction, and clustering—and compare them with the reader failures in InternalParquetRecordReader and RecordReaderImplementation. Done means reproducing the corruption or narrowing it to a deterministic write/read boundary; the issue names no Hudi source file or test to run.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- hadoop, java, spark
- Domain
- data-engineering, distributed-systems
- Issue type
- Bug
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Active
- Clarity
- Needs clarification
- Newbie friendliness
- 30/100