diff --git a/docs/docs/primary-key-table/changelog-producer.md b/docs/docs/primary-key-table/changelog-producer.md
index 11260cbf5d80..f6f276581313 100644
--- a/docs/docs/primary-key-table/changelog-producer.md
+++ b/docs/docs/primary-key-table/changelog-producer.md
@@ -103,6 +103,24 @@ changelog for the same record. It also supports `changelog-producer.ignore-updat
records and `changelog-producer.ignore-delete` to exclude DELETE (-D) records from changelog files. These options are
useful when downstream consumers only need the latest state (e.g. upsert sinks) and do not require retraction.
+By setting `'changelog-producer.preserve-sequence-on-retract'` to a comma-separated list of column names,
+retraction records (`-U`, `-D`) will take those columns' values from the incoming event instead of the
+stored row. This is useful when delete or update events carry an event timestamp that downstream consumers
+need, such as external systems like Cassandra that rely on `WRITETIME` for conflict resolution.
+This option is only supported by the `lookup` changelog producer.
+
+```sql
+CREATE TABLE my_table (
+ id INT PRIMARY KEY NOT ENFORCED,
+ data STRING,
+ event_ts BIGINT
+) WITH (
+ 'changelog-producer' = 'lookup',
+ 'sequence.field' = 'event_ts',
+ 'changelog-producer.preserve-sequence-on-retract' = 'event_ts'
+);
+```
+
(Note: Please increase `'execution.checkpointing.max-concurrent-checkpoints'` Flink configuration, this is very
important for performance).
diff --git a/docs/generated/core_configuration.html b/docs/generated/core_configuration.html
index 6a4955280dfb..133dc105a925 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -224,6 +224,12 @@
Boolean |
Whether to ignore update-before records in the changelog. When set to true, UPDATE_BEFORE (-U) records will not be written to changelog files. This configuration is only valid for the changelog-producer is lookup or full-compaction. |
+
+ changelog-producer.preserve-sequence-on-retract |
+ (none) |
+ String |
+ A comma-separated list of column names whose values should be taken from the incoming event rather than the stored row when producing changelog retraction records (-U, -D). This is useful when delete or update events carry their own event timestamp and you want that timestamp preserved in the changelog. Only valid when changelog-producer is lookup. |
+
changelog-producer.row-deduplicate |
false |
@@ -392,12 +398,6 @@
MemorySize |
When incremental size is bigger than this threshold, force a full compaction. |
-
- continuous-compaction.initial-scan-mode |
- earliest |
- Enum |
- Initial snapshot mode for dedicated streaming compaction. When set to 'earliest' (the default), compaction starts from the earliest available snapshot if no COMPACT snapshot exists; when a COMPACT snapshot exists, compaction always resumes from the snapshot after it. When set to 'latest', the latest snapshot is read in ALL mode as the initial baseline and subsequent scans start from the next snapshot. The 'latest' mode skips historical snapshot changes and should only be used when historical changelog replay is not required.
Possible values:- "earliest": Read snapshots from the earliest available snapshot.
- "latest": Read the latest snapshot as the initial full baseline.
|
-
compaction.max-size-amplification-percent |
200 |
@@ -488,6 +488,12 @@
Enum |
Specify the consumer consistency mode for table.
Possible values:- "exactly-once": Readers consume data at snapshot granularity, and strictly ensure that the snapshot-id recorded in the consumer is the snapshot-id + 1 that all readers have exactly consumed.
- "at-least-once": Each reader consumes snapshots at a different rate, and the snapshot with the slowest consumption progress among all readers will be recorded in the consumer.
|
+
+ continuous-compaction.initial-scan-mode |
+ earliest |
+ Enum |
+ Initial snapshot mode for dedicated streaming compaction. When set to 'earliest' (the default), compaction starts from the earliest available snapshot if no COMPACT snapshot exists; when a COMPACT snapshot exists, compaction always resumes from the snapshot after it. When set to 'latest', the latest snapshot is read in ALL mode as the initial baseline and subsequent scans start from the next snapshot. The 'latest' mode skips historical snapshot changes and should only be used when historical changelog replay is not required.
Possible values:- "earliest": Read snapshots from the earliest available snapshot.
- "latest": Read the latest snapshot as the initial full baseline.
|
+
continuous.discovery-interval |
10 s |
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index 298126f3c66d..9f75c23d7556 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -1116,6 +1116,18 @@ public InlineElement getDescription() {
.withDescription(
"Fields that are ignored for comparison while generating -U, +U changelog for the same record. This configuration is only valid for the changelog-producer.row-deduplicate is true.");
+ public static final ConfigOption CHANGELOG_PRODUCER_PRESERVE_SEQUENCE_ON_RETRACT =
+ key("changelog-producer.preserve-sequence-on-retract")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "A comma-separated list of column names whose values should be taken from the "
+ + "incoming event rather than the stored row when producing changelog "
+ + "retraction records (-U, -D). This is useful when delete or update "
+ + "events carry their own event timestamp and you want that timestamp "
+ + "preserved in the changelog. "
+ + "Only valid when changelog-producer is lookup.");
+
public static final ConfigOption TABLE_READ_SEQUENCE_NUMBER_ENABLED =
key("table-read.sequence-number.enabled")
.booleanType()
@@ -3914,6 +3926,16 @@ public List changelogRowDeduplicateIgnoreFields() {
.orElse(Collections.emptyList());
}
+ public List changelogPreserveSequenceOnRetract() {
+ return options.getOptional(CHANGELOG_PRODUCER_PRESERVE_SEQUENCE_ON_RETRACT)
+ .map(
+ s ->
+ Arrays.stream(s.split(","))
+ .map(String::trim)
+ .collect(Collectors.toList()))
+ .orElse(Collections.emptyList());
+ }
+
public boolean tableReadSequenceNumberEnabled() {
return options.get(TABLE_READ_SEQUENCE_NUMBER_ENABLED);
}
diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
index 7283a3030d01..b67f321015b4 100644
--- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
+++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupChangelogMergeFunctionWrapper.java
@@ -20,7 +20,15 @@
import org.apache.paimon.KeyValue;
import org.apache.paimon.codegen.RecordEqualiser;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.Blob;
+import org.apache.paimon.data.Decimal;
+import org.apache.paimon.data.InternalArray;
+import org.apache.paimon.data.InternalMap;
import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.data.InternalVector;
+import org.apache.paimon.data.Timestamp;
+import org.apache.paimon.data.variant.Variant;
import org.apache.paimon.deletionvectors.BucketedDvMaintainer;
import org.apache.paimon.lookup.LookupStrategy;
import org.apache.paimon.mergetree.lookup.FilePosition;
@@ -64,6 +72,7 @@ public class LookupChangelogMergeFunctionWrapper
private final LookupStrategy lookupStrategy;
private final @Nullable BucketedDvMaintainer deletionVectorsMaintainer;
private final Comparator comparator;
+ @Nullable private final SequenceFieldOverwriteRow reusedOverwriteRow;
public LookupChangelogMergeFunctionWrapper(
MergeFunctionFactory mergeFunctionFactory,
@@ -72,6 +81,24 @@ public LookupChangelogMergeFunctionWrapper(
LookupStrategy lookupStrategy,
@Nullable BucketedDvMaintainer deletionVectorsMaintainer,
@Nullable UserDefinedSeqComparator userDefinedSeqComparator) {
+ this(
+ mergeFunctionFactory,
+ lookup,
+ valueEqualiser,
+ lookupStrategy,
+ deletionVectorsMaintainer,
+ userDefinedSeqComparator,
+ null);
+ }
+
+ public LookupChangelogMergeFunctionWrapper(
+ MergeFunctionFactory mergeFunctionFactory,
+ Function lookup,
+ @Nullable RecordEqualiser valueEqualiser,
+ LookupStrategy lookupStrategy,
+ @Nullable BucketedDvMaintainer deletionVectorsMaintainer,
+ @Nullable UserDefinedSeqComparator userDefinedSeqComparator,
+ @Nullable int[] preserveFieldIndices) {
MergeFunction mergeFunction = mergeFunctionFactory.create();
checkArgument(
mergeFunction instanceof LookupMergeFunction,
@@ -88,6 +115,10 @@ public LookupChangelogMergeFunctionWrapper(
this.lookupStrategy = lookupStrategy;
this.deletionVectorsMaintainer = deletionVectorsMaintainer;
this.comparator = createSequenceComparator(userDefinedSeqComparator);
+ this.reusedOverwriteRow =
+ preserveFieldIndices != null && preserveFieldIndices.length > 0
+ ? new SequenceFieldOverwriteRow(preserveFieldIndices)
+ : null;
}
@Override
@@ -152,18 +183,27 @@ private void setChangelog(@Nullable KeyValue before, KeyValue after) {
}
} else {
if (!after.isAdd()) {
- reusedResult.addChangelog(replaceBefore(RowKind.DELETE, before));
+ reusedResult.addChangelog(
+ replaceBeforeWithSequenceOverwrite(RowKind.DELETE, before, after));
} else if (valueEqualiser == null
|| !valueEqualiser.equals(before.value(), after.value())) {
reusedResult
- .addChangelog(replaceBefore(RowKind.UPDATE_BEFORE, before))
+ .addChangelog(
+ replaceBeforeWithSequenceOverwrite(
+ RowKind.UPDATE_BEFORE, before, after))
.addChangelog(replaceAfter(RowKind.UPDATE_AFTER, after));
}
}
}
- private KeyValue replaceBefore(RowKind valueKind, KeyValue from) {
- return replace(reusedBefore, valueKind, from);
+ private KeyValue replaceBeforeWithSequenceOverwrite(
+ RowKind valueKind, KeyValue before, KeyValue after) {
+ if (reusedOverwriteRow != null) {
+ reusedOverwriteRow.replace(before.value(), after.value());
+ return reusedBefore.replace(
+ before.key(), after.sequenceNumber(), valueKind, reusedOverwriteRow);
+ }
+ return replace(reusedBefore, valueKind, before);
}
private KeyValue replaceAfter(RowKind valueKind, KeyValue from) {
@@ -188,4 +228,143 @@ private Comparator createSequenceComparator(
return Long.compare(o1.sequenceNumber(), o2.sequenceNumber());
};
}
+
+ /**
+ * An {@link InternalRow} that delegates to a primary row for all fields, except for specified
+ * sequence field positions which are read from a secondary (event) row. This allows changelog
+ * before-image records to carry the incoming event's sequence field value while preserving the
+ * rest of the stored row's data.
+ */
+ static class SequenceFieldOverwriteRow implements InternalRow {
+
+ private final boolean[] isSequenceField;
+ private InternalRow primaryRow;
+ private InternalRow eventRow;
+
+ SequenceFieldOverwriteRow(int[] sequenceFieldIndices) {
+ int maxIndex = 0;
+ for (int idx : sequenceFieldIndices) {
+ maxIndex = Math.max(maxIndex, idx);
+ }
+ this.isSequenceField = new boolean[maxIndex + 1];
+ for (int idx : sequenceFieldIndices) {
+ this.isSequenceField[idx] = true;
+ }
+ }
+
+ SequenceFieldOverwriteRow replace(InternalRow primaryRow, InternalRow eventRow) {
+ this.primaryRow = primaryRow;
+ this.eventRow = eventRow;
+ return this;
+ }
+
+ private InternalRow rowFor(int pos) {
+ return pos < isSequenceField.length && isSequenceField[pos] ? eventRow : primaryRow;
+ }
+
+ @Override
+ public int getFieldCount() {
+ return primaryRow.getFieldCount();
+ }
+
+ @Override
+ public RowKind getRowKind() {
+ return primaryRow.getRowKind();
+ }
+
+ @Override
+ public void setRowKind(RowKind kind) {
+ primaryRow.setRowKind(kind);
+ }
+
+ @Override
+ public boolean isNullAt(int pos) {
+ return rowFor(pos).isNullAt(pos);
+ }
+
+ @Override
+ public boolean getBoolean(int pos) {
+ return rowFor(pos).getBoolean(pos);
+ }
+
+ @Override
+ public byte getByte(int pos) {
+ return rowFor(pos).getByte(pos);
+ }
+
+ @Override
+ public short getShort(int pos) {
+ return rowFor(pos).getShort(pos);
+ }
+
+ @Override
+ public int getInt(int pos) {
+ return rowFor(pos).getInt(pos);
+ }
+
+ @Override
+ public long getLong(int pos) {
+ return rowFor(pos).getLong(pos);
+ }
+
+ @Override
+ public float getFloat(int pos) {
+ return rowFor(pos).getFloat(pos);
+ }
+
+ @Override
+ public double getDouble(int pos) {
+ return rowFor(pos).getDouble(pos);
+ }
+
+ @Override
+ public BinaryString getString(int pos) {
+ return rowFor(pos).getString(pos);
+ }
+
+ @Override
+ public Decimal getDecimal(int pos, int precision, int scale) {
+ return rowFor(pos).getDecimal(pos, precision, scale);
+ }
+
+ @Override
+ public Timestamp getTimestamp(int pos, int precision) {
+ return rowFor(pos).getTimestamp(pos, precision);
+ }
+
+ @Override
+ public byte[] getBinary(int pos) {
+ return rowFor(pos).getBinary(pos);
+ }
+
+ @Override
+ public Variant getVariant(int pos) {
+ return rowFor(pos).getVariant(pos);
+ }
+
+ @Override
+ public Blob getBlob(int pos) {
+ return rowFor(pos).getBlob(pos);
+ }
+
+ @Override
+ public InternalArray getArray(int pos) {
+ return rowFor(pos).getArray(pos);
+ }
+
+ @Override
+ public InternalVector getVector(int pos) {
+ return rowFor(pos).getVector(pos);
+ }
+
+ @Override
+ public InternalMap getMap(int pos) {
+ return rowFor(pos).getMap(pos);
+ }
+
+ @Override
+ public InternalRow getRow(int pos, int numFields) {
+ return rowFor(pos).getRow(pos, numFields);
+ }
+ }
}
diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
index 36cc4cd6a4d5..0750ded291d0 100644
--- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
+++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/LookupMergeTreeCompactRewriter.java
@@ -199,14 +199,17 @@ public static class LookupMergeFunctionWrapperFactory
@Nullable private final RecordEqualiser valueEqualiser;
private final LookupStrategy lookupStrategy;
@Nullable private final UserDefinedSeqComparator userDefinedSeqComparator;
+ @Nullable private final int[] preserveFieldIndices;
public LookupMergeFunctionWrapperFactory(
@Nullable RecordEqualiser valueEqualiser,
LookupStrategy lookupStrategy,
- @Nullable UserDefinedSeqComparator userDefinedSeqComparator) {
+ @Nullable UserDefinedSeqComparator userDefinedSeqComparator,
+ @Nullable int[] preserveFieldIndices) {
this.valueEqualiser = valueEqualiser;
this.lookupStrategy = lookupStrategy;
this.userDefinedSeqComparator = userDefinedSeqComparator;
+ this.preserveFieldIndices = preserveFieldIndices;
}
@Override
@@ -227,7 +230,8 @@ public MergeFunctionWrapper create(
valueEqualiser,
lookupStrategy,
deletionVectorsMaintainer,
- userDefinedSeqComparator);
+ userDefinedSeqComparator,
+ preserveFieldIndices);
}
}
diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
index 0ef421b0cbbf..863861c4614c 100644
--- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
+++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/MergeTreeCompactManagerFactory.java
@@ -323,11 +323,36 @@ private MergeTreeCompactRewriter createRewriter(
} else {
processorFactory = PersistValueProcessor.factory(valueType);
}
+ List preserveColumns = options.changelogPreserveSequenceOnRetract();
+ int[] preserveFieldIndices = null;
+ if (!preserveColumns.isEmpty()) {
+ List fieldNames = valueType.getFieldNames();
+ preserveFieldIndices =
+ preserveColumns.stream()
+ .mapToInt(
+ name -> {
+ int idx = fieldNames.indexOf(name);
+ if (idx < 0) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Column '%s' specified in '%s' not found in value type. "
+ + "Available columns: %s",
+ name,
+ CoreOptions
+ .CHANGELOG_PRODUCER_PRESERVE_SEQUENCE_ON_RETRACT
+ .key(),
+ fieldNames));
+ }
+ return idx;
+ })
+ .toArray();
+ }
wrapperFactory =
new LookupMergeFunctionWrapperFactory<>(
logDedupEqualSupplier.get(),
lookupStrategy,
- UserDefinedSeqComparator.create(valueType, options));
+ UserDefinedSeqComparator.create(valueType, options),
+ preserveFieldIndices);
}
LookupLevels> lookupLevels =
createLookupLevels(
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 57d99557ca5f..436a47fe37fd 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
@@ -556,4 +556,139 @@ public void testKeepLowestHighLevel() {
kv = result.result();
assertThat(kv.value().getInt(0)).isEqualTo(3);
}
+
+ @Test
+ public void testPreserveSequenceOnRetractDelete() {
+ // Schema: value has two fields: f0 (data), f1 (event_ts to preserve)
+ Map highLevel = new HashMap<>();
+ RowType valueType =
+ RowType.builder()
+ .fields(
+ new DataType[] {DataTypes.INT(), DataTypes.INT()},
+ new String[] {"f0", "f1"})
+ .build();
+ UserDefinedSeqComparator userDefinedSeqComparator =
+ UserDefinedSeqComparator.create(
+ valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1")));
+ assert userDefinedSeqComparator != null;
+
+ // preserve f1 (index 1) on retraction records
+ LookupChangelogMergeFunctionWrapper function =
+ new LookupChangelogMergeFunctionWrapper(
+ LookupMergeFunction.wrap(
+ DeduplicateMergeFunction.factory(), null, null, null),
+ highLevel::get,
+ null,
+ LookupStrategy.from(false, true, false, false),
+ null,
+ userDefinedSeqComparator,
+ new int[] {1});
+
+ // Delete with higher sequence field: changelog -D should carry the delete event's
+ // f1=100, not the old row's f1=50
+ highLevel.put(row(1), new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(2));
+ function.reset();
+ function.add(new KeyValue().replace(row(1), 2, DELETE, row(10, 100)).setLevel(0));
+ ChangelogResult result = function.getResult();
+ assertThat(result).isNotNull();
+ List changelogs = result.changelogs();
+ assertThat(changelogs).hasSize(1);
+ assertThat(changelogs.get(0).valueKind()).isEqualTo(DELETE);
+ // f0 should come from the old row (before image data)
+ assertThat(changelogs.get(0).value().getInt(0)).isEqualTo(10);
+ // f1 (preserved column) should come from the delete event
+ assertThat(changelogs.get(0).value().getInt(1)).isEqualTo(100);
+ // system sequence number should come from the delete event
+ assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(2);
+ }
+
+ @Test
+ public void testPreserveSequenceOnRetractUpdate() {
+ // Schema: value has two fields: f0 (data), f1 (event_ts to preserve)
+ Map highLevel = new HashMap<>();
+ RowType valueType =
+ RowType.builder()
+ .fields(
+ new DataType[] {DataTypes.INT(), DataTypes.INT()},
+ new String[] {"f0", "f1"})
+ .build();
+ UserDefinedSeqComparator userDefinedSeqComparator =
+ UserDefinedSeqComparator.create(
+ valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1")));
+ assert userDefinedSeqComparator != null;
+
+ // preserve f1 (index 1) on retraction records
+ LookupChangelogMergeFunctionWrapper function =
+ new LookupChangelogMergeFunctionWrapper(
+ LookupMergeFunction.wrap(
+ DeduplicateMergeFunction.factory(), null, null, null),
+ highLevel::get,
+ null,
+ LookupStrategy.from(false, true, false, false),
+ null,
+ userDefinedSeqComparator,
+ new int[] {1});
+
+ // Update: -U changelog should carry the new event's f1=100, not the old row's f1=50
+ function.reset();
+ function.add(new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(1));
+ function.add(new KeyValue().replace(row(1), 2, INSERT, row(20, 100)).setLevel(0));
+ ChangelogResult result = function.getResult();
+ assertThat(result).isNotNull();
+ List changelogs = result.changelogs();
+ assertThat(changelogs).hasSize(2);
+
+ // -U (UPDATE_BEFORE): f0 from old row, f1 from new event
+ assertThat(changelogs.get(0).valueKind()).isEqualTo(UPDATE_BEFORE);
+ assertThat(changelogs.get(0).value().getInt(0)).isEqualTo(10);
+ assertThat(changelogs.get(0).value().getInt(1)).isEqualTo(100);
+ assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(2);
+
+ // +U (UPDATE_AFTER): entirely from new event
+ assertThat(changelogs.get(1).valueKind()).isEqualTo(UPDATE_AFTER);
+ assertThat(changelogs.get(1).value().getInt(0)).isEqualTo(20);
+ assertThat(changelogs.get(1).value().getInt(1)).isEqualTo(100);
+ }
+
+ @Test
+ public void testPreserveSequenceOnRetractNotConfigured() {
+ // Verify that the old behavior is preserved when no columns are specified
+ Map highLevel = new HashMap<>();
+ RowType valueType =
+ RowType.builder()
+ .fields(
+ new DataType[] {DataTypes.INT(), DataTypes.INT()},
+ new String[] {"f0", "f1"})
+ .build();
+ UserDefinedSeqComparator userDefinedSeqComparator =
+ UserDefinedSeqComparator.create(
+ valueType, CoreOptions.fromMap(ImmutableMap.of("sequence.field", "f1")));
+ assert userDefinedSeqComparator != null;
+
+ // no preserve columns (null)
+ LookupChangelogMergeFunctionWrapper function =
+ new LookupChangelogMergeFunctionWrapper(
+ LookupMergeFunction.wrap(
+ DeduplicateMergeFunction.factory(), null, null, null),
+ highLevel::get,
+ null,
+ LookupStrategy.from(false, true, false, false),
+ null,
+ userDefinedSeqComparator,
+ null);
+
+ // Delete: changelog -D should use old row's values (original behavior)
+ highLevel.put(row(1), new KeyValue().replace(row(1), 1, INSERT, row(10, 50)).setLevel(2));
+ function.reset();
+ function.add(new KeyValue().replace(row(1), 2, DELETE, row(10, 100)).setLevel(0));
+ ChangelogResult result = function.getResult();
+ assertThat(result).isNotNull();
+ List changelogs = result.changelogs();
+ assertThat(changelogs).hasSize(1);
+ assertThat(changelogs.get(0).valueKind()).isEqualTo(DELETE);
+ assertThat(changelogs.get(0).value().getInt(0)).isEqualTo(10);
+ // f1 should be from the OLD row (original behavior)
+ assertThat(changelogs.get(0).value().getInt(1)).isEqualTo(50);
+ assertThat(changelogs.get(0).sequenceNumber()).isEqualTo(1);
+ }
}