From 5a7eec929e7757f35580b10c608fcf8f30e4ea0d Mon Sep 17 00:00:00 2001 From: litiliu <38579068+litiliu@users.noreply.github.com> Date: Sun, 20 Sep 2026 17:29:06 +0800 Subject: [PATCH] [server] Support update_if_changed merge engine for primary-key tables --- .../fluss/client/table/FlussTableITCase.java | 145 +++++++++ .../apache/fluss/config/ConfigOptions.java | 6 +- .../fluss/metadata/MergeEngineType.java | 25 +- .../fluss/flink/sink/FlinkTableSink.java | 6 +- .../flink/sink/FlinkTableSinkITCase.java | 34 +++ .../fluss/server/kv/KvWriteProcessor.java | 5 + .../fluss/server/kv/rowmerger/RowMerger.java | 5 +- .../rowmerger/UpdateIfChangedRowMerger.java | 194 ++++++++++++ .../kv/rowmerger/RowMergerCreateTest.java | 34 +++ .../UpdateIfChangedRowMergerTest.java | 283 ++++++++++++++++++ website/docs/engine-flink/options.md | 2 +- website/docs/engine-flink/writes.md | 2 +- .../docs/table-design/merge-engines/index.md | 1 + .../merge-engines/update-if-changed.md | 77 +++++ .../docs/table-design/table-types/pk-table.md | 1 + 15 files changed, 808 insertions(+), 12 deletions(-) create mode 100644 fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMerger.java create mode 100644 fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMergerTest.java create mode 100644 website/docs/table-design/merge-engines/update-if-changed.md diff --git a/fluss-client/src/test/java/org/apache/fluss/client/table/FlussTableITCase.java b/fluss-client/src/test/java/org/apache/fluss/client/table/FlussTableITCase.java index 9dd5ee45345..f275361e04d 100644 --- a/fluss-client/src/test/java/org/apache/fluss/client/table/FlussTableITCase.java +++ b/fluss-client/src/test/java/org/apache/fluss/client/table/FlussTableITCase.java @@ -36,6 +36,7 @@ import org.apache.fluss.config.MemorySize; import org.apache.fluss.fs.FsPath; import org.apache.fluss.fs.TestFileSystem; +import org.apache.fluss.metadata.ChangelogImage; import org.apache.fluss.metadata.DataLakeFormat; import org.apache.fluss.metadata.KvFormat; import org.apache.fluss.metadata.LogFormat; @@ -63,6 +64,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; +import org.junit.jupiter.params.provider.EnumSource; import org.junit.jupiter.params.provider.ValueSource; import javax.annotation.Nullable; @@ -1561,6 +1563,149 @@ void testFirstRowMergeEngine(boolean doProjection) throws Exception { } } + @ParameterizedTest + @EnumSource(ChangelogImage.class) + void testUpdateIfChangedMergeEngine(ChangelogImage changelogImage) throws Exception { + TableDescriptor tableDescriptor = + TableDescriptor.builder() + .schema(DATA1_SCHEMA_PK) + .property( + ConfigOptions.TABLE_MERGE_ENGINE, MergeEngineType.UPDATE_IF_CHANGED) + .property(ConfigOptions.TABLE_CHANGELOG_IMAGE, changelogImage) + .build(); + RowType rowType = DATA1_SCHEMA_PK.getRowType(); + TablePath tablePath = + TablePath.of( + "test_db_1", + "test_update_if_changed_merge_engine_" + changelogImage.name()); + createTable(tablePath, tableDescriptor, false); + + try (Table table = conn.getTable(tablePath)) { + UpsertWriter upsertWriter = table.newUpsert().createWriter(); + // insert a row + upsertWriter.upsert(row(0, "v0")); + // value-identical upserts: should be no-ops that emit no changelog + upsertWriter.upsert(row(0, "v0")); + upsertWriter.upsert(row(0, "v0")); + // a field differs: should emit an update changelog + upsertWriter.upsert(row(0, "v1")); + // delete the row: should emit a delete changelog + upsertWriter.delete(row(0, "v1")); + upsertWriter.flush(); + + // No records should be emitted for the value-identical upserts. + List expected = new ArrayList<>(); + expected.add(new ScanRecord(-1, -1, ChangeType.INSERT, row(0, "v0"))); + if (changelogImage == ChangelogImage.FULL) { + expected.add(new ScanRecord(-1, -1, ChangeType.UPDATE_BEFORE, row(0, "v0"))); + } + expected.add(new ScanRecord(-1, -1, ChangeType.UPDATE_AFTER, row(0, "v1"))); + expected.add(new ScanRecord(-1, -1, ChangeType.DELETE, row(0, "v1"))); + + LogScanner logScanner = table.newScan().createLogScanner(); + logScanner.subscribeFromBeginning(0); + List actualLogRecords = new ArrayList<>(expected.size()); + while (actualLogRecords.size() < expected.size()) { + ScanRecords scanRecords = logScanner.poll(Duration.ofSeconds(1)); + scanRecords.forEach(actualLogRecords::add); + } + assertThat(logScanner.poll(Duration.ofSeconds(1))).isEmpty(); + logScanner.close(); + + assertThat(actualLogRecords).hasSize(expected.size()); + for (int i = 0; i < actualLogRecords.size(); i++) { + ScanRecord actual = actualLogRecords.get(i); + assertThat(actual.getChangeType()).isEqualTo(expected.get(i).getChangeType()); + assertThatRow(actual.getRow()) + .withSchema(rowType) + .isEqualTo(expected.get(i).getRow()); + } + } + } + + @ParameterizedTest + @EnumSource(ChangelogImage.class) + void testUpdateIfChangedMergeEngineWithPartialUpdate(ChangelogImage changelogImage) + throws Exception { + Schema schema = + Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("name", DataTypes.STRING()) + .column("data", DataTypes.STRING()) + .primaryKey("id") + .build(); + TableDescriptor tableDescriptor = + TableDescriptor.builder() + .schema(schema) + .distributedBy(1, "id") + .property( + ConfigOptions.TABLE_MERGE_ENGINE, MergeEngineType.UPDATE_IF_CHANGED) + .property(ConfigOptions.TABLE_CHANGELOG_IMAGE, changelogImage) + .build(); + RowType rowType = schema.getRowType(); + TablePath tablePath = + TablePath.of( + "test_db_1", + "test_update_if_changed_partial_update_" + changelogImage.name()); + createTable(tablePath, tableDescriptor, false); + + try (Table table = conn.getTable(tablePath)) { + UpsertWriter fullWriter = table.newUpsert().createWriter(); + fullWriter.upsert(row(0, "v0", "kept")).get(); + + UpsertWriter partialWriter = + table.newUpsert().partialUpdate("id", "name").createWriter(); + // unchanged partial upsert: no changelog + partialWriter.upsert(row(0, "v0", null)).get(); + // changed partial upsert: normal update changelog + partialWriter.upsert(row(0, "v1", null)).get(); + // unchanged partial upsert: no changelog + partialWriter.upsert(row(0, "v1", null)).get(); + // changed partial delete: clears name and emits an update changelog + partialWriter.delete(row(0, "v1", null)).get(); + // unchanged partial delete: no changelog + partialWriter.delete(row(0, null, null)).get(); + partialWriter.flush(); + + List expected = new ArrayList<>(); + expected.add(new ScanRecord(-1, -1, ChangeType.INSERT, row(0, "v0", "kept"))); + if (changelogImage == ChangelogImage.FULL) { + expected.add( + new ScanRecord(-1, -1, ChangeType.UPDATE_BEFORE, row(0, "v0", "kept"))); + } + expected.add(new ScanRecord(-1, -1, ChangeType.UPDATE_AFTER, row(0, "v1", "kept"))); + if (changelogImage == ChangelogImage.FULL) { + expected.add( + new ScanRecord(-1, -1, ChangeType.UPDATE_BEFORE, row(0, "v1", "kept"))); + } + expected.add(new ScanRecord(-1, -1, ChangeType.UPDATE_AFTER, row(0, null, "kept"))); + + LogScanner logScanner = table.newScan().createLogScanner(); + logScanner.subscribeFromBeginning(0); + List actualLogRecords = new ArrayList<>(expected.size()); + while (actualLogRecords.size() < expected.size()) { + ScanRecords scanRecords = logScanner.poll(Duration.ofSeconds(1)); + scanRecords.forEach(actualLogRecords::add); + } + assertThat(logScanner.poll(Duration.ofSeconds(1))).isEmpty(); + logScanner.close(); + + assertThat(actualLogRecords).hasSize(expected.size()); + for (int i = 0; i < actualLogRecords.size(); i++) { + ScanRecord actual = actualLogRecords.get(i); + assertThat(actual.getChangeType()).isEqualTo(expected.get(i).getChangeType()); + assertThatRow(actual.getRow()) + .withSchema(rowType) + .isEqualTo(expected.get(i).getRow()); + } + + Lookuper lookuper = table.newLookup().createLookuper(); + assertThatRow(lookuper.lookup(row(0)).get().getSingletonRow()) + .withSchema(rowType) + .isEqualTo(row(0, null, "kept")); + } + } + @ParameterizedTest @CsvSource({"none,3", "lz4_frame,3", "zstd,3", "zstd,9"}) void testArrowCompressionAndProject(String compression, String level) throws Exception { diff --git a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java index d7f0afd7967..8f5f75b91b3 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java @@ -2010,10 +2010,12 @@ public class ConfigOptions { .noDefaultValue() .withDescription( "Defines the merge engine for the primary key table. By default, primary key table doesn't have merge engine. " - + "The supported merge engines are `first_row`, `versioned`, and `aggregation`. " + + "The supported merge engines are `first_row`, `versioned`, `aggregation`, and `update_if_changed`. " + "The `first_row` merge engine will keep the first row of the same primary key. " + "The `versioned` merge engine will keep the row with the largest version of the same primary key. " - + "The `aggregation` merge engine will aggregate rows with the same primary key using field-level aggregate functions."); + + "The `aggregation` merge engine will aggregate rows with the same primary key using field-level aggregate functions. " + + "The `update_if_changed` merge engine keeps last-row upsert semantics but suppresses value-identical writes: " + + "when an incoming row is logically equal to the stored row, the write is a no-op and no changelog is emitted."); public static final ConfigOption TABLE_MERGE_ENGINE_VERSION_COLUMN = // we may need to introduce "del-column" in the future to support delete operation diff --git a/fluss-common/src/main/java/org/apache/fluss/metadata/MergeEngineType.java b/fluss-common/src/main/java/org/apache/fluss/metadata/MergeEngineType.java index f4f94d5a744..dd6369fce0f 100644 --- a/fluss-common/src/main/java/org/apache/fluss/metadata/MergeEngineType.java +++ b/fluss-common/src/main/java/org/apache/fluss/metadata/MergeEngineType.java @@ -22,10 +22,10 @@ * *

A primary key table with a merge engine is a special kind of table, called "merge table". * Fluss provides 3 kinds of table: "primary key table", "log table", and "merge table". Merge table - * is a primary key table that has a primary key definition but doesn't directly UPDATE and DELETE - * rows in the table, and instead, it merges the append rows into a new data set according to the - * defined {@link MergeEngineType}. Therefore, it doesn't support direct UPDATE (also - * partial-update) and DELETE operations and only supports INSERT or APPEND operations. + * is a primary key table that merges incoming rows into a new data set according to the defined + * {@link MergeEngineType}. Most merge engines accept only INSERT or APPEND operations, rather than + * direct UPDATE (including partial update) and DELETE operations. {@link #UPDATE_IF_CHANGED} is an + * exception that retains normal primary-key table update and delete semantics. * *

Note: A primary key table doesn't have a merge engine by default. * @@ -61,7 +61,20 @@ public enum MergeEngineType { * * @since 0.9 */ - AGGREGATION; + AGGREGATION, + + /** + * A merge engine that keeps last-row upsert semantics but suppresses value-identical writes. + * When an incoming row is logically equal to the currently stored row (compared by logical + * field values, not raw bytes), the write is a no-op and no changelog is emitted. When at least + * one field differs, the incoming row replaces the stored row and a normal update changelog is + * emitted. Unlike {@link #FIRST_ROW}, legitimate updates are still applied; unlike {@link + * #VERSIONED}, no version column is required. Partial updates and delete operations are + * supported. + * + * @since 1.1 + */ + UPDATE_IF_CHANGED; /** Creates a {@link MergeEngineType} from the given string. */ public static MergeEngineType fromString(String type) { @@ -72,6 +85,8 @@ public static MergeEngineType fromString(String type) { return VERSIONED; case "AGGREGATION": return AGGREGATION; + case "UPDATE_IF_CHANGED": + return UPDATE_IF_CHANGED; default: throw new IllegalArgumentException("Unsupported merge engine type: " + type); } diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/sink/FlinkTableSink.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/sink/FlinkTableSink.java index a0d4aacc946..a6aa69e2d34 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/sink/FlinkTableSink.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/sink/FlinkTableSink.java @@ -156,7 +156,9 @@ public SinkRuntimeProvider getSinkRuntimeProvider(Context context) { "Fluss table sink does not support partial updates for table without primary key. Please make sure the " + "number of specified columns in INSERT INTO matches columns of the Fluss table."); } - if (mergeEngineType != null && mergeEngineType != MergeEngineType.AGGREGATION) { + if (mergeEngineType != null + && mergeEngineType != MergeEngineType.AGGREGATION + && mergeEngineType != MergeEngineType.UPDATE_IF_CHANGED) { throw new ValidationException( String.format( "Table %s uses the '%s' merge engine which does not support partial updates. Please make sure the " @@ -376,7 +378,7 @@ private void validateUpdatable() { "Table %s is a Log Table. Log Table doesn't support DELETE and UPDATE statements.", tablePath)); } - if (mergeEngineType != null) { + if (mergeEngineType != null && mergeEngineType != MergeEngineType.UPDATE_IF_CHANGED) { throw new UnsupportedOperationException( String.format( "Table %s uses the '%s' merge engine which does not support DELETE or UPDATE statements.", diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/FlinkTableSinkITCase.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/FlinkTableSinkITCase.java index 09a59de57ff..00ae0775a32 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/FlinkTableSinkITCase.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/FlinkTableSinkITCase.java @@ -1319,6 +1319,40 @@ void testUnsupportedStmtOnVersionMergeEngine() { tablePath); } + @Test + void testUpdateIfChangedMergeEngineSupportsMutations() throws Exception { + String tableName = "updateIfChangedMergeEngineTable"; + tBatchEnv.executeSql( + String.format( + "create table %s (" + + " a int not null," + + " b bigint null, " + + " c string null, " + + " primary key (a) not enforced" + + ") with ('table.merge-engine' = 'update_if_changed')", + tableName)); + + tBatchEnv + .executeSql( + String.format( + "INSERT INTO %s VALUES (1, 10, 'initial'), (2, 20, 'delete')", + tableName)) + .await(); + + // Verify that the Flink sink accepts partial updates for this merge engine. + tBatchEnv + .executeSql(String.format("INSERT INTO %s (a, c) VALUES (1, 'partial')", tableName)) + .await(); + + // Verify that row-level UPDATE and DELETE statements are accepted as well. + tBatchEnv.executeSql(String.format("UPDATE %s SET b = 11 WHERE a = 1", tableName)).await(); + tBatchEnv.executeSql(String.format("DELETE FROM %s WHERE a = 2", tableName)).await(); + + CloseableIterator rowIter = + tBatchEnv.executeSql(String.format("SELECT * FROM %s", tableName)).collect(); + assertResultsIgnoreOrder(rowIter, Collections.singletonList("+I[1, 11, partial]"), true); + } + @Test void testVersionMergeEngineWithTypeBigint() throws Exception { tEnv.executeSql( diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvWriteProcessor.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvWriteProcessor.java index f67bfaa627b..4411c8e4da6 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/KvWriteProcessor.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/KvWriteProcessor.java @@ -374,6 +374,11 @@ private long processDeletion( BinaryValue newValue = currentMerger.delete(oldValue); + if (newValue == oldValue) { + // no actual change, skip this record + return logOffset; + } + // if newValue is null, it means the row should be deleted if (newValue == null) { return applyDelete( diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/RowMerger.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/RowMerger.java index 3f2f5708865..020941b5bfe 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/RowMerger.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/RowMerger.java @@ -50,7 +50,8 @@ public interface RowMerger { * DeleteBehavior#ALLOW}. * * @param oldRow the old row. - * @return the merged row, or null if the row is deleted. + * @return the merged row, or null if the row is deleted. Returning the same instance as {@code + * oldRow} means that nothing happens to the row. */ @Nullable BinaryValue delete(BinaryValue oldRow); @@ -100,6 +101,8 @@ static RowMerger create(TableConfig tableConf, KvFormat kvFormat, SchemaGetter s return new VersionedRowMerger(versionColumn.get(), deleteBehavior); case AGGREGATION: return new AggregateRowMerger(tableConf, kvFormat, schemaGetter); + case UPDATE_IF_CHANGED: + return new UpdateIfChangedRowMerger(kvFormat, schemaGetter, deleteBehavior); default: throw new IllegalArgumentException( "Unsupported merge engine type: " + mergeEngineType.get()); diff --git a/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMerger.java b/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMerger.java new file mode 100644 index 00000000000..f825708c7e9 --- /dev/null +++ b/fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMerger.java @@ -0,0 +1,194 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.kv.rowmerger; + +import org.apache.fluss.metadata.DeleteBehavior; +import org.apache.fluss.metadata.KvFormat; +import org.apache.fluss.metadata.MergeEngineType; +import org.apache.fluss.metadata.Schema; +import org.apache.fluss.metadata.SchemaGetter; +import org.apache.fluss.record.BinaryValue; +import org.apache.fluss.row.BinaryRow; +import org.apache.fluss.row.InternalRow; +import org.apache.fluss.types.RowType; + +import javax.annotation.Nullable; + +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** + * A merger that keeps last-row upsert semantics but suppresses value-identical writes. + * + *

When the complete candidate row is logically equal to the currently stored row, the merger + * returns the {@code oldValue} instance so that the KV write path treats the write as a no-op (no + * state update, no changelog). When at least one field differs, it returns the candidate row so + * that a normal update changelog is emitted. + * + *

Equality is based on logical field values rather than raw serialized bytes: nulls compare with + * SQL semantics, binary values compare by content, and rows written with an older schema are + * aligned to the latest schema by stable column IDs (missing fields are treated as null) before + * comparison. + * + *

The default merge engine is used as a delegate to preserve its full-row, partial-update, and + * partial-delete semantics. This class intentionally does not extend {@link DefaultRowMerger}, so + * writes still look up the stored value before merging and can perform the equality check. + * + *

This class is not thread-safe: it caches the latest-schema equalizer between {@link + * #configureTargetColumns} and {@link #merge} calls, and is guaranteed to be accessed by a single + * thread at a time (protected by KvTablet's write lock). + * + * @see MergeEngineType#UPDATE_IF_CHANGED + */ +public class UpdateIfChangedRowMerger implements RowMerger { + + private final RowMerger delegate; + private final SchemaGetter schemaGetter; + + private RowEqualizer rowEqualizer; + + public UpdateIfChangedRowMerger( + KvFormat kvFormat, SchemaGetter schemaGetter, @Nullable DeleteBehavior deleteBehavior) { + this(new DefaultRowMerger(kvFormat, deleteBehavior), schemaGetter, null); + } + + private UpdateIfChangedRowMerger( + RowMerger delegate, SchemaGetter schemaGetter, @Nullable RowEqualizer rowEqualizer) { + this.delegate = delegate; + this.schemaGetter = schemaGetter; + this.rowEqualizer = rowEqualizer; + } + + @Nullable + @Override + public BinaryValue merge(@Nullable BinaryValue oldValue, BinaryValue newValue) { + return suppressUnchanged(oldValue, delegate.merge(oldValue, newValue)); + } + + @Nullable + @Override + public BinaryValue delete(BinaryValue oldRow) { + return suppressUnchanged(oldRow, delegate.delete(oldRow)); + } + + @Override + public DeleteBehavior deleteBehavior() { + return delegate.deleteBehavior(); + } + + @Override + public RowMerger configureTargetColumns( + @Nullable int[] targetColumns, short latestSchemaId, Schema latestSchema) { + RowMerger configuredMerger = + delegate.configureTargetColumns(targetColumns, latestSchemaId, latestSchema); + if (rowEqualizer == null || latestSchemaId != rowEqualizer.latestSchemaId) { + this.rowEqualizer = new RowEqualizer(schemaGetter, latestSchemaId, latestSchema); + } + return new UpdateIfChangedRowMerger(configuredMerger, schemaGetter, rowEqualizer); + } + + @Nullable + private BinaryValue suppressUnchanged( + @Nullable BinaryValue oldValue, @Nullable BinaryValue candidate) { + if (oldValue != null && candidate != null && rowEqualizer.equals(oldValue, candidate)) { + // return the old value (same instance) so the write path treats this as a no-op + return oldValue; + } + return candidate; + } + + /** Compares two rows by logical field values aligned to the latest schema. */ + static final class RowEqualizer { + + private final SchemaGetter schemaGetter; + private final short latestSchemaId; + private final Schema latestSchema; + private final List latestColumns; + private final Map alignedFieldGetters; + + RowEqualizer(SchemaGetter schemaGetter, short latestSchemaId, Schema latestSchema) { + this.schemaGetter = schemaGetter; + this.latestSchemaId = latestSchemaId; + this.latestSchema = latestSchema; + this.latestColumns = latestSchema.getColumns(); + this.alignedFieldGetters = new HashMap<>(); + } + + boolean equals(BinaryValue left, BinaryValue right) { + InternalRow.FieldGetter[] leftGetters = getOrCreateAlignedFieldGetters(left.schemaId); + InternalRow.FieldGetter[] rightGetters = getOrCreateAlignedFieldGetters(right.schemaId); + for (int i = 0; i < latestColumns.size(); i++) { + Object leftField = getFieldOrNull(leftGetters[i], left.row); + Object rightField = getFieldOrNull(rightGetters[i], right.row); + if (!fieldEquals(leftField, rightField)) { + return false; + } + } + return true; + } + + private InternalRow.FieldGetter[] getOrCreateAlignedFieldGetters(short schemaId) { + InternalRow.FieldGetter[] fieldGetters = alignedFieldGetters.get(schemaId); + if (fieldGetters != null) { + return fieldGetters; + } + + Schema schema = + schemaId == latestSchemaId ? latestSchema : schemaGetter.getSchema(schemaId); + Map columnIdToIndex = new HashMap<>(); + List columns = schema.getColumns(); + for (int i = 0; i < columns.size(); i++) { + columnIdToIndex.put(columns.get(i).getColumnId(), i); + } + + RowType rowType = schema.getRowType(); + fieldGetters = new InternalRow.FieldGetter[latestColumns.size()]; + for (int i = 0; i < latestColumns.size(); i++) { + Integer sourceIndex = columnIdToIndex.get(latestColumns.get(i).getColumnId()); + if (sourceIndex != null) { + fieldGetters[i] = + InternalRow.createFieldGetter( + rowType.getTypeAt(sourceIndex), sourceIndex); + } + } + alignedFieldGetters.put(schemaId, fieldGetters); + return fieldGetters; + } + + @Nullable + private static Object getFieldOrNull( + @Nullable InternalRow.FieldGetter fieldGetter, BinaryRow row) { + return fieldGetter == null ? null : fieldGetter.getFieldOrNull(row); + } + + private static boolean fieldEquals(@Nullable Object left, @Nullable Object right) { + if (left == right) { + return true; + } + if (left == null || right == null) { + return false; + } + if (left instanceof byte[] && right instanceof byte[]) { + return Arrays.equals((byte[]) left, (byte[]) right); + } + return left.equals(right); + } + } +} diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/RowMergerCreateTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/RowMergerCreateTest.java index e6d7f484995..f98247e96b2 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/RowMergerCreateTest.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/RowMergerCreateTest.java @@ -182,6 +182,40 @@ void testCreateDefaultRowMerger() { assertThat(merger).isInstanceOf(DefaultRowMerger.class); } + @Test + void testCreateUpdateIfChangedRowMerger() { + Schema schema = + Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("value", DataTypes.STRING()) + .primaryKey("id") + .build(); + + Configuration conf = new Configuration(); + conf.setString( + ConfigOptions.TABLE_MERGE_ENGINE.key(), MergeEngineType.UPDATE_IF_CHANGED.name()); + TableConfig tableConfig = new TableConfig(conf); + + RowMerger merger = + RowMerger.create(tableConfig, KvFormat.COMPACTED, createSchemaGetter(schema)); + merger.configureTargetColumns(null, SCHEMA_ID, schema); + + assertThat(merger).isInstanceOf(UpdateIfChangedRowMerger.class); + // delete is supported and defaults to ALLOW + assertThat(merger.deleteBehavior()).isEqualTo(DeleteBehavior.ALLOW); + + // value-identical write is a no-op (returns the old value instance) + BinaryRow oldRow = compactedRow(schema.getRowType(), new Object[] {1, "a"}); + BinaryRow newRow = compactedRow(schema.getRowType(), new Object[] {1, "a"}); + BinaryValue oldValue = toBinaryValue(oldRow); + assertThat(merger.merge(oldValue, toBinaryValue(newRow))).isSameAs(oldValue); + + // changed write returns the new value + BinaryValue changed = + toBinaryValue(compactedRow(schema.getRowType(), new Object[] {1, "b"})); + assertThat(merger.merge(oldValue, changed)).isSameAs(changed); + } + @Test void testCreateAggregateRowMergerWithCompositePrimaryKeyAndMultipleAggTypes() { // Create schema with composite primary key and various aggregation function types diff --git a/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMergerTest.java b/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMergerTest.java new file mode 100644 index 00000000000..571d8fb4c8b --- /dev/null +++ b/fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMergerTest.java @@ -0,0 +1,283 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.fluss.server.kv.rowmerger; + +import org.apache.fluss.metadata.DeleteBehavior; +import org.apache.fluss.metadata.KvFormat; +import org.apache.fluss.metadata.Schema; +import org.apache.fluss.metadata.SchemaInfo; +import org.apache.fluss.record.BinaryValue; +import org.apache.fluss.record.TestingSchemaGetter; +import org.apache.fluss.row.BinaryRow; +import org.apache.fluss.types.DataTypes; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; + +import java.util.Arrays; + +import static org.apache.fluss.testutils.DataTestUtils.compactedRow; +import static org.apache.fluss.testutils.DataTestUtils.indexedRow; +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link UpdateIfChangedRowMerger}. */ +class UpdateIfChangedRowMergerTest { + + private static final short SCHEMA_ID = 1; + private static final short SCHEMA_2_ID = 2; + private static final short SCHEMA_AFTER_DROP_ID = 3; + + private static final Schema SCHEMA = + Schema.newBuilder() + .column("id", DataTypes.INT()) + .column("name", DataTypes.STRING()) + .column("data", DataTypes.BYTES()) + .primaryKey("id") + .build(); + + private static final Schema SCHEMA_2 = + Schema.newBuilder() + .fromColumns( + Arrays.asList( + new Schema.Column("id", DataTypes.INT(), null, (short) 0), + new Schema.Column("name", DataTypes.STRING(), null, (short) 1), + new Schema.Column("data", DataTypes.BYTES(), null, (short) 2), + // add new column at end + new Schema.Column("age", DataTypes.INT(), null, (short) 3))) + .primaryKey("id") + .build(); + + private static final Schema SCHEMA_AFTER_DROP = + Schema.newBuilder() + .fromColumns( + Arrays.asList( + new Schema.Column("id", DataTypes.INT(), null, (short) 0), + new Schema.Column("data", DataTypes.BYTES(), null, (short) 2))) + .primaryKey("id") + .build(); + + private BinaryValue value(int id, String name, byte[] data) { + return value(KvFormat.COMPACTED, id, name, data); + } + + private BinaryValue value(KvFormat kvFormat, int id, String name, byte[] data) { + return binaryValue(SCHEMA_ID, SCHEMA, kvFormat, new Object[] {id, name, data}); + } + + private BinaryValue value2(int id, String name, byte[] data, Integer age) { + return binaryValue( + SCHEMA_2_ID, SCHEMA_2, KvFormat.COMPACTED, new Object[] {id, name, data, age}); + } + + private BinaryValue valueAfterDrop(int id, byte[] data) { + return binaryValue( + SCHEMA_AFTER_DROP_ID, + SCHEMA_AFTER_DROP, + KvFormat.COMPACTED, + new Object[] {id, data}); + } + + private static BinaryValue binaryValue( + short schemaId, Schema schema, KvFormat kvFormat, Object[] fields) { + BinaryRow row = + kvFormat == KvFormat.COMPACTED + ? compactedRow(schema.getRowType(), fields) + : indexedRow(schema.getRowType(), fields); + return new BinaryValue(schemaId, row); + } + + private static UpdateIfChangedRowMerger createMerger( + KvFormat kvFormat, SchemaInfo... schemaInfos) { + TestingSchemaGetter schemaGetter = new TestingSchemaGetter(schemaInfos[0]); + for (int i = 1; i < schemaInfos.length; i++) { + schemaGetter.updateLatestSchemaInfo(schemaInfos[i]); + } + return new UpdateIfChangedRowMerger(kvFormat, schemaGetter, DeleteBehavior.ALLOW); + } + + private static UpdateIfChangedRowMerger createMerger() { + return createMerger(KvFormat.COMPACTED, new SchemaInfo(SCHEMA, SCHEMA_ID)); + } + + @Test + void testInsertWhenNoOldValue() { + UpdateIfChangedRowMerger merger = createMerger(); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + BinaryValue newValue = value(1, "a", new byte[] {1, 2}); + assertThat(merger.merge(null, newValue)).isSameAs(newValue); + } + + @ParameterizedTest + @EnumSource(KvFormat.class) + void testNoOpWhenLogicallyEqual(KvFormat kvFormat) { + UpdateIfChangedRowMerger merger = createMerger(kvFormat, new SchemaInfo(SCHEMA, SCHEMA_ID)); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(kvFormat, 1, "a", new byte[] {1, 2}); + // different instances, same logical content (including binary content) + BinaryValue newValue = value(kvFormat, 1, "a", new byte[] {1, 2}); + + // returns the old value instance so the write path treats it as a no-op + assertThat(merger.merge(oldValue, newValue)).isSameAs(oldValue); + } + + @Test + void testUpdateWhenFieldDiffers() { + UpdateIfChangedRowMerger merger = createMerger(); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {1, 2}); + BinaryValue newValue = value(1, "b", new byte[] {1, 2}); + + assertThat(merger.merge(oldValue, newValue)).isSameAs(newValue); + } + + @Test + void testUpdateWhenBinaryContentDiffers() { + UpdateIfChangedRowMerger merger = createMerger(); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {1, 2}); + BinaryValue newValue = value(1, "a", new byte[] {1, 3}); + + assertThat(merger.merge(oldValue, newValue)).isSameAs(newValue); + } + + @Test + void testNullFieldEquality() { + UpdateIfChangedRowMerger merger = createMerger(); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + // both name null -> equal + BinaryValue oldValue = value(1, null, new byte[] {1}); + BinaryValue newValue = value(1, null, new byte[] {1}); + assertThat(merger.merge(oldValue, newValue)).isSameAs(oldValue); + + // one null, one non-null -> differs + BinaryValue newValue2 = value(1, "a", new byte[] {1}); + assertThat(merger.merge(oldValue, newValue2)).isSameAs(newValue2); + } + + @Test + void testSchemaEvolutionAlignment() { + UpdateIfChangedRowMerger merger = + createMerger( + KvFormat.COMPACTED, + new SchemaInfo(SCHEMA, SCHEMA_ID), + new SchemaInfo(SCHEMA_2, SCHEMA_2_ID)); + // latest schema has the extra nullable "age" column + merger.configureTargetColumns(null, SCHEMA_2_ID, SCHEMA_2); + + // old row uses the older schema (no age field), incoming row uses latest schema with age + // null; missing trailing fields are aligned to null -> logically equal + BinaryValue oldValue = value(1, "a", new byte[] {1, 2}); + BinaryValue newValue = value2(1, "a", new byte[] {1, 2}, null); + assertThat(merger.merge(oldValue, newValue)).isSameAs(oldValue); + + // incoming row sets age -> differs + BinaryValue newValue2 = value2(1, "a", new byte[] {1, 2}, 20); + assertThat(merger.merge(oldValue, newValue2)).isSameAs(newValue2); + + // A newly arrived request may still use the old schema. It is also aligned by column ID. + BinaryValue latestValue = value2(1, "a", new byte[] {1, 2}, null); + BinaryValue oldSchemaRequest = value(1, "a", new byte[] {1, 2}); + assertThat(merger.merge(latestValue, oldSchemaRequest)).isSameAs(latestValue); + } + + @Test + void testSchemaEvolutionAlignmentAfterDroppingMiddleColumn() { + UpdateIfChangedRowMerger merger = + createMerger( + KvFormat.COMPACTED, + new SchemaInfo(SCHEMA, SCHEMA_ID), + new SchemaInfo(SCHEMA_AFTER_DROP, SCHEMA_AFTER_DROP_ID)); + merger.configureTargetColumns(null, SCHEMA_AFTER_DROP_ID, SCHEMA_AFTER_DROP); + + BinaryValue oldValue = value(1, "removed", new byte[] {1, 2}); + BinaryValue sameValue = valueAfterDrop(1, new byte[] {1, 2}); + assertThat(merger.merge(oldValue, sameValue)).isSameAs(oldValue); + + BinaryValue changedValue = valueAfterDrop(1, new byte[] {1, 3}); + assertThat(merger.merge(oldValue, changedValue)).isSameAs(changedValue); + } + + @Test + void testDeleteReturnsNull() { + UpdateIfChangedRowMerger merger = createMerger(); + merger.configureTargetColumns(null, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {1}); + assertThat(merger.delete(oldValue)).isNull(); + assertThat(merger.deleteBehavior()).isEqualTo(DeleteBehavior.ALLOW); + } + + @Test + void testDefaultDeleteBehaviorIsAllow() { + UpdateIfChangedRowMerger merger = + new UpdateIfChangedRowMerger( + KvFormat.COMPACTED, + new TestingSchemaGetter(new SchemaInfo(SCHEMA, SCHEMA_ID)), + null); + assertThat(merger.deleteBehavior()).isEqualTo(DeleteBehavior.ALLOW); + } + + @Test + void testPartialUpdateNoOpWhenUnchanged() { + UpdateIfChangedRowMerger merger = createMerger(); + // partial update on id + name (omit data) + RowMerger partial = merger.configureTargetColumns(new int[] {0, 1}, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {9}); + // partial row only carries id + name equal to stored, data absent -> candidate keeps old + // data, so candidate equals old -> no-op + BinaryValue partialRow = value(1, "a", null); + assertThat(partial.merge(oldValue, partialRow)).isSameAs(oldValue); + } + + @Test + void testPartialUpdateAppliesWhenChanged() { + UpdateIfChangedRowMerger merger = createMerger(); + RowMerger partial = merger.configureTargetColumns(new int[] {0, 1}, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {9}); + BinaryValue partialRow = value(1, "b", null); + BinaryValue expected = value(1, "b", new byte[] {9}); + assertThat(partial.merge(oldValue, partialRow)).isEqualTo(expected); + } + + @Test + void testPartialDeleteNoOpWhenUnchanged() { + UpdateIfChangedRowMerger merger = createMerger(); + RowMerger partial = merger.configureTargetColumns(new int[] {0, 1}, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, null, new byte[] {9}); + assertThat(partial.delete(oldValue)).isSameAs(oldValue); + } + + @Test + void testPartialDeleteAppliesWhenChanged() { + UpdateIfChangedRowMerger merger = createMerger(); + RowMerger partial = merger.configureTargetColumns(new int[] {0, 1}, SCHEMA_ID, SCHEMA); + + BinaryValue oldValue = value(1, "a", new byte[] {9}); + BinaryValue expected = value(1, null, new byte[] {9}); + assertThat(partial.delete(oldValue)).isEqualTo(expected); + } +} diff --git a/website/docs/engine-flink/options.md b/website/docs/engine-flink/options.md index 62fe1e1ee05..af7e8ae1d40 100644 --- a/website/docs/engine-flink/options.md +++ b/website/docs/engine-flink/options.md @@ -94,7 +94,7 @@ See more details about [ALTER TABLE ... SET](engine-flink/ddl.md#set-properties) | table.datalake.freshness | Duration | 3min | It defines the maximum amount of time that the datalake table's content should lag behind updates to the Fluss table. Based on this target freshness, the Fluss service automatically moves data from the Fluss table and updates to the datalake table, so that the data in the datalake table is kept up to date within this target. If the data does not need to be as fresh, you can specify a longer target freshness time to reduce costs. | | table.datalake.auto-compaction | Boolean | false | If true, compaction will be triggered automatically when tiering service writes to the datalake. It is disabled by default. | | table.datalake.auto-expire-snapshot | Boolean | false | If true, snapshot expiration will be triggered automatically when tiering service commits to the datalake. It is disabled by default. | -| table.merge-engine | Enum | (None) | Defines the merge engine for the primary key table. By default, primary key table uses the [default merge engine(last_row)](table-design/merge-engines/default.md). It also supports two merge engines are `first_row`, `versioned` and `aggregation`. The [first_row merge engine](table-design/merge-engines/first-row.md) will keep the first row of the same primary key. The [versioned merge engine](table-design/merge-engines/versioned.md) will keep the row with the largest version of the same primary key. The `aggregation` merge engine will aggregate rows with the same primary key using field-level aggregate functions. | +| table.merge-engine | Enum | (None) | Defines the merge engine for the primary key table. By default, primary key table uses the [default merge engine(last_row)](table-design/merge-engines/default.md). It also supports the merge engines `first_row`, `versioned`, `aggregation` and `update_if_changed`. The [first_row merge engine](table-design/merge-engines/first-row.md) will keep the first row of the same primary key. The [versioned merge engine](table-design/merge-engines/versioned.md) will keep the row with the largest version of the same primary key. The `aggregation` merge engine will aggregate rows with the same primary key using field-level aggregate functions. The [update_if_changed merge engine](table-design/merge-engines/update-if-changed.md) keeps last-row upsert semantics but suppresses value-identical writes so that unchanged rows emit no changelog. | | table.merge-engine.versioned.ver-column | String | (None) | The column name of the version column for the `versioned` merge engine. If the merge engine is set to `versioned`, the version column must be set. | | table.delete.behavior | Enum | ALLOW | Controls the behavior of delete operations on primary key tables. Three modes are supported: `ALLOW` (default for default merge engine) - allows normal delete operations; `IGNORE` - silently ignores delete requests without errors; `DISABLE` - rejects delete requests and throws explicit errors. This configuration provides system-level guarantees for some downstream pipelines (e.g., Flink Delta Join) that must not receive any delete events in the changelog of the table. For tables with `first_row` or `versioned` or `aggregation` merge engines, this option is automatically set to `IGNORE` and cannot be overridden. Note: For `aggregation` merge engine, when set to `allow`, delete operations will remove the entire record. This configuration only applicable to primary key tables. | | table.changelog.image | Enum | FULL | Defines the changelog image mode for primary key tables. This configuration is inspired by similar settings in database systems like MySQL's `binlog_row_image` and PostgreSQL's `replica identity`. Two modes are supported: `FULL` (default) - produces both UPDATE_BEFORE and UPDATE_AFTER records for update operations, capturing complete information about updates and allowing tracking of previous values; `WAL` - does not produce UPDATE_BEFORE records. Only INSERT, UPDATE_AFTER (and DELETE if allowed) records are emitted. When WAL mode is enabled, the default merge engine is used (no merge engine configured), updates are full row updates (not partial update), and there is no auto-increment column, an optimization is applied to skip looking up old values, and in this case INSERT operations are converted to UPDATE_AFTER events. This mode reduces storage and transmission costs but loses the ability to track previous values. Only applicable to primary key tables. | diff --git a/website/docs/engine-flink/writes.md b/website/docs/engine-flink/writes.md index 2bbbcf98020..0d9a4e4d5b6 100644 --- a/website/docs/engine-flink/writes.md +++ b/website/docs/engine-flink/writes.md @@ -92,7 +92,7 @@ SELECT shop_id, user_id, num_orders FROM source; Fluss supports deleting data for primary-key tables in batch mode via `DELETE FROM` statement. The `WHERE` clause can be any condition and does not need to cover all primary key columns. :::note -`DELETE FROM` and `UPDATE` are only supported for primary-key tables using the default merge engine. Tables with the `first_row`, `versioned`, or `aggregation` merge engine reject both statements. +`DELETE FROM` and `UPDATE` are only supported for primary-key tables using the default or `update_if_changed` merge engine. Tables with the `first_row`, `versioned`, or `aggregation` merge engine reject both statements. ::: ```sql title="Flink SQL" diff --git a/website/docs/table-design/merge-engines/index.md b/website/docs/table-design/merge-engines/index.md index 1fc7f9bb13e..72b00072484 100644 --- a/website/docs/table-design/merge-engines/index.md +++ b/website/docs/table-design/merge-engines/index.md @@ -15,3 +15,4 @@ The following merge engines are supported: 2. [FirstRow Merge Engine](first-row.md) 3. [Versioned Merge Engine](versioned.md) 4. [Aggregation Merge Engine](aggregation.md) +5. [UpdateIfChanged Merge Engine](update-if-changed.md) diff --git a/website/docs/table-design/merge-engines/update-if-changed.md b/website/docs/table-design/merge-engines/update-if-changed.md new file mode 100644 index 00000000000..9138292f8aa --- /dev/null +++ b/website/docs/table-design/merge-engines/update-if-changed.md @@ -0,0 +1,77 @@ +--- +sidebar_label: UpdateIfChanged +title: UpdateIfChanged Merge Engine +sidebar_position: 6 +--- + +# UpdateIfChanged Merge Engine + +The **UpdateIfChanged Merge Engine** keeps the last-row upsert semantics of the [Default Merge Engine](default.md) +but suppresses value-identical writes. By setting `'table.merge-engine' = 'update_if_changed'` in the table +properties, an upsert whose logical field values are all equal to the currently stored row becomes a no-op: +the stored row is kept and no changelog record is emitted. When at least one field differs, the incoming row +replaces the stored row and a normal update changelog is emitted. + +This is useful for reducing unnecessary KV writes and changelog amplification when upstream data is replayed, +retried, or backfilled, and for making value-identical records that are replicated between primary-key tables a +no-op instead of producing new changelog records. + +Compared with the other merge engines: + +- `first_row` ignores every subsequent row for an existing primary key, including legitimate updates. +- `versioned` requires a version column and accepts rows whose version is equal to the stored version. +- `aggregation` applies field-level aggregation rather than last-row upsert semantics. + +## Semantics + +| Stored row | Incoming operation | Result | +|------------|-------------------------------------|-----------------------------------------------------------| +| Absent | Insert/upsert | Store the incoming row and emit an insert changelog | +| Present | Every logical field is equal | Keep the stored row and emit no changelog | +| Present | At least one logical field differs | Store the incoming row and emit the normal update changelog| +| Present | Delete | Delete the row and emit a delete changelog | +| Absent | Delete | No-op and emit no changelog | + +Equality is based on logical field values, not raw serialized bytes: + +- Null values compare using SQL/logical equality. +- Binary values are compared by content. +- Rows written with an older schema are aligned to the latest schema by stable column IDs. Fields that exist in + the latest schema but are absent from a compared row are treated as null. +- For partial updates, the incoming update is first merged into a complete candidate row and then compared with + the stored row. + +:::note +This is value-based duplicate suppression, not version ordering. Replaying a row that is identical to the current +stored value is a no-op, while replaying a historical row whose fields differ from the current value is still +accepted. Unlike `first_row`, `versioned`, and `aggregation`, this engine supports `UPDATE`, `DELETE`, and Partial +Update, just like the default merge engine. +::: + +## Example + +```sql title="Flink SQL" +CREATE TABLE T ( + k INT, + v1 DOUBLE, + v2 STRING, + PRIMARY KEY (k) NOT ENFORCED +) WITH ( + 'table.merge-engine' = 'update_if_changed' +); + +INSERT INTO T VALUES (1, 2.0, 't1'); +-- value-identical upsert: no changelog is emitted +INSERT INTO T VALUES (1, 2.0, 't1'); +-- a field differs: a normal update changelog is emitted +INSERT INTO T VALUES (1, 3.0, 't1'); + +SELECT * FROM T WHERE k = 1; + +-- Output +-- +---+-----+------+ +-- | k | v1 | v2 | +-- +---+-----+------+ +-- | 1 | 3.0 | t1 | +-- +---+-----+------+ +``` diff --git a/website/docs/table-design/table-types/pk-table.md b/website/docs/table-design/table-types/pk-table.md index 54e22d50208..43a823f1ed9 100644 --- a/website/docs/table-design/table-types/pk-table.md +++ b/website/docs/table-design/table-types/pk-table.md @@ -86,6 +86,7 @@ The following merge engines are supported: 2. [FirstRow Merge Engine](/table-design/merge-engines/first-row.md) 3. [Versioned Merge Engine](/table-design/merge-engines/versioned.md) 4. [Aggregation Merge Engine](/table-design/merge-engines/aggregation.md) +5. [UpdateIfChanged Merge Engine](/table-design/merge-engines/update-if-changed.md) ## Change Data Feed