This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 4e10b49440 [core] Take the stream position outside the finally so
close always runs (#9233)
4e10b49440 is described below
commit 4e10b4944003c9ad68b2cdf6f90ea8256db7ddbb
Author: ZIHAN DAI <[email protected]>
AuthorDate: Mon Aug 17 15:31:58 2026 +1000
[core] Take the stream position outside the finally so close always runs
(#9233)
---
.../java/org/apache/paimon/utils/ObjectsFile.java | 7 +-
.../paimon/utils/ObjectsFileWriteFailureTest.java | 160 +++++++++++++++++++++
2 files changed, 163 insertions(+), 4 deletions(-)
diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsFile.java
b/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsFile.java
index 7d1a5d483a..e6c923ef7c 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsFile.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsFile.java
@@ -200,17 +200,16 @@ public abstract class ObjectsFile<T> implements
SimpleFileReader<T> {
}
return Pair.of(path.getName(), fileIO.getFileSize(path));
} else {
- PositionOutputStream out = fileIO.newOutputStream(path, false);
long pos;
- try {
+ try (PositionOutputStream out = fileIO.newOutputStream(path,
false)) {
+ // Nested rather than a single resource list: the position
has to be read
+ // after the writer has flushed, and before the stream
itself is closed.
try (FormatWriter writer = writerFactory.create(out,
compression)) {
while (records.hasNext()) {
writer.addElement(serializer.toRow(records.next()));
}
}
- } finally {
pos = out.getPos();
- out.close();
}
return Pair.of(path.getName(), pos);
}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/utils/ObjectsFileWriteFailureTest.java
b/paimon-core/src/test/java/org/apache/paimon/utils/ObjectsFileWriteFailureTest.java
new file mode 100644
index 0000000000..5c06fe58ab
--- /dev/null
+++
b/paimon-core/src/test/java/org/apache/paimon/utils/ObjectsFileWriteFailureTest.java
@@ -0,0 +1,160 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.utils;
+
+import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.format.FormatWriter;
+import org.apache.paimon.fs.FileIO;
+import org.apache.paimon.fs.Path;
+import org.apache.paimon.fs.PositionOutputStream;
+import org.apache.paimon.fs.local.LocalFileIO;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.io.IOException;
+import java.util.Collections;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/**
+ * When both the write and the closing of the output stream fail, {@link
+ * ObjectsFile#writeWithoutRolling} must report the write failure and keep the
close failure as a
+ * suppressed exception rather than letting the latter replace the former.
+ */
+class ObjectsFileWriteFailureTest {
+
+ @Test
+ void writeFailureSurvivesAFailingStreamClose(@TempDir java.nio.file.Path
tempDir) {
+ FailingStream stream = new FailingStream();
+ ObjectsFile<String> file = objectsFile(tempDir, stream, true);
+
+ assertThatThrownBy(() ->
file.writeWithoutRolling(Collections.emptyIterator()))
+ .isInstanceOf(RuntimeException.class)
+ .cause()
+ .hasMessage("writer close failed")
+ .satisfies(
+ cause ->
+ assertThat(cause.getSuppressed())
+ .extracting(Throwable::getMessage)
+ .containsExactly("stream close
failed"));
+
+ assertThat(stream.closed).isTrue();
+ }
+
+ /** The stream is closed even when the write succeeds and the position
read is the last step. */
+ @Test
+ void streamIsClosedOnTheSuccessPath(@TempDir java.nio.file.Path tempDir)
throws Exception {
+ FailingStream stream = new FailingStream();
+ stream.failOnClose = false;
+ ObjectsFile<String> file = objectsFile(tempDir, stream, false);
+
+
assertThat(file.writeWithoutRolling(Collections.emptyIterator()).getValue()).isEqualTo(7L);
+ assertThat(stream.closed).isTrue();
+ }
+
+ private static ObjectsFile<String> objectsFile(
+ java.nio.file.Path tempDir, PositionOutputStream stream, boolean
writerCloseFails) {
+ Path path = new Path(tempDir.toUri().toString(), "manifest-0");
+ FileIO fileIO =
+ new LocalFileIO() {
+ @Override
+ public PositionOutputStream newOutputStream(Path file,
boolean overwrite) {
+ return stream;
+ }
+ };
+ return new ObjectsFile<String>(
+ fileIO,
+ null,
+ null,
+ (f, size) -> {
+ throw new UnsupportedOperationException();
+ },
+ (out, compression) -> new StubWriter(writerCloseFails),
+ "none",
+ new PathFactory() {
+ @Override
+ public Path newPath() {
+ return path;
+ }
+
+ @Override
+ public Path toPath(String fileName) {
+ return path;
+ }
+ },
+ null) {};
+ }
+
+ private static class StubWriter implements FormatWriter {
+
+ private final boolean failOnClose;
+
+ private StubWriter(boolean failOnClose) {
+ this.failOnClose = failOnClose;
+ }
+
+ @Override
+ public void addElement(InternalRow element) {}
+
+ @Override
+ public boolean reachTargetSize(boolean suggestedCheck, long
targetSize) {
+ return false;
+ }
+
+ @Override
+ public void close() throws IOException {
+ if (failOnClose) {
+ throw new IOException("writer close failed");
+ }
+ }
+ }
+
+ private static class FailingStream extends PositionOutputStream {
+
+ private boolean failOnClose = true;
+ private boolean closed = false;
+
+ @Override
+ public long getPos() {
+ return 7L;
+ }
+
+ @Override
+ public void write(int b) {}
+
+ @Override
+ public void write(byte[] b) {}
+
+ @Override
+ public void write(byte[] b, int off, int len) {}
+
+ @Override
+ public void flush() {}
+
+ @Override
+ public void close() throws IOException {
+ closed = true;
+ if (failOnClose) {
+ throw new IOException("stream close failed");
+ }
+ }
+ }
+}