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
5 changes: 5 additions & 0 deletions docs/docs/primary-key-table/merge-engine/partial-update.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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) {
Expand Down Expand Up @@ -193,6 +193,7 @@ public void add(KeyValue kv) {

latestSequenceNumber = kv.sequenceNumber();
if (fieldSeqComparators.isEmpty()) {
currentDeleteRow = false;
updateNonNullFields(kv);
} else {
updateWithSequenceGroup(kv);
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,15 @@
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;
import org.apache.paimon.mergetree.compact.aggregate.AggregateMergeFunction;
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;
Expand Down Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<KeyValue> 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<KeyValue> 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<KeyValue> 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<KeyValue> 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<KeyValue> factory =
PartialUpdateMergeFunction.factory(options, rowType, ImmutableList.of("id"));
RowType readType =
factory.adjustReadType(
rowType.project(
selectedFields.isEmpty()
? new String[0]
: selectedFields.split(",")));
MergeFunction<KeyValue> 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<KeyValue> 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<KeyValue> function,
ProjectedRow projection,
Expand Down
Loading