apache / apache/paimon

[Bug] Failed to finish checkpoint due to the expection: 'Failed to read 4 bytes'

Open
#5,434 8 comments 0 reactions 0 assignees View on GitHub
bug
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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.