danzhewuju opened a new issue, #9932: URL: https://github.com/apache/paimon/issues/9932
### Search before asking - [x] I searched in the [issues](https://github.com/apache/paimon/issues) and found nothing similar. ### Paimon version The problem was originally observed on a Paimon 1.4-based build. ### Compute Engine Apache Flink 2.2. The problem is in `paimon-core` and can also be reproduced with a unit test without starting a Flink cluster. ### Minimal reproduce step Create a `RollingFileWriterImpl` whose `SingleFileWriter`: 1. Creates a physical output file. 2. Throws an `OutOfMemoryError` from `FormatWriter.addElement`. 3. Uses a `result()` implementation that reads the generated file metadata, for example by calling `FileIO#getFileSize`. Then execute the same lifecycle used by `MergeTreeWriter#flushWriteBuffer`: ```java OutOfMemoryError expected = new OutOfMemoryError("expected"); RollingFileWriterImpl<InternalRow, Void> writer = createRollingWriterWhoseAddElementThrows(expected); assertThatThrownBy(() -> writer.write(GenericRow.of(1))) .isSameAs(expected); // Simulates the close from MergeTreeWriter's finally block. writer.close(); ``` The relevant production sequence is: MergeTreeWriter.flushWriteBuffer -> dataWriter.write -> RollingFileWriterImpl.write -> SingleFileWriter.write -> FormatWriter.addElement throws OutOfMemoryError -> SingleFileWriter.abort deletes the uncommitted file -> RollingFileWriterImpl.abort leaves currentWriter referenced -> MergeTreeWriter finally calls dataWriter.close -> RollingFileWriterImpl.closeCurrentWriter -> currentWriter.result -> SingleFileWriter.outputBytes -> FileIO.getFileSize on the deleted file -> FileNotFoundException On master, RollingFileWriterImpl#abort aborts currentWriter but does not clear it or mark the rolling writer as closed: ```java public void abort() { if (currentWriter != null) { currentWriter.abort(); } for (FileWriterAbortExecutor abortExecutor : closedWriters) { abortExecutor.abort(); } } ``` MergeTreeWriter#flushWriteBuffer always closes the writer in finally: ```java try { writeBuffer.forEach(..., dataWriter::write); } finally { writeBuffer.clear(); IOUtils.closeAll(changelogWriter, dataWriter); } ``` The subsequent close still calls result() on the aborted writer: ```java currentWriter.close(); currentWriter.abortExecutor().ifPresent(closedWriters::add); results.add(currentWriter.result()); currentWriter = null; ``` A typical final stack trace is therefore misleading: ``` java java.io.IOException: Could not perform checkpoint 871 for operator Writer(write-only) : *** (10/16)#34. at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1420) 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$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:182) at org.apache.flink.streaming.runtime.io.checkpointing.CheckpointedInputGate.pollNext(CheckpointedInputGate.java:160) at org.apache.flink.streaming.runtime.io.AbstractStreamTaskNetworkInput.emitNext(AbstractStreamTaskNetworkInput.java:171) at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65) at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:650) at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231) at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:992) at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:929) at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:977) at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:959) at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:764) at org.apache.flink.runtime.taskmanager.Task.run(Task.java:569) at java.base/java.lang.Thread.run(Thread.java:1583) Caused by: java.io.IOException: java.io.FileNotFoundException: Caused by error 30001: File *** [RequestId]: 6AAB99BD49D64C3235BB31D2 at org.apache.paimon.flink.sink.StoreSinkWriteImpl.prepareCommit(StoreSinkWriteImpl.java:147) at org.apache.paimon.flink.sink.TableWriteOperator.prepareCommit(TableWriteOperator.java:143) at org.apache.paimon.flink.sink.RowDataStoreWriteOperator.prepareCommit(RowDataStoreWriteOperator.java:70) at org.apache.paimon.flink.sink.PrepareCommitOperator.emitCommittables(PrepareCommitOperator.java:115) at org.apache.paimon.flink.sink.PrepareCommitOperator.prepareSnapshotPreBarrier(PrepareCommitOperator.java:95) at org.apache.flink.streaming.runtime.tasks.RegularOperatorChain.prepareSnapshotPreBarrier(RegularOperatorChain.java:89) at org.apache.flink.streaming.runtime.tasks.SubtaskCheckpointCoordinatorImpl.checkpointState(SubtaskCheckpointCoordinatorImpl.java:333) at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$performCheckpoint$18(StreamTask.java:1463) at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$1.runThrowing(StreamTaskActionExecutor.java:50) at org.apache.flink.streaming.runtime.tasks.StreamTask.performCheckpoint(StreamTask.java:1451) at org.apache.flink.streaming.runtime.tasks.StreamTask.triggerCheckpointOnBarrier(StreamTask.java:1408) ... 21 more Caused by: java.io.FileNotFoundException: Caused by error 30001: File **** [RequestId]: 6AAB99BD49D64C3235BB31D2 at com.aliyun.jindodata.api.spec.JdoNativeResult.get(JdoNativeResult.java:54) at com.aliyun.jindodata.api.spec.protos.coder.JdoGetFileStatusReplyDecoder.decode(JdoGetFileStatusReplyDecoder.java:23) at com.aliyun.jindodata.api.JindoCommonApis.getFileStatus(JindoCommonApis.java:89) at com.aliyun.jindodata.call.JindoGetFileStatusCall.execute(JindoGetFileStatusCall.java:55) at com.aliyun.jindodata.common.JindoHadoopSystem.getFileStatus(JindoHadoopSystem.java:1011) at com.ctrip.di.oss.fs.hdfs.TripJindoHdfsFileSystem.getFileStatus(TripJindoHdfsFileSystem.java:302) at org.apache.paimon.fs.hadoop.HadoopSecuredFileSystem.lambda$getFileStatus$15(HadoopSecuredFileSystem.java:169) at java.base/javax.security.auth.Subject.lambda$callAs$0(Subject.java:379) at java.base/java.security.AccessController.doPrivileged(AccessController.java:714) at java.base/javax.security.auth.Subject.doAs(Subject.java:525) at java.base/javax.security.auth.Subject.callAs(Subject.java:381) at org.apache.hadoop.util.SubjectUtil.callAs(SubjectUtil.java:161) at org.apache.hadoop.util.SubjectUtil.doAs(SubjectUtil.java:233) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1741) at org.apache.paimon.fs.hadoop.HadoopSecuredFileSystem.runSecuredWithIOException(HadoopSecuredFileSystem.java:190) at org.apache.paimon.fs.hadoop.HadoopSecuredFileSystem.getFileStatus(HadoopSecuredFileSystem.java:169) at org.apache.paimon.fs.hadoop.HadoopFileIO.getFileStatus(HadoopFileIO.java:110) at org.apache.paimon.fs.FileIO.getFileSize(FileIO.java:295) at org.apache.paimon.io.SingleFileWriter.outputBytes(SingleFileWriter.java:217) at org.apache.paimon.io.KeyValueDataFileWriter.result(KeyValueDataFileWriter.java:148) at org.apache.paimon.io.KeyValueDataFileWriter.result(KeyValueDataFileWriter.java:55) at org.apache.paimon.io.RollingFileWriterImpl.closeCurrentWriter(RollingFileWriterImpl.java:133) at org.apache.paimon.io.RollingFileWriterImpl.close(RollingFileWriterImpl.java:165) at org.apache.paimon.mergetree.MergeTreeWriter.flushWriteBuffer(MergeTreeWriter.java:234) at org.apache.paimon.mergetree.MergeTreeWriter.prepareCommit(MergeTreeWriter.java:253) at org.apache.paimon.operation.AbstractFileStoreWrite.prepareCommit(AbstractFileStoreWrite.java:235) at org.apache.paimon.operation.MemoryFileStoreWrite.prepareCommit(MemoryFileStoreWrite.java:152) at org.apache.paimon.table.sink.TableWriteImpl.prepareCommit(TableWriteImpl.java:252) at org.apache.paimon.flink.sink.StoreSinkWriteImpl.prepareCommit(StoreSinkWriteImpl.java:143) ... 31 more ``` The original failure was OOM. ### What doesn't meet your expectations? Deleting the uncommitted file during abort() is expected. However, an aborted rolling writer should be in a terminal state and a later close() should be idempotent. It should not call result() for a file that has already been deleted. More importantly, the original write failure must remain the primary exception. In the Flink sink path, the FileNotFoundException thrown from the finally block replaces the original OutOfMemoryError. As a result: - Flink reports that the checkpoint failed because a data file does not exist. - The actual JVM heap OOM is hidden. - Users may incorrectly investigate filesystem consistency or external file deletion. - The same file can be deleted and queried repeatedly by nested abort/close paths. Expected behavior: `OutOfMemoryError: Java heap space` ### Anything else? _No response_ ### Are you willing to submit a PR? - [x] I'm willing to submit a PR! -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
