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 @@ -87,7 +87,7 @@ public void notifyNewDeletion(String fileName, DeletionVector deletionVector) {
public void mergeNewDeletion(String fileName, DeletionVector deletionVector) {
DeletionVector old = deletionVectors.get(fileName);
if (old != null) {
deletionVector.merge(old);
deletionVector = DeletionVector.mergeVectors(deletionVector, old);
}
deletionVectors.put(fileName, deletionVector);
modified = true;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,27 @@ static DeletionVector deserializeFromBytes(byte[] bytes) {
}
}

/**
* Merges {@code stored} into {@code fresh} and returns the vector holding the union.
*
* <p>{@code deletion-vectors.bitmap64} can be flipped on an existing table, so a freshly
* created vector and a previously stored one may not share the same bitmap width. When they
* differ, the {@link BitmapDeletionVector} side is promoted to {@link Bitmap64DeletionVector}
* so the merge runs at the wider format instead of {@link #merge} throwing on the type
* mismatch. A bitmap64 result stays readable under a bitmap32 table because vectors are
* dispatched by magic number on read.
*/
static DeletionVector mergeVectors(DeletionVector fresh, DeletionVector stored) {
if (fresh instanceof Bitmap64DeletionVector && stored instanceof BitmapDeletionVector) {
stored = Bitmap64DeletionVector.fromBitmapDeletionVector((BitmapDeletionVector) stored);
} else if (fresh instanceof BitmapDeletionVector
&& stored instanceof Bitmap64DeletionVector) {
fresh = Bitmap64DeletionVector.fromBitmapDeletionVector((BitmapDeletionVector) fresh);
}
fresh.merge(stored);
return fresh;
}

/** Interface to create {@link DeletionVector}. */
interface Factory {
Optional<DeletionVector> create(String fileName) throws IOException;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,8 @@ public DeletionFile notifyRemovedDeletionVector(String dataFile) {
public void notifyNewDeletionVector(String dataFile, DeletionVector deletionVector) {
DeletionFile previous = notifyRemovedDeletionVector(dataFile);
if (previous != null) {
deletionVector.merge(dvIndexFile.readDeletionVector(previous));
DeletionVector stored = dvIndexFile.readDeletionVector(previous);
deletionVector = DeletionVector.mergeVectors(deletionVector, stored);
}
deletionVectors.put(dataFile, deletionVector);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -283,6 +283,52 @@ public void testReadAndWriteMixedDv(boolean bitmap64) {
assertThat(dvs.get("f3").getCardinality()).isEqualTo(2);
}

@ParameterizedTest
@ValueSource(booleans = {true, false})
public void testMergeNewDeletionAcrossBitmap64Flip(boolean initialBitmap64) {
// write a stored dv under one bitmap64 setting, then flip the option and merge a fresh
// vector of the opposite type into it: the merge must not crash on the type mismatch
initIndexHandler(initialBitmap64);
BucketedDvMaintainer.Factory factory1 = BucketedDvMaintainer.factory(fileHandler);
BucketedDvMaintainer dvMaintainer1 = factory1.create(partition, 0, new HashMap<>());
dvMaintainer1.notifyNewDeletion("f1", 1);
dvMaintainer1.notifyNewDeletion("f1", 3);
commitIndexFile(dvMaintainer1.writeDeletionVectorsIndex().get());

initIndexHandler(!initialBitmap64);
BucketedDvMaintainer.Factory factory2 = BucketedDvMaintainer.factory(fileHandler);
List<IndexFileMeta> indexFiles =
fileHandler.scan(
table.latestSnapshot().get(), DELETION_VECTORS_INDEX, partition, 0);
BucketedDvMaintainer dvMaintainer2 = factory2.create(partition, 0, indexFiles);

DeletionVector fresh = createDeletionVector(!initialBitmap64);
fresh.delete(5);
dvMaintainer2.mergeNewDeletion("f1", fresh);

DeletionVector merged = dvMaintainer2.deletionVectorOf("f1").get();
assertThat(merged.isDeleted(1)).isTrue();
assertThat(merged.isDeleted(3)).isTrue();
assertThat(merged.isDeleted(5)).isTrue();
assertThat(merged.getCardinality()).isEqualTo(3);
}

private void commitIndexFile(IndexFileMeta file) {
CommitMessage commitMessage =
new CommitMessageImpl(
partition,
0,
1,
DataIncrement.emptyIncrement(),
new CompactIncrement(
Collections.emptyList(),
Collections.emptyList(),
Collections.emptyList(),
Collections.singletonList(file),
Collections.emptyList()));
table.newBatchWriteBuilder().newCommit().commit(Collections.singletonList(commitMessage));
}

private DeletionVector createDeletionVector(boolean bitmap64) {
return bitmap64 ? new Bitmap64DeletionVector() : new BitmapDeletionVector();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
import org.apache.paimon.CoreOptions;
import org.apache.paimon.TestAppendFileStore;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.deletionvectors.Bitmap64DeletionVector;
import org.apache.paimon.deletionvectors.BitmapDeletionVector;
import org.apache.paimon.deletionvectors.DeletionVector;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.local.LocalFileIO;
Expand All @@ -32,6 +34,7 @@
import org.apache.paimon.table.sink.CommitMessageImpl;
import org.apache.paimon.table.source.DeletionFile;

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
Expand All @@ -48,6 +51,88 @@ class AppendDeletionFileMaintainerTest {

@TempDir java.nio.file.Path tempDir;

@Test
public void testMergeStoredBitmap32IntoBitmap64() throws Exception {
// write DVs as bitmap32, then flip deletion-vectors.bitmap64 on: the next
// notification must merge the stored vector instead of crashing on the type
Map<String, String> options = new HashMap<>();
options.put(CoreOptions.DELETION_VECTOR_BITMAP64.key(), "false");
TestAppendFileStore store = TestAppendFileStore.createAppendStore(tempDir, options);

CommitMessageImpl commitMessage =
store.writeDVIndexFiles(
BinaryRow.EMPTY_ROW,
0,
Collections.singletonMap("f1", Arrays.asList(1, 3, 5)));
store.commit(commitMessage);

IndexPathFactory indexPathFactory =
store.pathFactory().indexFileFactory(BinaryRow.EMPTY_ROW, 0);
Map<String, DeletionFile> dataFileToDeletionFiles =
createDeletionFileMapFromIndexFileMetas(
indexPathFactory, commitMessage.newFilesIncrement().newIndexFiles());

Map<String, String> flipped = new HashMap<>();
flipped.put(CoreOptions.DELETION_VECTOR_BITMAP64.key(), "true");
TestAppendFileStore flippedStore = TestAppendFileStore.createAppendStore(tempDir, flipped);
AppendDeleteFileMaintainer dvIFMaintainer =
flippedStore.createDVIFMaintainer(BinaryRow.EMPTY_ROW, dataFileToDeletionFiles);

Bitmap64DeletionVector fresh = new Bitmap64DeletionVector();
fresh.delete(7);
dvIFMaintainer.notifyNewDeletionVector("f1", fresh);

List<IndexManifestEntry> res = dvIFMaintainer.persist();
assertThat(res).hasSize(2);
// the old index file is replaced by one holding the merged vector: stored 3
// deletions plus the fresh one
assertThat(res).anyMatch(entry -> entry.kind() == FileKind.DELETE);
IndexManifestEntry added =
res.stream().filter(entry -> entry.kind() == FileKind.ADD).findAny().get();
assertThat(added.indexFile().dvRanges()).containsKey("f1");
assertThat(added.indexFile().dvRanges().get("f1").cardinality()).isEqualTo(4);
}

@Test
public void testMergeStoredBitmap64IntoBitmap32() throws Exception {
// reverse of the case above: write DVs as bitmap64, then flip deletion-vectors.bitmap64
// off. A fresh bitmap32 vector must still merge the stored bitmap64 one instead of crashing
Map<String, String> options = new HashMap<>();
options.put(CoreOptions.DELETION_VECTOR_BITMAP64.key(), "true");
TestAppendFileStore store = TestAppendFileStore.createAppendStore(tempDir, options);

CommitMessageImpl commitMessage =
store.writeDVIndexFiles(
BinaryRow.EMPTY_ROW,
0,
Collections.singletonMap("f1", Arrays.asList(1, 3, 5)));
store.commit(commitMessage);

IndexPathFactory indexPathFactory =
store.pathFactory().indexFileFactory(BinaryRow.EMPTY_ROW, 0);
Map<String, DeletionFile> dataFileToDeletionFiles =
createDeletionFileMapFromIndexFileMetas(
indexPathFactory, commitMessage.newFilesIncrement().newIndexFiles());

Map<String, String> flipped = new HashMap<>();
flipped.put(CoreOptions.DELETION_VECTOR_BITMAP64.key(), "false");
TestAppendFileStore flippedStore = TestAppendFileStore.createAppendStore(tempDir, flipped);
AppendDeleteFileMaintainer dvIFMaintainer =
flippedStore.createDVIFMaintainer(BinaryRow.EMPTY_ROW, dataFileToDeletionFiles);

BitmapDeletionVector fresh = new BitmapDeletionVector();
fresh.delete(7);
dvIFMaintainer.notifyNewDeletionVector("f1", fresh);

List<IndexManifestEntry> res = dvIFMaintainer.persist();
assertThat(res).hasSize(2);
assertThat(res).anyMatch(entry -> entry.kind() == FileKind.DELETE);
IndexManifestEntry added =
res.stream().filter(entry -> entry.kind() == FileKind.ADD).findAny().get();
assertThat(added.indexFile().dvRanges()).containsKey("f1");
assertThat(added.indexFile().dvRanges().get("f1").cardinality()).isEqualTo(4);
}

@ParameterizedTest
@ValueSource(booleans = {true, false})
public void test(boolean bitmap64) throws Exception {
Expand Down
Loading