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 @@ -74,7 +74,6 @@
import java.util.stream.IntStream;

import static org.apache.paimon.CoreOptions.LOOKUP_CACHE_ROWS;
import static org.apache.paimon.utils.Preconditions.checkArgument;

/** Assign UPDATE_BEFORE and bucket for the input record, output record with bucket. */
public class GlobalIndexAssigner implements Serializable, Closeable {
Expand Down Expand Up @@ -186,15 +185,24 @@ public void open(
}

public void bootstrapKey(InternalRow value) throws IOException {
checkArgument(inBoostrap());
BinaryRow partition = keyPartExtractor.partition(value);
BinaryRow key = keyPartExtractor.trimmedPrimaryKey(value);
int partId = partMapping.index(partition);
int bucket = value.getInt(bucketIndex);
bucketAssigner.bootstrapBucket(partition, bucket);
PositiveIntInt partAndBucket = new PositiveIntInt(partId, bucket);
bootstrapKeys.write(
GenericRow.of(keyIndex.serializeKey(key), keyIndex.serializeValue(partAndBucket)));
if (inBoostrap()) {
bootstrapKeys.write(
GenericRow.of(
keyIndex.serializeKey(key), keyIndex.serializeValue(partAndBucket)));
} else if (keyIndex.get(key) == null) {
// with unaligned checkpoints the bootstrap can end on the barrier while KEY_PART
// records are still queued: register a late key directly instead of crashing, but
// only when no input record has already assigned it -- a newer input must win over
// the bootstrapped (pre-checkpoint) state, otherwise this stale put would leave
// keyIndex pointing at a bucket that no longer holds the record
keyIndex.put(key, partAndBucket);
}
}

public boolean inBoostrap() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,89 @@ public void testFirstRow() throws Exception {
assigner.close();
}

@Test
public void testLateBootstrapKeyAfterEndRegistersKey() throws Exception {
// with unaligned checkpoints the bootstrap can end on the barrier while KEY_PART
// records are still queued: the late key must register instead of crashing
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE);
List<List<Integer>> output = new ArrayList<>();
assigner.open(
0,
null,
ioManager(),
2,
0,
(row, bucket) ->
output.add(
Arrays.asList(
row.getInt(0), row.getInt(1), row.getInt(2), bucket)));

assigner.endBoostrap(false);
// a KEY_PART record queued before the barrier now arrives: (pk, pt, bucket)
assigner.bootstrapKey(GenericRow.of(1, 2, 2));

assigner.processInput(GenericRow.of(2, 1, 2));

assertThat(output).containsExactlyInAnyOrder(Arrays.asList(2, 1, 2, 2));
assigner.close();
}

@Test
public void testLateBootstrapKeyDoesNotOverwriteAssignedKey() throws Exception {
// with unaligned checkpoints an input ROW can be processed before a late KEY_PART for the
// same primary key arrives; the newer input must win, the late bootstrap state must not
// overwrite the freshly assigned index entry
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE);
List<List<Integer>> output = new ArrayList<>();
assigner.open(
0,
null,
ioManager(),
2,
0,
(row, bucket) ->
output.add(
Arrays.asList(
row.getInt(0), row.getInt(1), row.getInt(2), bucket)));

assigner.endBoostrap(false);
// input ROW for pk=1 in partition 9 is processed first: registers keyIndex[1]=(pt=9, 0)
assigner.processInput(GenericRow.of(9, 1, 1));
// a KEY_PART for the same pk=1 carrying the stale pre-checkpoint location (pt=2, bucket=2)
// arrives late: (pk, pt, bucket). It must not overwrite the freshly assigned entry
assigner.bootstrapKey(GenericRow.of(1, 2, 2));
// a later same-pk record in partition 9 must still route to its assigned bucket (0),
// proving keyIndex still points at the input location, not the stale bootstrap one
assigner.processInput(GenericRow.of(9, 1, 5));

assertThat(output).containsExactly(Arrays.asList(9, 1, 1, 0), Arrays.asList(9, 1, 5, 0));
assigner.close();
}

@Test
public void testLateBootstrapKeyRegistersAndRetractsOldPartition() throws Exception {
// a late KEY_PART arriving before any same-pk input must register the pre-checkpoint
// location, so a following cross-partition input retracts the old row instead of
// leaving a cross-partition duplicate -- this pins the late-key keyIndex.put itself
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE);
List<Pair<InternalRow, Integer>> output = new ArrayList<>();
assigner.open(
0, null, ioManager(), 2, 0, (row, bucket) -> output.add(Pair.of(row, bucket)));

assigner.endBoostrap(false);
// late KEY_PART registers keyIndex[pk=1]=(pt=2, bucket=2): (pk, pt, bucket)
assigner.bootstrapKey(GenericRow.of(1, 2, 2));
// a cross-partition input for the same pk must see the registered old location (pt=2)
// and retract it before assigning the new partition (pt=9)
assigner.processInput(GenericRow.of(9, 1, 5));

Assertions.assertThat(output)
.containsExactly(
Pair.of(GenericRow.ofKind(RowKind.DELETE, 2, 1, 5), 2),
Pair.of(GenericRow.of(9, 1, 5), 0));
assigner.close();
}

@Test
public void testBootstrapRecords() throws Exception {
GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE);
Expand Down
Loading