Repository navigation
Conversation
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>
fenos
requested review from
RobertIndie,
coderzc,
mattisonchao and
merlimat
as code owners
October 8, 2026 12:44
Contributor
There was a problem hiding this comment.
🟡 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.
… 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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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:
Batcher window (
oxia/batch). A batcher withMaxBatchesInFlightsendsAsyncBatchbatches 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). ABarrierfires 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.Async write stream (
oxia/internal).ExecuteWriteAsyncsends 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:
ErrNodeIsNotLeaderwith a leader hint, which the server only sends before reading anything;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
ExecuteWritepath is untouched.Client wiring and options. Write batches are
AsyncBatch.WithMaxWriteBatchesInFlight(default 4, as in the Java client;0keeps today's behavior) andWithMaxBatchSize(the byte cap already had a validation error but no option). Reads are unchanged; a read window can follow separately.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.
Server: end a stream as unprocessed. The leader controller tags a write it rejects before appending it (
lead.IsNotAppended). The write stream then: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.
Server: graceful close drains, and setup rejections are marked.
NewTerm) still stops it at once.unprocessedwhatever the error.Together these make a rolling restart transparent:
TestCoordinator_LeaderFailoverWithOpenClientkeeps 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).
SyncClientwith 64 goroutines, 100-byte values. Window 0 is today's path, same build:BenchmarkE2Efrom #1443 (async client, 5,000 operations outstanding, fsync off),mainvs this branch:Put throughput is about +60%. The higher latency is queueing at saturation, not slower requests: on
mainthe 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
unprocessedend resends them in order;ErrShardNotFound;oxiamodule and thetestsmodule pass with-racewith the window on by default, including the split, leader-hint and failover e2e tests.