Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -80,6 +81,10 @@ public class MergeTreeWriter implements RecordWriter<KeyValue>, MemoryOwner {
private final LinkedHashMap<String, DataFileMeta> compactBefore;
private final LinkedHashSet<DataFileMeta> compactAfter;
private final LinkedHashSet<DataFileMeta> compactChangelog;
private final List<IndexFileMeta> newIndexFiles;
private final List<IndexFileMeta> deletedIndexFiles;
private final List<IndexFileMeta> compactNewIndexFiles;
private final List<IndexFileMeta> compactDeletedIndexFiles;

@Nullable private CompactDeletionFile compactDeletionFile;

Expand Down Expand Up @@ -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());
}
}
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -39,13 +42,15 @@
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;

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;

import java.io.IOException;
import java.util.Collections;
import java.util.Comparator;
import java.util.function.Function;

Expand Down Expand Up @@ -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())));

Expand Down Expand Up @@ -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));
}
Expand All @@ -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);
}
}
Loading