Implementing Client Streaming RPCs
1. Receiving Client Stream
Example: Server reads client stream
func (s *server) UploadFile(stream filev1.FileService_UploadFileServer) error {
var total int64
for {
chunk, err := stream.Recv()
if err == io.EOF {
return stream.SendAndClose(&filev1.UploadResult{Bytes: total})
}
if err != nil { return err }
total += int64(len(chunk.Data))
}
}
| Method | Detail |
|---|---|
stream.Recv() | Next request or io.EOF |
stream.SendAndClose(res) | Send final response and close |
2. Processing Incoming Messages
| Pattern | Detail |
|---|---|
| Streaming aggregate | Sum/count/hash incrementally |
| Buffered write | Write each chunk to disk/blob store |
| Validate per message | Fail fast — return error mid-stream |
3. Sending Final Response
| Rule | Detail |
|---|---|
| Call once | Exactly one SendAndClose |
| After EOF only | Must consume all client messages first |
4. Aggregating Stream Data
Example: Hash + size summary
h := sha256.New()
for {
c, err := stream.Recv()
if err == io.EOF {
return stream.SendAndClose(&filev1.UploadResult{
Sha256: hex.EncodeToString(h.Sum(nil)),
})
}
if err != nil { return err }
h.Write(c.Data)
}
5. Creating Client Stream
| Note | Detail |
|---|---|
| Method takes no request | Request comes via Send calls |
Stream is bound to ctx | Canceling ctx aborts the stream |
6. Sending Multiple Requests
| Method | Detail |
|---|---|
stream.Send(msg) | Push request |
| Not concurrency-safe | Serialize sends |
7. Closing Client Stream
| Method | Action |
|---|---|
stream.CloseAndRecv() | Close send-side and read final response |
| Returns | Final response or error |
8. Handling Send Errors
| Error | Detail |
|---|---|
io.EOF from Send | Server already closed; call CloseAndRecv to get status |
| Other errors | Connection level — abort |
9. Receiving Final Response
Example: Close + receive
for _, chunk := range chunks {
if err := stream.Send(chunk); err != nil { break }
}
res, err := stream.CloseAndRecv()
if err != nil { return err }
fmt.Println("uploaded", res.Bytes)
10. Using Stream Context
| Use | API |
|---|---|
| Deadline | Set on ctx passed to method |
| Server-side ctx | stream.Context() |
| Metadata in | metadata.FromIncomingContext(stream.Context()) |