Repository navigation
Conversation
9ebf163 to
5da67ad
Compare
5da67ad to
5a7eec9
Compare
|
@luoyuxia PTAL, when you have time thanks! |
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
There are a couple of correctness issues (locale-sensitive config parsing and a potential NPE on missing schema lookup) that should be addressed before approval.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 1
Open (1)
What changed in this PR
Adds a new primary-key-table merge engine (update_if_changed) that preserves last-row upsert semantics while suppressing value-identical writes (no KV update and no changelog emission when logical row content is unchanged), and wires it through server write path, Flink sink validation, docs, and tests.
Changes:
- Introduce
MergeEngineType.UPDATE_IF_CHANGEDand expose it via existingtable.merge-engineconfiguration/documentation. - Implement
UpdateIfChangedRowMergerwith logical row equality aligned by stable column IDs, including partial update/delete support via delegation to the default merger. - Add unit + IT coverage and update website docs to describe the new engine and supported Flink mutations.
| File | Description |
|---|---|
| website/docs/table-design/table-types/pk-table.md | Add UpdateIfChanged to the supported PK table merge-engine list. |
| website/docs/table-design/merge-engines/update-if-changed.md | New documentation page describing semantics, equality rules, and examples. |
| website/docs/table-design/merge-engines/index.md | Add UpdateIfChanged to merge-engine index listing. |
| website/docs/engine-flink/writes.md | Document that Flink DELETE FROM / UPDATE are supported for default + update_if_changed. |
| website/docs/engine-flink/options.md | Extend table.merge-engine option docs to include update_if_changed. |
| fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMergerTest.java | Unit tests for logical equality, schema alignment, and partial mutations. |
| fluss-server/src/test/java/org/apache/fluss/server/kv/rowmerger/RowMergerCreateTest.java | Validate RowMerger.create() instantiates UpdateIfChanged and its semantics. |
| fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/UpdateIfChangedRowMerger.java | New merger implementation with schema-aligned logical row equality. |
| fluss-server/src/main/java/org/apache/fluss/server/kv/rowmerger/RowMerger.java | Wire new merge engine into RowMerger.create() and clarify delete no-op semantics. |
| fluss-server/src/main/java/org/apache/fluss/server/kv/KvWriteProcessor.java | Treat delete no-ops consistently by skipping when merger returns the old-value instance. |
| fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/sink/FlinkTableSinkITCase.java | Integration test that Flink sink accepts partial updates + UPDATE/DELETE for this engine. |
| fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/sink/FlinkTableSink.java | Allow partial updates and UPDATE/DELETE statements for update_if_changed tables. |
| fluss-common/src/main/java/org/apache/fluss/metadata/MergeEngineType.java | Add enum constant + parsing support + updated merge-table semantics documentation. |
| fluss-common/src/main/java/org/apache/fluss/config/ConfigOptions.java | Document update_if_changed as a supported merge engine in config description. |
| fluss-client/src/test/java/org/apache/fluss/client/table/FlussTableITCase.java | End-to-end tests verifying changelog suppression and partial mutation behavior under FULL/WAL. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
|
@coding-otter-bot review |
Automated code review could not completeThe review for commit Failure code: After the problem is resolved, post a new comment to retry:
|
|
@coding-otter-bot review |
Automated code reviewProblemNot provided; the supplied review result was empty. ApproachNot provided; the supplied review result was empty. Review coverageLimited: no additional finding is implied, but some material paths remain unresolved.
FindingsNo actionable findings survived the final audit. |

Summary
Introduce an
UPDATE_IF_CHANGEDmerge engine for primary-key tables ('table.merge-engine' = 'update_if_changed').It keeps last-row upsert semantics but suppresses value-identical writes:
Equality is based on logical field values, not raw bytes: null values use logical equality, binary values are compared by content, and rows from different schema versions are aligned to the latest schema by stable column IDs. Fields that exist in the latest schema but are absent from either compared row are treated as null. Partial updates and partial deletes are first applied to the stored row to produce a complete candidate row, which is then compared with the stored row.
Because it must read the stored value before deciding, it does not use the WAL full-row fast path that skips old-value lookup (the merger is not a
DefaultRowMerger, socanSkipOldValueLookupstays false).Fixes #4343
Changes
MergeEngineType.UPDATE_IF_CHANGED, handle it inMergeEngineType.fromString(), and documentupdate_if_changedas a supported value of the existingtable.merge-engineoptionUpdateIfChangedRowMerger, supporting full-row, partial-update, and partial-delete mutations with logical row equality, wired intoRowMerger.create()FlinkTableSink: allow partial updates and UPDATE/DELETE for this engineTest Plan
UpdateIfChangedRowMergerTestandRowMergerCreateTest(unit: logical equality, schema alignment, full-row and partial mutations, and merger creation)FlussTableITCase#testUpdateIfChangedMergeEngine(end-to-end: value-identical full-row upserts emit no changelog record while changes and deletes do)FlussTableITCase#testUpdateIfChangedMergeEngineWithPartialUpdate(end-to-end: value-identical partial updates/deletes are no-ops while changed mutations emit changelog records)FULLandWALchangelog imagesFlinkTableSinkITCase#testUpdateIfChangedMergeEngineSupportsMutations(Flink sink integration: partial INSERT, UPDATE, and DELETE)🤖 AI-assisted changes - reviewed by human developer