[Bug] Failed to finish checkpoint due to the expection: 'Failed to read 4 bytes'
- Dominant language
- Java
- Stars
- 3.4k
- Forks
- 1.4k
- Avg merge
- 1d 11h
- Merged PRs (30d)
- 396
Description
### Search before asking
- [x] I searched in the [issues](https://github.com/apache/paimon/issues) and found nothing similar.
### Paimon version
Paimon Version: 1.0.1
Storage: OSS
### Compute Engine
Flink 1.20.1
### Minimal reproduce step
Some of my flinksql tasks throw out 'Failed to read 4 bytes' during doing checkpoint after running for several hours.The task is to read from some paimon tables and write to a paimon table of partial-update engine.
The source code of the task:
`
SET 'parallelism.default' = '2';
SET 'table.exec.sink.not-null-enforcer'='DROP';
CREATE TABLE if not exists masked_schema.masked_table (
masked_field_1 STRING,
masked_field_2 STRING,
masked_field_3 STRING,
masked_field_4 STRING,
masked_field_5 STRING,
masked_field_6 STRING,
masked_field_7 STRING,
masked_field_8 STRING,
masked_field_9 TIMESTAMP,
masked_field_10 ARRAY>,
PRIMARY KEY (masked_field_1, masked_field_2) NOT ENFORCED
) WITH (
'bucket' = '10',
'merge-engine' = 'partial-update',
'changelog-producer' = 'full-compaction',
'full-compaction.delta-commits'='3',
'snapshot.time-retained' = '24h',
'fields.masked_field_10.aggregate-function' = 'nested_update',
'fields.masked_field_10.nested-key' = 'masked_subfield_1,masked_subfield_2,masked_subfield_3',
'fields.masked_field_9.sequence-group' = 'masked_field_10'
);
INSERT INTO masked_schema.masked_table
SELECT
masked_field_1,
masked_field_2,
masked_field_3,
masked_field_4,
masked_field_5,
masked_field_6,
masked_field_7,
masked_field_8,
CAST(NULL AS TIMESTAMP) AS masked_field_9,
CAST(NULL AS ARRAY>) AS masked_field_10
FROM
masked_schema.masked_source_table_1 /*+ OPTIONS('scan.infer-parallelism' = 'false', 'consumer-id' = 'masked_table', 'consumer.expiration-time' = '2d', 'consumer.mode' = 'exactly-once') */
WHERE
masked_flag ='0'
UNION ALL
SELECT
masked_field_1,
masked_field_2,
CAST(NULL AS STRING) AS masked_field_3,
CAST(NULL AS STRING) AS masked_field_4,
CAST(NULL AS STRING) AS masked_field_5,
CAST(NULL AS STRING) AS masked_field_6,
CAST(NULL AS STRING) AS masked_field_7,
CAST(NULL AS STRING) AS masked_field_8,
IFNULL(IFNULL(masked_update_time, masked_create_time), NOW()) AS masked_field_9,
ARRAY[ROW(masked_subfield_1,masked_subfield_2,masked_subfield_3)] AS masked_field_10
FROM
masked_schema.masked_source_table_2 /*+ OPTIONS('scan.infer-parallelism' = 'false', 'consumer-id' = 'masked_table', 'consumer.expiration-time' = '2d', 'consumer.mode' = 'exactly-once') */
WHERE
`masked_subfield_1` NOT IN ('-','','0') AND `masked_field_1` NOT IN ('-','','0');
`
### What doesn't meet your expectations?
`
2025-04-09 17:32:53
java.io.IOException: Could not perform checkpoint 323 for operator Writer : tmp_table_name (2/2)#156.
at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1394)
at org.apache.flink.streaming.runtime.io.checkpointing.CheckpointBarrierHandler.notifyCheckpoint(CheckpointBarrierHandler.java:147)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler.triggerCheckpoint(SingleCheckpointBarrierHandler.java:287)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler.access$100(SingleCheckpointBarrierHandler.java:64)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler$ControllerImpl.triggerGlobalCheckpoint(SingleCheckpointBarrierHandler.java:488)
at org.apache.flink.streaming.runtime.io.checkpointing.AbstractAlignedBarrierHandlerState.triggerGlobalCheckpoint(AbstractAlignedBarrierHandlerState.java:74)
at org.apache.flink.streaming.runtime.io.checkpointing.AbstractAlignedBarrierHandlerState.barrierReceived(AbstractAlignedBarrierHandlerState.java:66)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler.lambda$processBarrier$2(SingleCheckpointBarrierHandler.java:234)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler.markCheckpointAlignedAndTransformState(SingleCheckpointBarrierHandler.java:262)
at org.apache.flink.streaming.runtime.io.checkpointing.SingleCheckpointBarrierHandler.processBarrier(SingleCheckpointBarrierHandler.java:231)
at org.apache.flink.streaming.runtime.io.checkpointing.CheckpointedInputGate.handleEvent(CheckpointedInputGate.java:181)
at org.apache.flink.streaming.runtime.io.checkpointing.CheckpointedInputGate.pollNext(CheckpointedInputGate.java:159)
at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:122)
at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65)
at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:638)
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231)
at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:973)
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:917)
at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:970)
at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:949)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:763)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:575)
at java.lang.Thread.run(Thread.java:750)
Caused by: java.io.IOException: java.util.concurrent.ExecutionException: java.lang.RuntimeException: org.apache.paimon.shade.org.apache.parquet.io.ParquetDecodingException: Failed to read 4 bytes
at org.apache.paimon.flink.sink.StoreSinkWriteImpl.prepareCommit(StoreSinkWriteImpl.java:234)
at org.apache.paimon.flink.sink.GlobalFullCompactionSinkWrite.prepareCommit(GlobalFullCompactionSinkWrite.java:184)
at org.apache.paimon.flink.sink.TableWriteOperator.prepareCommit(TableWriteOperator.java:127)
at org.apache.paimon.flink.sink.RowDataStoreWriteOperator.prepareCommit(RowDataStoreWriteOperator.java:198)
at org.apache.paimon.flink.sink.PrepareCommitOperator.emitCommittables(PrepareCommitOperator.java:104)
at org.apache.paimon.flink.sink.PrepareCommitOperator.prepareSnapshotPreBarrier(PrepareCommitOperator.java:84)
at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.prepareSnapshotPreBarrier(RegularOperatorChain.java:89)
at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:332)
at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$18(StreamTask.java:1437)
at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50)
at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:1425)
at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1382)
... 22 more
Caused by: java.util.concurrent.ExecutionException: java.lang.RuntimeException: org.apache.paimon.shade.org.apache.parquet.io.ParquetDecodingException: Failed to read 4 bytes
at java.util.concurrent.FutureTask.report(FutureTask.java:122)
at java.util.concurrent.FutureTask.get(FutureTask.java:192)
at org.apache.paimon.compact.CompactFutureManager.obtainCompactResult(CompactFutureManager.java:67)
at org.apache.paimon.compact.CompactFutureManager.innerGetCompactionResult(CompactFutureManager.java:53)
at org.apache.paimon.mergetree.compact.MergeTreeCompactManager.getCompactionResult(MergeTreeCompactManager.java:223)
at org.apache.paimon.mergetree.MergeTreeWriter.trySyncLatestCompaction(MergeTreeWriter.java:328)
at org.apache.paimon.mergetree.MergeTreeWriter.prepareCommit(MergeTreeWriter.java:276)
at org.apache.paimon.operation.AbstractFileStoreWrite.prepareCommit(AbstractFileStoreWrite.java:218)
at org.apache.paimon.operation.MemoryFileStoreWrite.prepareCommit(MemoryFileStoreWrite.java:155)
at org.apache.paimon.table.sink.TableWriteImpl.prepareCommit(TableWriteImpl.java:253)
at org.apache.paimon.flink.sink.StoreSinkWriteImpl.prepareCommit(StoreSinkWriteImpl.java:229)
... 33 more
Caused by: java.lang.RuntimeException: org.apache.paimon.shade.org.apache.parquet.io.ParquetDecodingException: Failed to read 4 bytes
at org.apache.paimon.reader.RecordReaderIterator.(RecordReaderIterator.java:40)
at org.apache.paimon.mergetree.compact.MergeTreeCompactRewriter.rewriteCompaction(MergeTreeCompactRewriter.java:89)
at org.apache.paimon.mergetree.compact.ChangelogMergeTreeRewriter.rewrite(ChangelogMergeTreeRewriter.java:108)
at org.apache.paimon.mergetree.compact.MergeTreeCompactTask.rewriteImpl(MergeTreeCompactTask.java:157)
at org.apache.paimon.mergetree.compact.MergeTreeCompactTask.rewrite(MergeTreeCompactTask.java:152)
at org.apache.paimon.mergetree.compact.MergeTreeCompactTask.doCompact(MergeTreeCompactTask.java:105)
at org.apache.paimon.compact.CompactTask.call(CompactTask.java:49)
at org.apache.paimon.compact.CompactTask.call(CompactTask.java:34)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
... 1 more
Caused by: org.apache.paimon.shade.org.apache.parquet.io.ParquetDecodingException: Failed to read 4 bytes
at org.apache.paimon.format.parquet.newreader.VectorizedPlainValuesReader.getBuffer(VectorizedPlainValuesReader.java:116)
at org.apache.paimon.format.parquet.newreader.VectorizedPlainValuesReader.readInteger(VectorizedPlainValuesReader.java:250)
at org.apache.paimon.format.parquet.newreader.VectorizedPlainValuesReader.readBinary(VectorizedPlainValuesReader.java:281)
at org.apache.paimon.format.parquet.newreader.ParquetVectorUpdaterFactory$BinaryUpdater.readValues(ParquetVectorUpdaterFactory.java:541)
at org.apache.paimon.format.parquet.newreader.ParquetVectorUpdaterFactory$BinaryUpdater.readValues(ParquetVectorUpdaterFactory.java:534)
at org.apache.paimon.format.parquet.newreader.VectorizedRleValuesReader.readBatchInternal(VectorizedRleValuesReader.java:268)
at org.apache.paimon.format.parquet.newreader.VectorizedRleValuesReader.readBatch(VectorizedRleValuesReader.java:188)
at org.apache.paimon.format.parquet.newreader.VectorizedColumnReader.readBatch(VectorizedColumnReader.java:213)
at org.apache.paimon.format.parquet.newreader.VectorizedParquetRecordReader.nextBatch(VectorizedParquetRecordReader.java:277)
at org.apache.paimon.format.parquet.newreader.VectorizedParquetRecordReader.readBatch(VectorizedParquetRecordReader.java:335)
at org.apache.paimon.io.DataFileRecordReader.readBatch(DataFileRecordReader.java:66)
at org.apache.paimon.io.KeyValueDataFileRecordReader.readBatch(KeyValueDataFileRecordReader.java:50)
at org.apache.paimon.io.KeyValueDataFileRecordReader.readBatch(KeyValueDataFileRecordReader.java:34)
at org.apache.paimon.mergetree.compact.LoserTree$LeafIterator.advanceIfAvailable(LoserTree.java:315)
at org.apache.paimon.mergetree.compact.LoserTree.initializeIfNeeded(LoserTree.java:87)
at org.apache.paimon.mergetree.compact.SortMergeReaderWithLoserTree.readBatch(SortMergeReaderWithLoserTree.java:71)
at org.apache.paimon.reader.RecordReaderIterator.(RecordReaderIterator.java:37)
... 13 more
Caused by: java.io.EOFException
at org.apache.paimon.shade.org.apache.parquet.bytes.SingleBufferInputStream.slice(SingleBufferInputStream.java:116)
at org.apache.paimon.format.parquet.newreader.VectorizedPlainValuesReader.getBuffer(VectorizedPlainValuesReader.java:114)
... 29 more
`
### Anything else?
_No response_
### Are you willing to submit a PR?
- [ ] I'm willing to submit a PR!
Contributor guide
No contributing guide indexed for this repository
Research direction
Start with the supplied Flink SQL and reproduce the failure around checkpointing. Trace the stack from StoreSinkWriteImpl.prepareCommit through MergeTreeCompactRewriter.rewriteCompaction to VectorizedPlainValuesReader.getBuffer and the reported EOFException. Done means the workload completes checkpoints without the Parquet decoding error.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- databases
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 28/100