Implementing Server Streaming RPCs
1. Creating Streaming Handler
Example: Server stream handler
func (s *server) ListEvents(r *eventv1.ListEventsRequest, stream eventv1.EventService_ListEventsServer) error {
for _, e := range s.store.Query(r.Filter) {
if err := stream.Send(e); err != nil { return err }
if err := stream.Context().Err(); err != nil { return err }
}
return nil
}
| Signature | Detail |
|---|---|
No ctx param | Use stream.Context() |
Return nil | Signals successful end-of-stream |
2. Sending Multiple Responses
| Method | Behavior |
|---|---|
stream.Send(msg) | Push one message; respects flow control |
| Backpressure | Blocks if HTTP/2 window full |
| Concurrency | Send is not safe for concurrent calls |
3. Detecting Client Disconnect
Example: Check ctx in loop
select {
case <-stream.Context().Done():
return stream.Context().Err()
default:
}
| Signal | Detail |
|---|---|
Send returns error | EOF / connection closed |
ctx.Done() | Client canceled / deadline hit |
4. Handling Stream Errors
| Source | Action |
|---|---|
| Send error | Stop producing, return error |
| Upstream error | Wrap with status.Error |
| Partial success | Use trailer metadata or final summary message |
5. Creating Client Receiver
Example: Client receives stream
stream, err := client.ListEvents(ctx, &eventv1.ListEventsRequest{Filter: "all"})
if err != nil { return err }
for {
e, err := stream.Recv()
if err == io.EOF { break }
if err != nil { return err }
process(e)
}
6. Receiving Stream Messages
| Method | Returns |
|---|---|
stream.Recv() | Next message or io.EOF on end |
| After EOF | Stream is closed; do not call again |
7. Detecting End of Stream
| Signal | Meaning |
|---|---|
err == io.EOF | Successful clean end |
err != nil && err != io.EOF | RPC failed — check status |
8. Handling Receive Errors
Example: Status-aware receive
e, err := stream.Recv()
if err != nil && err != io.EOF {
if st, ok := status.FromError(err); ok {
log.Printf("stream failed: %s %s", st.Code(), st.Message())
}
}
9. Processing Stream Data
| Pattern | Detail |
|---|---|
| Process inline | Simple loop body |
| Channel + worker | Decouple receive from processing |
| Batch | Accumulate N then flush for throughput |
10. Closing Server Stream
| Method | Action |
|---|---|
Return nil | Clean end; client sees io.EOF |
| Return error | Sends trailer with status; client sees error |
| Set trailer | stream.SetTrailer(md) before return |