Skip to content

[server] Support update_if_changed merge engine for primary-key tables - #4439

Open
litiliu wants to merge 2 commits into
apache:mainfrom
litiliu:feature/update-if-changed-merge-engine
Open

litiliu wants to merge 2 commits into
apache:mainfrom
litiliu:feature/update-if-changed-merge-engine

Conversation

@litiliu

@litiliu litiliu commented Sep 20, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Introduce an UPDATE_IF_CHANGED merge engine for primary-key tables ('table.merge-engine' = 'update_if_changed').

It keeps last-row upsert semantics but suppresses value-identical writes:

Stored row Incoming operation Result
Absent Insert/upsert Store the row, emit an insert changelog record
Present Every logical field is equal Keep the stored row, emit no changelog record (no-op)
Present At least one field differs Store the row, emit a normal update changelog record
Present Delete Delete the row, emit a delete changelog record
Absent Delete No-op

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, so canSkipOldValueLookup stays false).

Fixes #4343

Changes

  • Add MergeEngineType.UPDATE_IF_CHANGED, handle it in MergeEngineType.fromString(), and document update_if_changed as a supported value of the existing table.merge-engine option
  • New UpdateIfChangedRowMerger, supporting full-row, partial-update, and partial-delete mutations with logical row equality, wired into RowMerger.create()
  • Schema-version alignment by stable column IDs, including schemas with dropped columns
  • FlinkTableSink: allow partial updates and UPDATE/DELETE for this engine
  • Docs: new merge-engine page + listings/options updates

Test Plan

  • UpdateIfChangedRowMergerTest and RowMergerCreateTest (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)
  • Both client end-to-end tests run with FULL and WAL changelog images
  • FlinkTableSinkITCase#testUpdateIfChangedMergeEngineSupportsMutations (Flink sink integration: partial INSERT, UPDATE, and DELETE)

🤖 AI-assisted changes - reviewed by human developer

@litiliu
litiliu marked this pull request as draft September 20, 2026 06:38
@litiliu
litiliu force-pushed the feature/update-if-changed-merge-engine branch 7 times, most recently from 9ebf163 to 5da67ad Compare September 20, 2026 10:03
@litiliu
litiliu force-pushed the feature/update-if-changed-merge-engine branch from 5da67ad to 5a7eec9 Compare September 21, 2026 03:22
@litiliu
litiliu marked this pull request as ready for review September 21, 2026 06:05
@litiliu

litiliu commented Sep 21, 2026

Copy link
Copy Markdown
Contributor Author

@luoyuxia PTAL, when you have time thanks!

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 Medium severity

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_CHANGED and expose it via existing table.merge-engine configuration/documentation.
  • Implement UpdateIfChangedRowMerger with 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.

@luoyuxia

Copy link
Copy Markdown
Contributor

@coding-otter-bot review

@coding-otter-bot

Copy link
Copy Markdown

Automated code review could not complete

The review for commit 5a7eec929e7757f35580b10c608fcf8f30e4ea0d stopped before results could be published.

Failure code: REVIEW_DRY_RUN_FAILED
Invocation: inv_0cb1d0d0-ec78-4186-9e23-9e84986f822d

After the problem is resolved, post a new comment to retry:

@coding-otter-bot review

@luoyuxia

Copy link
Copy Markdown
Contributor

@coding-otter-bot review

@coding-otter-bot

Copy link
Copy Markdown

Automated code review

Problem

Not provided; the supplied review result was empty.

Approach

Not provided; the supplied review result was empty.

Review coverage

Limited: no additional finding is implied, but some material paths remain unresolved.

  • The reviewer response was empty, so no summary, coverage, or findings were available to preserve.

Findings

No actionable findings survived the final audit.

@litiliu

litiliu commented Oct 8, 2026

Copy link
Copy Markdown
Contributor Author

@luoyuxia The locale-sensitive parsing issue noted by Copilot is fixed in 014e28b (Locale.ROOT). I also explained and resolved the schema-null thread based on the tablet-server SchemaGetter path. All 10 CI checks are green. Could you please review this PR when you have time? Thanks!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Feature] Support an update-if-changed merge engine for primary-key tables

4 participants