From 6dc24794ba97aab262a50e747add984b285e6fcd Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 20:54:46 +0800 Subject: [PATCH 1/3] [core] Register late bootstrap keys after unaligned checkpoints With execution.checkpointing.unaligned=true the checkpoint barrier can run the assigner's endBootstrap while KEY_PART records from slower channels are still queued. The first late bootstrapKey then failed the in-bootstrap assertion, failing the task and likely restart-looping under the same backpressure that motivated unaligned checkpoints. A late KEY_PART record's purpose is only to register its key in the index, which by then is bulk-loaded: put the key directly instead of crashing. Assisted-by: GLM-5.3 --- .../crosspartition/GlobalIndexAssigner.java | 13 ++++++--- .../GlobalIndexAssignerTest.java | 27 +++++++++++++++++++ 2 files changed, 36 insertions(+), 4 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java index eb56a537fcf0..c1358ee70f28 100644 --- a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java +++ b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java @@ -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 { @@ -186,15 +185,21 @@ 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 { + // with unaligned checkpoints the bootstrap can end on the barrier while KEY_PART + // records are still queued: register the key directly instead of crashing + keyIndex.put(key, partAndBucket); + } } public boolean inBoostrap() { diff --git a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java index 821521d2554b..4dc660aa0aae 100644 --- a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java @@ -216,6 +216,33 @@ 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> 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 testBootstrapRecords() throws Exception { GlobalIndexAssigner assigner = createAssigner(MergeEngine.DEDUPLICATE); From 40703699c18f4005b233a524833be43112c47d5d Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Mon, 28 Sep 2026 15:49:37 +0800 Subject: [PATCH 2/3] fix: let a newer input win over a late bootstrap key Guard the late-KEY_PART put with keyIndex.get(key)==null so a KEY_PART arriving after endBoostrap under unaligned checkpoints cannot overwrite an index entry an input record already assigned. Add tests pinning both the no-overwrite guard and the late-key registration (cross-partition retraction). --- .../crosspartition/GlobalIndexAssigner.java | 7 ++- .../GlobalIndexAssignerTest.java | 57 +++++++++++++++++++ 2 files changed, 62 insertions(+), 2 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java index c1358ee70f28..bd7344c33e7b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java +++ b/paimon-core/src/main/java/org/apache/paimon/crosspartition/GlobalIndexAssigner.java @@ -195,9 +195,12 @@ public void bootstrapKey(InternalRow value) throws IOException { bootstrapKeys.write( GenericRow.of( keyIndex.serializeKey(key), keyIndex.serializeValue(partAndBucket))); - } else { + } else if (keyIndex.get(key) == null) { // with unaligned checkpoints the bootstrap can end on the barrier while KEY_PART - // records are still queued: register the key directly instead of crashing + // 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); } } diff --git a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java index 4dc660aa0aae..40a7e9d1606c 100644 --- a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java @@ -243,6 +243,63 @@ public void testLateBootstrapKeyAfterEndRegistersKey() throws Exception { 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> 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> 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); From 3134c2e293ea75aab183124db8ac2857781881b1 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Mon, 28 Sep 2026 16:55:13 +0800 Subject: [PATCH 3/3] style: apply spotless formatting to GlobalIndexAssignerTest --- .../apache/paimon/crosspartition/GlobalIndexAssignerTest.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java index 40a7e9d1606c..b6648971dd03 100644 --- a/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/crosspartition/GlobalIndexAssignerTest.java @@ -271,8 +271,7 @@ public void testLateBootstrapKeyDoesNotOverwriteAssignedKey() throws Exception { // 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)); + assertThat(output).containsExactly(Arrays.asList(9, 1, 1, 0), Arrays.asList(9, 1, 5, 0)); assigner.close(); }