plugins/arrow-flight-rpc/docs/backpressure.md
FlightServerChannel.sendResponseBatch honours gRPC's isReady() contract: the
producer thread parks on BackpressureStrategy.waitForListener before each batch
is queued, and resumes only once gRPC reports the per-stream outbound buffer has
drained below setOnReadyThreshold. Slow consumers throttle producer wall-clock
rather than allowing buffers to accumulate to the point of allocator OOM.
Multiple producer threads may call channel.sendResponseBatch(...) concurrently
on the same stream (concurrent segment search, parallel batch generation). Arrow
Flight's ServerStreamListener.putNext is not safe for concurrent calls and
the contract requires start → putNext × N → completed in strict order on a
single thread; out-of-order or interleaved calls corrupt the stream.
Each FlightServerChannel therefore funnels submissions through a per-channel
single-threaded executor (the "eventloop"). Producer threads enqueue
BatchTasks; the eventloop dequeues them and performs the actual zero-copy
transfer plus putNext. This serialises ordering without making producer threads
contend on a per-channel mutex.
The back-pressure gate runs on the producer thread before the eventloop submission, so a slow consumer throttles allocation rather than letting the queue grow.
| Setting | Default | Property |
|---|---|---|
arrow.flight.channel.ready_timeout | 60s | node-scope |
arrow.flight.channel.outbound_buffer_threshold | 64mb | node-scope |
ready_timeout caps how long the producer thread parks before failing the
batch with StreamErrorCode.TIMED_OUT. 100ms minimum.outbound_buffer_threshold is the per-stream gRPC buffered-bytes watermark.
isReady() flips false once the per-stream outbound queue crosses this size.
Must be strictly smaller than native.allocator.pool.flight.max so the
gate engages before the allocator runs out (typical headroom: ~16 MiB).The readiness check is essentially free: the producer's call to
awaitReadyOrThrow returns immediately because gRPC's outbound buffer stays
below threshold.
isReady() flips false once gRPC's per-stream outbound buffered-bytes crosses
the threshold. The producer's next sendResponseBatch parks.OnReadyHandler after the consumer drains some bytes from the wire
(HTTP/2 WINDOW_UPDATE arriving). The producer wakes and resumes.ready_timeout causes the batch to fail
with TIMED_OUT; the producer thread is freed and the stream terminates.Client-cancel propagates via gRPC's OnCancelHandler into the strategy's cancel
callback. Any thread parked on waitForListener wakes promptly with
StreamErrorCode.CANCELLED.
native.allocator.pool.flight.max is per node, shared across all active
streams. outbound_buffer_threshold is per stream. Rule of thumb:
flight pool max >= N concurrent streams × (threshold + ~16 MiB headroom)
For N=10 concurrent streams at the default 64 MiB threshold, that's roughly
800 MiB. The threshold itself is a watermark, not a max-message-size — a single
batch larger than the threshold ships in one shot, with the producer parking on
the next sendResponseBatch. The hard constraint on per-batch size is the pool
cap (allocation must fit).
The producer thread parks while waiting for the consumer. The action handler
runs on whichever thread pool it was registered against (typically SEARCH or
GENERIC); under N concurrent slow streams, N threads from that pool are
parked simultaneously. Once the pool is exhausted, new requests targeting it
queue up or are rejected.
Mitigations:
ready_timeout to fail unresponsive streams within an acceptable window.awaitReadyOrThrow checks listener.isReady(), which reflects only gRPC's
per-stream outbound buffer. The actual push to gRPC happens later, on the
per-channel eventloop, when the eventloop dequeues a BatchTask and calls
putNext. Between enqueue and putNext, the in-flight batch sits in our own
queue — invisible to gRPC, so isReady() doesn't account for it.
Implication: a producer that allocates batches significantly faster than the
eventloop drains can pile up retained batches in the queue before gRPC's
outbound crosses the threshold and isReady() flips false. The flight pool
allocator can be pushed past outbound_buffer_threshold by roughly the size
of the queued-but-not-yet-pushed batches. The current flight pool max sizing
guidance (threshold + ~16 MiB headroom per stream) absorbs typical bursts but
doesn't formally bound this.
A byte-aware bounded queue at the eventloop entry — that parks (or rejects) the producer when the sum of queued batch sizes crosses a per-channel cap — would close the gap for the async path. Open design questions:
awaitReadyOrThrow's
shape; same thread-pool-exhaustion caveat applies) vs reject with
RESOURCE_EXHAUSTED.flight pool max.Tracked separately; not addressed in this change.
A caller that needs a hard bound today can send synchronously:
channel.sendResponseBatch(response, /* sync */ true). The batch is still
serialized and written on the channel's send executor (the same executor the
async path uses), but the calling thread blocks until the batch has been
pushed to gRPC's per-stream outbound buffer, so it cannot queue the next batch
ahead of this one — outstanding batches are bounded to one, with no eventloop
queue to grow. When that outbound buffer is full, the isReady() back-pressure
gate (above) throttles the producer.
A caller that opts in must drive a single stream from a single thread and must
not call from the send-executor thread itself. The default
sendResponseBatch(response) is unchanged — async, via the eventloop.
For workloads where the producer is mostly bottlenecked on
awaitReadyOrThrow(lots of parking, little compute per batch), the code that registers the stream action can dispatch its handler on a virtual-thread executor it owns. Park time then doesn't consume a platform thread. CPU-bound work (e.g. building each batch) should still run on a sized platform-thread pool to avoid pinning the carrier — typically by submitting that work to a separate executor and awaiting the result from the virtual thread. The framework itself does not assume virtual threads; this is a decision for the action's registrant.