Skip to content

[server] Keep async KV flushes on WAL batch boundaries - #4268

Merged
platinumhamburg merged 7 commits into
apache:mainfrom
platinumhamburg:fix-kv-flush-wal-batch-boundaries
Sep 14, 2026
Merged

platinumhamburg merged 7 commits into
apache:mainfrom
platinumhamburg:fix-kv-flush-wal-batch-boundaries

Conversation

@platinumhamburg

Copy link
Copy Markdown
Contributor

Track pending WAL batch ends through append, recovery, flush completion and truncation. Keep each independently completed native write on whole WAL batches so storage backpressure cannot publish an interior high watermark.

Cover native write rejection, retry, byte and entry budgets, duplicate and empty batches, leader HW publication, follower truncation and WAL recovery with real storage tests.

Purpose

Linked issue: close #4267

Brief change log

Tests

API and Format

Documentation

Track pending WAL batch ends through append, recovery, flush completion and truncation. Keep each independently completed native write on whole WAL batches so storage backpressure cannot publish an interior high watermark.

Cover native write rejection, retry, byte and entry budgets, duplicate and empty batches, leader HW publication, follower truncation and WAL recovery with real storage tests.
Reuse boxed batch ends and bypass segment scanning for a single WAL batch. Allocate the segment list only when splitting is required. Consolidate byte and entry budget coverage in the real Replica recovery test and retain the complete truncation and recovery scenario.
Accumulate complete WAL batches before checking the soft entry and byte budgets. Remove lookahead batch counters and keep trailing empty batches in the final segment. Verify grouped native writes, partial completion and recovery through real storage tests.
Keep real append/backpressure/retry and WAL recovery/truncation scenarios. Remove duplicate parameter combinations and buffer regressions while retaining existing grouping tests.
Replace the boxed boundary queue and flush snapshots with a flag on the final KV mutation of each confirmed WAL batch. Split directly on marked entries and retain the requested flush target for empty tails. Boundary cleanup follows the existing entry completion and truncation lifecycle.
Name the tail mutation marker explicitly for WAL batches and describe its per-batch call ordering. Name native write budgets as entry and byte targets rather than hard maxima.
Continue writing an existing key after deleting the KV directory so RocksDB surfaces the storage error even when the large WAL batch is applied in one native write. Keep the existing failover and full restored-data assertions.

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.

🟡 Changes recommended

One or more issues must be addressed before approval.

Get a fresh assessment by requesting another Copilot review.

Pull request overview

This PR keeps asynchronous KV flush progress aligned with complete WAL batches, preserving high-watermark and recovery correctness.

Changes:

  • Tracks WAL batch boundaries during append and recovery.
  • Splits native KV writes only at batch boundaries.
  • Adds backpressure, recovery, truncation, and segmentation tests.
File summaries
File Description
fluss-server/src/test/java/org/apache/fluss/server/replica/KvReplicaRestoreITCase.java Updated as part of this pull request.
fluss-server/src/test/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBufferTest.java Updated as part of this pull request.
fluss-server/src/test/java/org/apache/fluss/server/kv/KvTabletTest.java Updated as part of this pull request.
fluss-server/src/test/java/org/apache/fluss/server/kv/KvFlushBatchBoundaryTest.java Updated as part of this pull request.
fluss-server/src/main/java/org/apache/fluss/server/kv/prewrite/KvPreWriteBuffer.java Updated as part of this pull request.
fluss-server/src/main/java/org/apache/fluss/server/kv/KvTablet.java Updated as part of this pull request.
fluss-server/src/main/java/org/apache/fluss/server/kv/KvRecoverHelper.java Updated as part of this pull request.
Review details

Suppressed comments (1)

fluss-server/src/test/java/org/apache/fluss/server/kv/KvFlushBatchBoundaryTest.java:180

  • The new marker logic has distinct paths for duplicated appends and empty/no-op WAL batches, but this scenario only appends two non-empty batches. Existing tests cover their log generation, not that a duplicate or an empty batch leaves the next independently completed flush at a valid batch boundary; please add a real-storage case covering those paths before relying on this integration test for the recovery/truncation guarantee.
        replica.putRecordsToLeader(genKvRecordBatch(rows), null, MergeMode.DEFAULT, 0);
        replica.putRecordsToLeader(
                genKvRecordBatch(new Object[] {600, "tail"}), null, MergeMode.DEFAULT, 0);
        assertThat(replica.getLocalLogEndOffset()).isEqualTo(601);
  • Files reviewed: 7/7 changed files
  • Comments generated: 1
  • Review effort level: Lite

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

newLeaderServer.set(restoreServer);
return true;
}
putRecordBatch(tableBucket, leaderServer, triggerBatch);

@luoyuxia luoyuxia 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.

+1

@platinumhamburg
platinumhamburg merged commit b4b1fb1 into apache:main Sep 14, 2026
20 checks passed
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.

[Bug] Asynchronous KV flush can publish a high watermark inside a WAL batch

3 participants