From cb0c2790e8bbea3ffcc66776f337c279a7234df1 Mon Sep 17 00:00:00 2001 From: QingweiYang Date: Mon, 28 Sep 2026 17:23:03 +0800 Subject: [PATCH] [core] Preserve partial-update sequence-group deletion state MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Port apache/paimon#10296 at b43e36137f2b5b3f51aaf7290c6f96d67744675a. Keep the production fix and deletion/restoration documentation unchanged. Reduce regression coverage to six methods and ten parameterized invocations: remove five-group permutations and retain empty/reordered projections. Co-Authored-By: 仟弋 Co-Authored-By: Codex AI-Model: unknown AI-Contributed/Feature: 11/11 AI-Contributed/UT: 223/223 --- .../merge-engine/partial-update.md | 5 + .../compact/PartialUpdateMergeFunction.java | 6 +- ...okupChangelogMergeFunctionWrapperTest.java | 54 ++++++ .../PartialUpdateMergeFunctionTest.java | 169 ++++++++++++++++++ 4 files changed, 233 insertions(+), 1 deletion(-) diff --git a/docs/docs/primary-key-table/merge-engine/partial-update.md b/docs/docs/primary-key-table/merge-engine/partial-update.md index 2430eee2932b..2299de1c1def 100644 --- a/docs/docs/primary-key-table/merge-engine/partial-update.md +++ b/docs/docs/primary-key-table/merge-engine/partial-update.md @@ -236,5 +236,10 @@ Set `partial-update.remove-record-on-sequence-group` to a comma-separated list o field names from the groups whose deletes should remove the whole row. For the `profiles` schema above, using `profile_version` allows an accepted profile delete to remove the row. +While merging a deleted row, subsequent retractions or updates to other groups do not cancel +the whole-row deletion. An `INSERT` or `UPDATE_AFTER` for a configured deletion group can +restore the row only when its group version passes the normal newer-or-equal comparison. +An older update or a group with all-null ordering fields cannot restore the row. + `partial-update.remove-record-on-delete` cannot be combined with sequence groups. Neither whole-row removal option can be combined with `ignore-delete`. diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunction.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunction.java index b97d7a29dac1..467524ee116d 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunction.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunction.java @@ -138,6 +138,7 @@ protected PartialUpdateMergeFunction( @Override public void reset() { this.currentKey = null; + this.currentDeleteRow = false; this.meetInsert = false; this.notNullColumnFilled = false; this.row = new GenericRow(getters.length); @@ -149,7 +150,6 @@ public void reset() { public void add(KeyValue kv) { // refresh key object to avoid reference overwritten currentKey = kv.key(); - currentDeleteRow = false; if (kv.valueKind().isRetract()) { if (!notNullColumnFilled) { @@ -193,6 +193,7 @@ public void add(KeyValue kv) { latestSequenceNumber = kv.sequenceNumber(); if (fieldSeqComparators.isEmpty()) { + currentDeleteRow = false; updateNonNullFields(kv); } else { updateWithSequenceGroup(kv); @@ -257,6 +258,9 @@ private void updateWithSequenceGroup(KeyValue kv) { if (Arrays.stream(seqComparator.compareFields()) .anyMatch(seqIndex -> seqIndex == index)) { for (int fieldIndex : seqComparator.compareFields()) { + if (sequenceGroupPartialDelete.contains(fieldIndex)) { + currentDeleteRow = false; + } row.setField( fieldIndex, getters[fieldIndex].getFieldOrNull(kv.value())); // Mark these sequence fields as processed diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java index db063462fb09..d126eb6b865e 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapperTest.java @@ -21,6 +21,7 @@ import org.apache.paimon.CoreOptions; import org.apache.paimon.KeyValue; import org.apache.paimon.codegen.RecordEqualiser; +import org.apache.paimon.data.GenericRow; import org.apache.paimon.data.InternalRow; import org.apache.paimon.data.InternalRow.FieldGetter; import org.apache.paimon.lookup.LookupStrategy; @@ -28,6 +29,7 @@ import org.apache.paimon.mergetree.compact.aggregate.FieldAggregator; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldLastValueAggFactory; import org.apache.paimon.mergetree.compact.aggregate.factory.FieldSumAggFactory; +import org.apache.paimon.options.Options; import org.apache.paimon.types.DataType; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; @@ -60,6 +62,58 @@ public class LookupChangelogMergeFunctionWrapperTest { private static final RecordEqualiser EQUALISER = (row1, row2) -> row1.getInt(0) == row2.getInt(0); + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testSequenceGroupDeleteProducesDeleteChangelog(boolean lookupOldRow) { + RowType valueType = + RowType.of( + DataTypes.INT().notNull(), + DataTypes.INT(), + DataTypes.INT(), + DataTypes.INT(), + DataTypes.INT()); + Options options = new Options(); + options.set("fields.f1.sequence-group", "f2"); + options.set("fields.f3.sequence-group", "f4"); + options.set("partial-update.remove-record-on-sequence-group", "f1"); + KeyValue oldRow = + new KeyValue() + .replace(row(1), 1, INSERT, GenericRow.of(1, 1, 10, 1, 20)) + .setLevel(2); + LookupChangelogMergeFunctionWrapper function = + new LookupChangelogMergeFunctionWrapper( + LookupMergeFunction.wrap( + PartialUpdateMergeFunction.factory( + options, valueType, Collections.singletonList("f0")), + null, + null, + null), + key -> lookupOldRow ? oldRow : null, + null, + LookupStrategy.from(false, true, false, false), + null, + null); + function.reset(); + if (!lookupOldRow) { + function.add(oldRow); + } + function.add( + new KeyValue() + .replace(row(1), 2, DELETE, GenericRow.of(1, 2, 10, null, null)) + .setLevel(0)); + function.add( + new KeyValue() + .replace(row(1), 2, DELETE, GenericRow.of(1, null, null, 2, 20)) + .setLevel(0)); + + ChangelogResult result = function.getResult(); + assertThat(result.result().valueKind()).isEqualTo(DELETE); + assertThat(result.result().value().getInt(1)).isEqualTo(2); + assertThat(result.changelogs()).hasSize(1); + assertThat(result.changelogs().get(0).valueKind()).isEqualTo(DELETE); + assertThat(result.changelogs().get(0).value()).isEqualTo(oldRow.value()); + } + @ParameterizedTest @ValueSource(booleans = {false, true}) public void testDeduplicate(boolean changelogRowDeduplicate) { diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunctionTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunctionTest.java index 3c7eca4e7281..d8c70797a544 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunctionTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/PartialUpdateMergeFunctionTest.java @@ -32,6 +32,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; import org.junit.jupiter.params.provider.ValueSource; import static org.apache.paimon.CoreOptions.FIELDS_DEFAULT_AGG_FUNC; @@ -1154,6 +1155,174 @@ public void testInitRowWithNullableFieldOnDelete() { validate(func, 1, 2, 2, null); } + @ParameterizedTest + @EnumSource( + value = RowKind.class, + names = {"DELETE", "UPDATE_BEFORE"}) + public void testSequenceGroupDeleteSurvivesSubsequentRetractions(RowKind retractKind) { + MergeFunction function = createSequenceGroupDeleteFunction(); + function.reset(); + add(function, 1, 1, 10, 1, 20); + add(function, RowKind.DELETE, 1, 2, 10, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + + add(function, retractKind, 1, null, null, 2, 20); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + add(function, retractKind, 1, 1, 10, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + add(function, retractKind, 1, null, null, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + assertThat(function.getResult().value().getInt(1)).isEqualTo(2); + } + + @ParameterizedTest + @EnumSource( + value = RowKind.class, + names = {"INSERT", "UPDATE_AFTER"}) + public void testSequenceGroupDeleteOnlyRevivedByAcceptedGroupUpdate(RowKind updateKind) { + MergeFunction function = createSequenceGroupDeleteFunction(); + function.reset(); + add(function, 1, 1, 10, 1, 20); + add(function, RowKind.DELETE, 1, 2, 10, null, null); + + add(function, updateKind, 1, null, null, 3, 30); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + add(function, updateKind, 1, 1, 99, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + assertThat(function.getResult().value().getInt(1)).isEqualTo(2); + + add(function, updateKind, 1, 3, 40, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.INSERT); + validate(function, 1, 3, 40, 3, 30); + + add(function, RowKind.DELETE, 1, 2, 10, null, null); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.INSERT); + } + + @Test + public void testSequenceGroupDeleteSurvivesMergeBoundaryAndReset() { + MergeFunction function = createSequenceGroupDeleteFunction(); + function.reset(); + add(function, 1, 1, 10, 1, 20); + add(function, RowKind.DELETE, 1, 2, 10, null, null); + KeyValue result = function.getResult(); + KeyValue deleted = + new KeyValue() + .replace( + result.key(), + result.sequenceNumber(), + result.valueKind(), + result.value()); + + function.reset(); + function.add(deleted); + add(function, 1, null, null, 3, 30); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + assertThat(function.getResult().value().getInt(1)).isEqualTo(2); + + function.reset(); + function.add( + new KeyValue() + .replace( + GenericRow.of(2), + sequence++, + RowKind.INSERT, + GenericRow.of(2, null, null, 1, 20))); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.INSERT); + validate(function, 2, null, null, 1, 20); + } + + @Test + public void testRemoveRecordOnDeleteIgnoresUpdateBeforeUntilInsert() { + Options options = new Options(); + options.set("partial-update.remove-record-on-delete", "true"); + MergeFunction function = + PartialUpdateMergeFunction.factory( + options, + RowType.of(DataTypes.INT(), DataTypes.INT()), + ImmutableList.of("f0")) + .create(); + function.reset(); + add(function, 1, 10); + add(function, RowKind.DELETE, 1, 10); + add(function, RowKind.UPDATE_BEFORE, 1, 10); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.DELETE); + + add(function, RowKind.INSERT, 1, 20); + assertThat(function.getResult().valueKind()).isEqualTo(RowKind.INSERT); + validate(function, 1, 20); + } + + @ParameterizedTest + @ValueSource(strings = {"", "otherValue,id"}) + public void testSequenceGroupDeleteStateWithProjectedCompositeSequence(String selectedFields) { + RowType rowType = + RowType.builder() + .field("id", DataTypes.INT()) + .field("sequence", DataTypes.INT()) + .field("subSequence", DataTypes.INT()) + .field("value", DataTypes.INT()) + .field("otherSequence", DataTypes.INT()) + .field("otherValue", DataTypes.INT()) + .build(); + Options options = new Options(); + options.set("fields.sequence,subSequence.sequence-group", "value"); + options.set("fields.otherSequence.sequence-group", "otherValue"); + options.set("partial-update.remove-record-on-sequence-group", "sequence"); + MergeFunctionFactory factory = + PartialUpdateMergeFunction.factory(options, rowType, ImmutableList.of("id")); + RowType readType = + factory.adjustReadType( + rowType.project( + selectedFields.isEmpty() + ? new String[0] + : selectedFields.split(","))); + MergeFunction function = factory.create(readType); + GenericRow[] records = { + GenericRow.of(1, 1, 1, 10, 1, 20), + GenericRow.of(1, 1, 2, 10, null, null), + GenericRow.of(1, null, null, null, 2, 20), + GenericRow.of(1, 1, 1, 99, null, null), + GenericRow.of(1, 1, 2, 30, null, null) + }; + RowKind[] kinds = { + RowKind.INSERT, RowKind.DELETE, RowKind.DELETE, RowKind.INSERT, RowKind.UPDATE_AFTER + }; + RowKind[] expected = { + RowKind.INSERT, RowKind.DELETE, RowKind.DELETE, RowKind.DELETE, RowKind.INSERT + }; + function.reset(); + for (int index = 0; index < records.length; index++) { + function.add( + new KeyValue() + .replace( + GenericRow.of(1), + sequence++, + kinds[index], + ProjectedRow.from(readType, rowType) + .replaceRow(records[index]))); + assertThat(function.getResult().valueKind()) + .as("Projection %s, input %s", selectedFields, index) + .isEqualTo(expected[index]); + } + } + + private MergeFunction createSequenceGroupDeleteFunction() { + Options options = new Options(); + options.set("fields.f1.sequence-group", "f2"); + options.set("fields.f3.sequence-group", "f4"); + options.set("partial-update.remove-record-on-sequence-group", "f1"); + RowType rowType = + RowType.of( + DataTypes.INT().notNull(), + DataTypes.INT(), + DataTypes.INT(), + DataTypes.INT(), + DataTypes.INT()); + return PartialUpdateMergeFunction.factory(options, rowType, ImmutableList.of("f0")) + .create(); + } + private void assertProjectedDelete( MergeFunction function, ProjectedRow projection,