From 28d3cc2d4e3a9c674a2465394f7d1e32c5ebeff7 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 01:54:50 +0800 Subject: [PATCH 1/2] [core] Skip empty changelog files in safelyGetAllChangelogs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The empty-file branch logged a warning but still fell through to Changelog.fromJson, which throws UncheckedIOException on the empty string and escaped the IOException-only catch, killing the whole enumeration — and with it the orphan files clean — over one torn or truncated changelog file. Skip the file after the warning, as the branch intended. Assisted-by: GLM-5.3 --- .../apache/paimon/utils/ChangelogManager.java | 5 +- .../paimon/utils/ChangelogManagerTest.java | 60 +++++++++++++++++++ 2 files changed, 64 insertions(+), 1 deletion(-) create mode 100644 paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java b/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java index 7d60649443bf..62fed024cc79 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java @@ -156,9 +156,12 @@ public List safelyGetAllChangelogs() throws IOException { try { String changelogStr = fileIO.readFileUtf8(path); if (StringUtils.isNullOrWhitespaceOnly(changelogStr)) { + // skip a torn or truncated changelog file instead of letting the + // empty parse failure kill the whole enumeration LOG.warn("Changelog file is empty, path: {}", path); + } else { + changelogs.add(Changelog.fromJson(changelogStr)); } - changelogs.add(Changelog.fromJson(changelogStr)); } catch (IOException e) { if (!(e instanceof FileNotFoundException)) { throw new RuntimeException(e); diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java new file mode 100644 index 000000000000..9322d62ea5d3 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java @@ -0,0 +1,60 @@ +/* + * 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.Changelog; +import org.apache.paimon.Snapshot; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.local.LocalFileIO; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import static org.apache.paimon.utils.SnapshotManagerTest.createSnapshotWithMillis; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link ChangelogManager}. */ +public class ChangelogManagerTest { + + @TempDir java.nio.file.Path tempDir; + + private FileIO fileIO; + private ChangelogManager changelogManager; + + @BeforeEach + public void before() { + fileIO = LocalFileIO.create(); + changelogManager = new ChangelogManager(fileIO, new Path(tempDir.toUri().toString()), null); + } + + @Test + public void testSafelyGetAllChangelogsSkipsEmptyFile() throws Exception { + Snapshot snapshot = createSnapshotWithMillis(1, 1000); + changelogManager.commitChangelog(new Changelog(snapshot), 1); + + // a torn write can leave an empty changelog file behind + fileIO.writeFile(changelogManager.longLivedChangelogPath(2), "", true); + + assertThat(changelogManager.safelyGetAllChangelogs()) + .extracting(Changelog::id) + .containsExactly(1L); + } +} From 7e577afcadb279d20843e7d5448cb5026c7b5037 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 02:04:58 +0800 Subject: [PATCH 2/2] [core] Write changelog files atomically commitChangelog wrote the JSON directly to the target file with overwrite, so a crash midway left readers with an empty or partial changelog file, which safelyGetAllChangelogs then had to tolerate. Write through a temp file and rename like snapshot and schema commits, and treat an already-existing file as an idempotent retry only when its content matches. Assisted-by: GLM-5.3 --- .../apache/paimon/utils/ChangelogManager.java | 14 +++++++- .../paimon/utils/ChangelogManagerTest.java | 33 +++++++++++++++++++ 2 files changed, 46 insertions(+), 1 deletion(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java b/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java index 62fed024cc79..2851ff6ea6fd 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/ChangelogManager.java @@ -121,7 +121,19 @@ public Path changelogDirectory() { } public void commitChangelog(Changelog changelog, long id) throws IOException { - fileIO.writeFile(longLivedChangelogPath(id), changelog.toJson(), true); + Path changelogPath = longLivedChangelogPath(id); + boolean committed = fileIO.tryToWriteAtomic(changelogPath, changelog.toJson()); + if (!committed) { + if (!fileIO.exists(changelogPath)) { + throw new IOException( + "Commit changelog " + id + " failed and " + changelogPath + " not found"); + } + // the file exists from a previous attempt; the same id commits the same content + if (!changelog.equals(Changelog.fromJson(fileIO.readFileUtf8(changelogPath)))) { + throw new IOException( + "Changelog file " + changelogPath + " exists with different content"); + } + } } public void commitLongLivedChangelogLatestHint(long snapshotId) throws IOException { diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java index 9322d62ea5d3..f482b7dd9275 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/ChangelogManagerTest.java @@ -27,9 +27,14 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import org.mockito.ArgumentMatchers; +import org.mockito.Mockito; + +import java.io.IOException; import static org.apache.paimon.utils.SnapshotManagerTest.createSnapshotWithMillis; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Tests for {@link ChangelogManager}. */ public class ChangelogManagerTest { @@ -45,6 +50,34 @@ public void before() { changelogManager = new ChangelogManager(fileIO, new Path(tempDir.toUri().toString()), null); } + @Test + public void testCommitChangelogWritesAtomically() throws Exception { + FileIO spyIO = Mockito.spy(fileIO); + ChangelogManager spyManager = + new ChangelogManager(spyIO, new Path(tempDir.toUri().toString()), null); + Changelog changelog = new Changelog(createSnapshotWithMillis(1, 1000)); + Path changelogPath = spyManager.longLivedChangelogPath(1); + + spyManager.commitChangelog(changelog, 1); + + // the target must never be opened for a direct overwrite: a crash midway would + // leave readers with an empty or partial changelog file + Mockito.verify(spyIO, Mockito.never()) + .writeFile( + ArgumentMatchers.eq(changelogPath), + ArgumentMatchers.anyString(), + ArgumentMatchers.eq(true)); + // readFileUtf8 strips newlines, so compare through a parse round-trip + assertThat(Changelog.fromJson(fileIO.readFileUtf8(changelogPath))).isEqualTo(changelog); + + // retrying the same commit is idempotent, a different content for the same id fails + spyManager.commitChangelog(changelog, 1); + Changelog other = new Changelog(createSnapshotWithMillis(1, 2000)); + assertThatThrownBy(() -> spyManager.commitChangelog(other, 1)) + .isInstanceOf(IOException.class) + .hasMessageContaining("exists with different content"); + } + @Test public void testSafelyGetAllChangelogsSkipsEmptyFile() throws Exception { Snapshot snapshot = createSnapshotWithMillis(1, 1000);