gRPC Server-Side Streaming Under Backpressure: Flow Control, Send Blocking, and the Buffer Accumulation Trap
How gRPC HTTP/2 flow control interacts with Go's send path, where buffers accumulate silently, and how to design streaming RPCs that degrade gracefully.
gRPC Server-Side Streaming Under Backpressure: Flow Control, Send Blocking, and the Buffer Accumulation Trap
Server-side streaming RPCs look deceptively simple: open one request, push many responses. In practice the gap between the server's production rate and the client's consumption rate becomes a silent operational liability. gRPC's HTTP/2 flow control is the mechanism that should regulate this gap, but its interaction with Go's goroutine scheduler, the grpc-go send buffer, and your application-level write loop produces failure modes that don't surface until load is meaningful.
This article covers the mechanics precisely, because imprecise mental models lead to real production incidents: OOM kills, zombie streaming goroutines, and latency spikes that are diagnosed as CPU problems when they are actually scheduler stalls.
HTTP/2 Flow Control Is Not Enough on Its Own
HTTP/2 defines two layers of flow control: connection-level and stream-level. Each is a credit system. The receiver advertises a window (in bytes); the sender may not transmit more data than the outstanding credit allows. When the window is exhausted, the sender must block.
In grpc-go, stream.Send() ultimately calls down into the transport layer's controlBuf, which enqueues a dataFrame. The write loop dequeues frames and calls the underlying net.Conn write. If the HTTP/2 stream window is zero, the write loop stalls waiting for a WINDOW_UPDATE from the client.
This is correct behavior by protocol. The problem is what happens above that stall.
Your streaming handler looks like this:
func (s *Server) WatchEvents(req *pb.WatchRequest, stream pb.Events_WatchEventsServer) error {
ctx := stream.Context()
ch := s.eventBus.Subscribe(req.TopicId)
defer s.eventBus.Unsubscribe(req.TopicId, ch)
for {
select {
case <-ctx.Done():
return ctx.Err()
case ev, ok := <-ch:
if !ok {
return nil
}
if err := stream.Send(ev); err != nil {
return err
}
}
}
}
When the client stops reading—network partition, slow consumer, paused process—the HTTP/2 window drains. stream.Send() blocks inside the transport. The goroutine running WatchEvents is now parked. Meanwhile s.eventBus continues delivering events into ch. If ch is buffered, it fills. If it's unbuffered or full, producers upstream block too, or events are dropped, depending on your bus design. Neither outcome is acceptable without an explicit policy.
Where Buffers Accumulate
grpc-go maintains a per-stream write buffer whose default is 32 KB (defaultWriteBufSize). Data written to stream.Send() is first serialized into a protobuf byte slice, then framed and queued in the control buffer before the HTTP/2 write loop processes it. The write loop itself has no bounded queue—it drains what it has and blocks on window credit.
This means buffering happens at three distinct layers simultaneously:
- Your application channel (
chabove)—bounded only if you sized it. - grpc-go's control buffer—effectively unbounded in practice because it grows with pending frames.
- The OS TCP send buffer—kernel space, typically 4–8 MB default on Linux, tunable via
SO_SNDBUF.
Under a slow client, data piles into layer 3 first (the kernel drains slowly), then layer 2 (frames queue waiting for window), then layer 1 (your channel fills). At no point does grpc-go signal backpressure to your application code. stream.Send() simply blocks. Your goroutine count climbs. Heap grows from the enqueued protobuf slices. The heap growth triggers more frequent GC, which increases tail latency for unrelated requests on the same process.
Applying an Explicit Send Timeout
The most direct mitigation is a per-send deadline using a context with timeout wrapped around the send operation. grpc-go's stream.Send() does not accept a context argument directly—the stream's context controls the stream lifetime, not individual sends. To bound a single send, you need to race against a timer yourself:
func sendWithTimeout(stream pb.Events_WatchEventsServer, ev *pb.Event, d time.Duration) error {
done := make(chan error, 1)
go func() {
done <- stream.Send(ev)
}()
// Note: this spawns a goroutine per send under backpressure—
// acceptable only as a last-resort detection mechanism.
// Prefer the select pattern below for production.
select {
case err := <-done:
return err
case <-time.After(d):
return status.Error(codes.DeadlineExceeded, "send timeout: client not consuming")
}
}
The goroutine-per-send approach is a debugging aid, not a production pattern. The goroutine spawned on the slow path leaks until the original stream.Send() eventually unblocks or the stream is torn down. The correct production pattern is to enforce the send budget at the channel selection layer:
case ev, ok := <-ch:
if !ok {
return nil
}
sendDone := make(chan error, 1)
go func() { sendDone <- stream.Send(ev) }()
select {
case err := <-sendDone:
if err != nil {
return err
}
case <-time.After(s.cfg.SendTimeout):
// Record metric: stream.send_timeout_total
return status.Errorf(codes.ResourceExhausted, "client not consuming: evicting stream")
case <-ctx.Done():
return ctx.Err()
}
Evicting the stream on timeout is an explicit policy choice. The alternative—dropping the event and continuing—is only appropriate for loss-tolerant use cases like metrics tails or live dashboards. For ordered event delivery, eviction and reconnect is the correct contract.
gRPC Flow Control Knobs in grpc-go
gRPC exposes initial window sizes via ServerOption:
grpc.NewServer(
grpc.InitialWindowSize(64 * 1024), // per-stream
grpc.InitialConnWindowSize(1024 * 1024), // per-connection
)
Smaller initial windows cause the client to issue WINDOW_UPDATE frames more frequently for the same throughput, increasing control-plane overhead. Larger windows allow more in-flight data before backpressure kicks in—useful for high-bandwidth streams over high-RTT paths, but they increase the buffer accumulation risk on a slow consumer. There is no universally correct value. Tune toward smaller windows when your consumer reliability is low (mobile clients, edge deployments); tune toward larger when your consumers are co-located services on the same datacenter fabric.
The MaxRecvMsgSize and MaxSendMsgSize options bound individual message size, not stream volume. They do not protect against the buffer accumulation problem.
Observing the Failure Before It Bites
Two metrics expose slow-consumer pressure without code instrumentation:
Goroutine count (runtime.NumGoroutine()): a streaming server with N active streams should have a goroutine count that scales predictably with N. If the count grows faster than stream count, your send goroutines are accumulating.
Heap live objects (runtime/metrics or go_memstats_heap_inuse_bytes): protobuf serialization allocates. Blocked send goroutines hold references to those allocations. A climbing heap that tracks with goroutine count and not with RPS is a slow-consumer signature.
At the gRPC layer, enable the channelz service in non-production environments:
import "google.golang.org/grpc/channelz/service"
service.RegisterChannelzServiceToServer(grpcServer)
Channelz exposes per-stream send and receive byte counts, window sizes, and last call timestamps via a gRPC service or HTTP JSON endpoint. It is the single most useful diagnostic surface for streaming RPC production issues and is almost never enabled.
Architecture Decision Framework
When designing a server-side streaming RPC, answer these four questions before writing the handler:
1. Is the event stream loss-tolerant?
If yes, use a dropping channel on the subscribe side and never block stream.Send() beyond a short timeout. Drop metrics are your health signal. If no, you must either evict slow consumers or implement a durable replay mechanism (e.g., a Redis stream cursor the client uses to resume).
2. What is the maximum acceptable goroutine lifetime for a stalled stream?
Set SendTimeout to that value. Evict and force reconnect. Document the reconnect contract in your proto comments, not just your runbook.
3. Who owns flow control—the gRPC layer or your application? If your events arrive from a message queue (Kafka, SQS, Redis Streams), the consumer group offset is your application-level flow control. Do not pull from the queue faster than you can confirm the client has consumed. Couple your queue pull rate to a semaphore counting in-flight sends per stream.
4. Do you need ordered delivery or just eventual delivery? Ordered delivery under reconnect requires sequence numbers in your proto and a resume-token mechanism. gRPC streams are ordered within a single stream lifetime but not across reconnects. This is not a gRPC limitation—it is a distributed systems invariant. Design the protocol accordingly.
Server-side streaming is a sharp tool. The HTTP/2 flow control it rides is precise and correct, but it pushes blocking behavior into your application goroutines without warning. Making that blocking visible, bounded, and policy-driven is the difference between a streaming RPC that degrades gracefully and one that takes down the process.