From d2cc75f51248e3c85e3f75ed9c54d943b3c3c76b Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 03:01:24 +0800 Subject: [PATCH] [core] Replay index file changes when restoring MergeTreeWriter The restore constructor replayed the pending increment's data, changelog and compaction files but dropped newIndexFiles and deletedIndexFiles. On the checkpoint path taken by a CDC schema-change replace, the primary-key index maintainers finish and accept pending builds and merge their index file changes into that increment; after restore the changes were silently lost, so accepted payloads never reached the index manifest while replaced payloads stayed as zombie entries. Carry the index file changes through the replay and re-emit them on the next prepareCommit. Assisted-by: GLM-5.3 --- .../paimon/mergetree/MergeTreeWriter.java | 28 ++++++++++++- .../MergeTreeWriterCloseFailureTest.java | 41 ++++++++++++++++++- 2 files changed, 65 insertions(+), 4 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/MergeTreeWriter.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/MergeTreeWriter.java index 4f7b2488700b..2d58056b1e24 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/MergeTreeWriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/MergeTreeWriter.java @@ -27,6 +27,7 @@ import org.apache.paimon.compression.CompressOptions; import org.apache.paimon.data.InternalRow; import org.apache.paimon.disk.IOManager; +import org.apache.paimon.index.IndexFileMeta; import org.apache.paimon.io.CompactIncrement; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataIncrement; @@ -80,6 +81,10 @@ public class MergeTreeWriter implements RecordWriter, MemoryOwner { private final LinkedHashMap compactBefore; private final LinkedHashSet compactAfter; private final LinkedHashSet compactChangelog; + private final List newIndexFiles; + private final List deletedIndexFiles; + private final List compactNewIndexFiles; + private final List compactDeletedIndexFiles; @Nullable private CompactDeletionFile compactDeletionFile; @@ -123,16 +128,27 @@ public MergeTreeWriter( this.compactBefore = new LinkedHashMap<>(); this.compactAfter = new LinkedHashSet<>(); this.compactChangelog = new LinkedHashSet<>(); + this.newIndexFiles = new ArrayList<>(); + this.deletedIndexFiles = new ArrayList<>(); + this.compactNewIndexFiles = new ArrayList<>(); + this.compactDeletedIndexFiles = new ArrayList<>(); if (increment != null) { newFiles.addAll(increment.newFilesIncrement().newFiles()); deletedFiles.addAll(increment.newFilesIncrement().deletedFiles()); newFilesChangelog.addAll(increment.newFilesIncrement().changelogFiles()); + // index file changes must survive the restore replay too: a payload accepted + // during checkpoint would otherwise never reach the manifest, and replaced + // payloads would stay as zombie entries + newIndexFiles.addAll(increment.newFilesIncrement().newIndexFiles()); + deletedIndexFiles.addAll(increment.newFilesIncrement().deletedIndexFiles()); increment .compactIncrement() .compactBefore() .forEach(f -> compactBefore.put(f.fileName(), f)); compactAfter.addAll(increment.compactIncrement().compactAfter()); compactChangelog.addAll(increment.compactIncrement().changelogFiles()); + compactNewIndexFiles.addAll(increment.compactIncrement().newIndexFiles()); + compactDeletedIndexFiles.addAll(increment.compactIncrement().deletedIndexFiles()); updateCompactDeletionFile(increment.compactDeletionFile()); } } @@ -285,20 +301,28 @@ private CommitIncrement drainIncrement() { new DataIncrement( new ArrayList<>(newFiles), new ArrayList<>(deletedFiles), - new ArrayList<>(newFilesChangelog)); + new ArrayList<>(newFilesChangelog), + new ArrayList<>(newIndexFiles), + new ArrayList<>(deletedIndexFiles)); CompactIncrement compactIncrement = new CompactIncrement( new ArrayList<>(compactBefore.values()), new ArrayList<>(compactAfter), - new ArrayList<>(compactChangelog)); + new ArrayList<>(compactChangelog), + new ArrayList<>(compactNewIndexFiles), + new ArrayList<>(compactDeletedIndexFiles)); CompactDeletionFile drainDeletionFile = compactDeletionFile; newFiles.clear(); deletedFiles.clear(); newFilesChangelog.clear(); + newIndexFiles.clear(); + deletedIndexFiles.clear(); compactBefore.clear(); compactAfter.clear(); compactChangelog.clear(); + compactNewIndexFiles.clear(); + compactDeletedIndexFiles.clear(); compactDeletionFile = null; return new CommitIncrement(dataIncrement, compactIncrement, drainDeletionFile); diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeWriterCloseFailureTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeWriterCloseFailureTest.java index b2f9c223cc3d..4eb15249f5f2 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeWriterCloseFailureTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/MergeTreeWriterCloseFailureTest.java @@ -30,6 +30,9 @@ import org.apache.paimon.fs.Path; import org.apache.paimon.fs.PositionOutputStream; import org.apache.paimon.fs.PositionOutputStreamWrapper; +import org.apache.paimon.index.IndexFileMeta; +import org.apache.paimon.io.CompactIncrement; +import org.apache.paimon.io.DataIncrement; import org.apache.paimon.io.KeyValueFileWriterFactory; import org.apache.paimon.memory.HeapMemorySegmentPool; import org.apache.paimon.mergetree.compact.DeduplicateMergeFunction; @@ -39,6 +42,7 @@ import org.apache.paimon.types.IntType; import org.apache.paimon.types.RowKind; import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.CommitIncrement; import org.apache.paimon.utils.FileStorePathFactory; import org.apache.paimon.utils.TraceableFileIO; @@ -46,6 +50,7 @@ import org.junit.jupiter.api.io.TempDir; import java.io.IOException; +import java.util.Collections; import java.util.Comparator; import java.util.function.Function; @@ -83,7 +88,7 @@ void dataWriterIsClosedWhenTheChangelogWriterFails() throws Exception { .isEmpty(); } - private MergeTreeWriter createWriter(Path path) { + private MergeTreeWriter createWriter(Path path, CommitIncrement increment) { RowType keyType = new RowType(singletonList(new DataField(0, "k", new IntType()))); RowType valueType = new RowType(singletonList(new DataField(0, "v", new IntType()))); @@ -122,13 +127,17 @@ private MergeTreeWriter createWriter(Path path) { writerFactory, false, INPUT, - null, + increment, null); writer.setMemoryPool( new HeapMemorySegmentPool(coreOptions.writeBufferSize(), coreOptions.pageSize())); return writer; } + private MergeTreeWriter createWriter(Path path) { + return createWriter(path, null); + } + private KeyValue kv(int k, int v) { return new KeyValue().replace(GenericRow.of(k), RowKind.INSERT, GenericRow.of(v)); } @@ -147,4 +156,32 @@ public void close() throws IOException { }; } } + + @Test + void testRestoreReplayKeepsIndexFiles() throws Exception { + IndexFileMeta payload = new IndexFileMeta("hash", "index-1", 0, 0, null, null, null); + IndexFileMeta removed = new IndexFileMeta("hash", "index-0", 0, 0, null, null, null); + CommitIncrement increment = + new CommitIncrement( + new DataIncrement( + Collections.emptyList(), + Collections.emptyList(), + Collections.emptyList(), + Collections.singletonList(payload), + Collections.emptyList()), + new CompactIncrement( + Collections.emptyList(), + Collections.emptyList(), + Collections.emptyList(), + Collections.emptyList(), + Collections.singletonList(removed)), + null); + + java.nio.file.Path folder = java.nio.file.Files.createTempDirectory(tempDir, "restore"); + MergeTreeWriter writer = createWriter(new Path(folder.toUri()), increment); + CommitIncrement drained = writer.prepareCommit(false); + + assertThat(drained.newFilesIncrement().newIndexFiles()).containsExactly(payload); + assertThat(drained.compactIncrement().deletedIndexFiles()).containsExactly(removed); + } }