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]

Reply via email to