Skip to content

feat(client): keep several write batches of a shard in flight (#1432) - #1537

Open
fenos wants to merge 6 commits into
oxia-db:mainfrom
fenos:feat/client-write-window
Open

fenos wants to merge 6 commits into
oxia-db:mainfrom
fenos:feat/client-write-window

Conversation

@fenos

@fenos fenos commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Fixes the write half of #1432: the Go client sent the next write batch of a shard only after the previous one was answered, so a shard got one batch per round trip no matter how many callers there were.

Changes

Six commits:

  1. Batcher window (oxia/batch). A batcher with MaxBatchesInFlight sends AsyncBatch batches without waiting for their responses, up to the window, and completes them in send order on its own goroutine. While the window is full, the ready batch keeps filling up to its limits (H2: Go client keeps one batch in flight per shard #1432's cheap fix comes for free: at linger 0, the calls already queued join the batch). A Barrier fires only once the batches sent before it have completed, so the split hand-off from feat: keep the write order across shard splits #1360 keeps its ordering. Without a window the batcher is unchanged.

  2. Async write stream (oxia/internal). ExecuteWriteAsync sends on a per-shard stream and matches responses in send order. When the stream ends, the writes in flight are settled together, in order:

    • Resent, in order on a new stream and ahead of any newer write, only when the end proves the server processed none of them:

      • the stream was rejected at setup: ErrNodeIsNotLeader with a leader hint, which the server only sends before reading anything;
      • or the server marked the end unprocessed (commit 5).

      Repeated rejections back off. An error a retry can't fix (e.g. the shard is gone after a split) fails them, and the batches reroute them.

    • Failed, in order, and never resent, in every other case: broken connection, EOF, or a leader that steps down mid-stream. Those writes may already be appended.

    • Partial resend: a resend whose stream breaks after taking some writes reads the stream status. Unless it proves them unprocessed, they fail instead of being sent a third time.

    • Deadline: a write with no response by its deadline aborts the stream (only while that write is still in flight), so the writes behind it fail at once.

    Only a request that was never sent is retried before sending. The synchronous ExecuteWrite path is untouched.

  3. Client wiring and options. Write batches are AsyncBatch. WithMaxWriteBatchesInFlight (default 4, as in the Java client; 0 keeps today's behavior) and WithMaxBatchSize (the byte cap already had a validation error but no option). Reads are unchanged; a read window can follow separately.

  4. Resend only what is provably unprocessed (review fix): the rules above, the partial-resend check, and the stale-abort race. A generation check and the abort it guards are now one critical section.

  5. Server: end a stream as unprocessed. The leader controller tags a write it rejects before appending it (lead.IsNotAppended). The write stream then:

    • stops reading;
    • answers every write it appended before the rejected one;
    • ends with the rejection plus ErrorMetadata["unprocessed"].

    If one of those earlier writes fails instead, the stream ends with that error, unmarked. Older clients ignore the key, and older servers never send it, so they get the safe fail path.

  6. Server: graceful close drains, and setup rejections are marked.

    • A closing leader controller stops accepting writes (rejecting new ones before append) and waits up to 5 s for the ones in flight to complete, instead of abandoning them. A fence (NewTerm) still stops it at once.
    • A stream rejected at setup is marked unprocessed whatever the error.
    • A stream whose leader closed ends marked once every appended write is answered, unless one of them failed.

    Together these make a rolling restart transparent: TestCoordinator_LeaderFailoverWithOpenClient keeps passing, with no write sent twice.

Relation to #1399

This follows #1399's rule: a write that may have reached a leader is never sent again. A resend only happens when the server proves it processed none of the writes. So a leader move, a split freeze or a graceful shutdown stays transparent when the server marks the stream, and a crash fails the writes in flight, without a duplicate in either case. Leader following on assignment changes and cancelling the stream setup on the request deadline are left to #1399; they apply to the async stream the same way.

Benchmarks

darwin/arm64, 10 cores, in-process 3-server cluster, one shard, RF=3, measured at the head of this PR (3 to 4 runs each).

SyncClient with 64 goroutines, 100-byte values. Window 0 is today's path, same build:

WAL fsync window 0 window 4
on 3.5k to 3.9k ops/s 39k to 43k ops/s
off 5.2k to 5.4k ops/s 65k to 66k ops/s

BenchmarkE2E from #1443 (async client, 5,000 operations outstanding, fsync off), main vs this branch:

main this PR
Put 7.5 to 7.7 µs/op, p50 5.6 to 5.9 ms, p99 15 to 19 ms 4.6 to 4.8 µs/op, p50 20 to 22 ms, p99 38 to 42 ms
Mixed80Read 4.7 to 5.2 µs/op 3.7 to 3.9 µs/op

Put throughput is about +60%. The higher latency is queueing at saturation, not slower requests: on main the client blocks the caller once a shard's queue is full, so only about 750 operations are actually outstanding (131k/s × 5.7 ms), while with the window nearly all 5,000 are (213k/s × 21 ms).

Tests

  • Batcher: the window overlaps batches; while full, a batch keeps filling; a barrier waits for the batches in flight; close fails the batch being formed and lets the ones in flight complete.
  • Stream (real gRPC servers):
    • writes go out before their responses;
    • a broken stream fails the writes in flight in order and never resends them;
    • a leader stepping down mid-stream (no hint) fails them, and never resends them;
    • a setup rejection resends them in order to the hinted leader, ahead of a newer write;
    • an unprocessed end resends them in order;
    • a partial resend that breaks with an unknown status fails them (no third send), and with a setup rejection retries them;
    • repeated rejections back off;
    • rejected writes of a gone shard fail with ErrShardNotFound;
    • a deadline aborts the stream.
  • Server:
    • a frozen or fenced leader tags its rejection as not appended;
    • a write stream ends unprocessed only after answering the writes appended before the rejected one, and reads nothing after it;
    • a failed earlier write leaves the end unmarked.
  • Client: two batches in flight; a split with writes in flight keeps per-key order (it fails if the barrier only waits for sends).
  • The whole oxia module and the tests module pass with -race with the window on by default, including the split, leader-hint and failover e2e tests.

fenos added 3 commits October 8, 2026 14:43
A batcher with MaxBatchesInFlight sends its batches that implement
AsyncBatch without waiting for their responses, up to the window size,
and completes them in the order they were sent on a separate goroutine.
While the window is full, the batch that is ready keeps taking the calls
that arrive until it is full too; without linger, the calls already
queued join the batch being sent instead of waiting for the next one.
A barrier is done once the batches sent before it are completed, so the
writes handed off by a split shard still follow them. Without a window
the batcher is unchanged.

Signed-off-by: fenos <fabri.feno@gmail.com>
… responses

ExecuteWriteAsync sends a write request on a stream dedicated to the
shard and returns a function that waits for its response, so that
several writes of a shard are in flight. The stream matches the
responses to the writes in the order they were sent, and settles the
writes still in flight together when it ends:

- The leader rejected them, e.g. it is no longer the leader: it applied
  none from the rejected one on, so they are sent again, in order, on a
  new stream steered by the leader hint, before any newer write. A new
  leader that rejects them too gets them with a growing delay, and an
  error that another attempt cannot fix, e.g. the shard is gone after a
  split, fails them.
- The stream broke otherwise: the leader may have applied any of them,
  so they fail, and are never sent again.

Only a request that was not sent is retried before it is sent. A write
that gets no response before its deadline aborts the stream, so the
writes behind it fail at once instead of each waiting for its own
deadline. The synchronous write path is unchanged.

Signed-off-by: fenos <fabri.feno@gmail.com>
The Go client sent the next write batch of a shard only once the
previous one was answered, so a shard got one batch per round trip,
whatever the number of callers (oxia-db#1432). A write batch is now an
AsyncBatch, and the write batchers keep up to MaxWriteBatchesInFlight
batches of a shard in flight, 4 by default as in the Java client. The
read batchers are unchanged.

The writes in flight to a shard when it splits fail with
ErrShardNotFound and are rerouted in order, and the child shards hold
the writes issued after the split until they are, as before.

WithMaxWriteBatchesInFlight sets the window, and zero keeps the previous
behavior. WithMaxBatchSize sets the byte cap of a write batch, which had
a validation error but no option.

Signed-off-by: fenos <fabri.feno@gmail.com>
Copilot AI balanced review requested due to automatic review settings October 8, 2026 12:44

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

Partial resend failures can duplicate writes, and generation aborts contain a race that can terminate a newer stream.

4 open findings
What changed in this PR

Adds pipelined per-shard write batching to improve client throughput while preserving write order and split hand-off semantics.

Changes:

  • Adds bounded asynchronous batch windows.
  • Introduces ordered asynchronous write-stream handling and retries.
  • Adds client options and coverage for batching, failures, and shard splits.
File Description
oxia/​options_client.go Adds write-window and batch-size options.
oxia/​options_client_test.go Tests option defaults and validation.
oxia/​internal/​write_stream_window_test.go Tests asynchronous stream behavior.
oxia/​internal/​write_stream_async.go Implements ordered asynchronous writes and recovery.
oxia/​internal/​rpc_provider.go Manages per-shard asynchronous streams.
oxia/​internal/​executor.go Extends the executor API.
oxia/​internal/​batch/​write_batch.go Splits batch sending from completion.
oxia/​internal/​batch/​batcher_factory.go Configures write windows.
oxia/​batch/​window.go Implements bounded batch pipelining.
oxia/​batch/​batcher.go Integrates windowed execution and shutdown.
oxia/​batch/​batcher_window_test.go Tests windows, barriers, and closure.
oxia/​batch/​batcher_factory.go Adds generic window configuration.
oxia/​batch/​batch.go Defines asynchronous batches.
oxia/​async_client_window_test.go Tests client pipelining and split ordering.
oxia/​async_client_impl.go Wires client options into write batching.

🧠 Review effort: Balanced


Give feedback about Copilot approvals in this survey to enter a drawing for a $150 gift card.

Comment thread oxia/internal/write_stream_async.go Outdated
Comment thread oxia/internal/write_stream_async.go
Comment thread oxia/internal/executor.go
Comment thread oxia/options_client.go Outdated
fenos added 3 commits October 8, 2026 17:48
… it unprocessed

The writes in flight on a shard stream were sent again whenever the
leader rejected the stream. A leader that steps down mid-stream rejects
one write, but it may have appended the writes before it, which may
still commit: sending them again could apply them twice. They are now
sent again only when the end of the stream proves that the server
processed none of them:

- a node that is not the leader rejected the stream at its setup,
  before reading any write, and hinted the leader of the shard; a
  rejection mid-stream carries no hint;
- the server marked the end as unprocessed (the new ErrorMetadata key
  "unprocessed"): it answered every write it appended before rejecting
  one, and appended none after.

Every other end fails the writes in flight, in order. A resend whose
stream breaks after taking some of the writes reads the stream status:
unless it proves them unprocessed, they fail instead of being sent a
third time. The hint of a stream rejected at its setup steers the next
stream also when no write was in flight.

The abort of a stream is now one critical section with its generation
check, and the abort on a missed deadline only happens while the write
is still in flight, so that a stale abort never ends a newer stream.
The executor and option docs state the delivery contract.

Signed-off-by: fenos <fabri.feno@gmail.com>
…d before its append

A write stream pipelines its writes: the leader appends several before
answering any, so when it rejects one and the stream ends, the writes
before it may still commit and the client cannot tell which of the
writes it has in flight were applied.

The leader controller now tags the error of a write it rejected before
appending it to its log (lead.IsNotAppended). A write stream that gets
such a rejection stops reading, answers every write it appended before
it, and then ends with the rejection marked unprocessed (the
ErrorMetadata key "unprocessed"): none of the writes left unanswered
was applied, so the client can send them again in order. If one of the
writes appended before it fails instead, its outcome is unknown and the
stream ends with its error, unmarked. A client that does not know the
key ignores it.

Signed-off-by: fenos <fabri.feno@gmail.com>
…te stream marks the ends that left nothing applied

A leader controller that closed abandoned the writes it had appended
and not committed yet: they could still commit on the next leader, so
their clients could not tell whether they were applied. Close now first
stops accepting writes, rejecting the new ones before appending them,
and waits up to closeDrainTimeout for the writes in flight to complete.
A fence (NewTerm) still stops the leader at once.

A write stream marks as unprocessed the ends after which none of its
unanswered writes was applied, beyond the write rejected before its
append:
- a stream rejected at its setup, whatever the error, read none of its
  writes;
- a stream whose leader controller closed ends once every write it had
  appended is answered, marked unprocessed unless one of them failed.

So a client sends its writes in flight to the next leader when a leader
shuts down gracefully, without ever applying one twice.

Signed-off-by: fenos <fabri.feno@gmail.com>

This branch has not been deployed

No deployments
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.

2 participants