gRPC Streaming: Real-Time Communication That Actually Works
WebSockets kept disconnecting. SSE was one-way. Long polling burned CPU. Then we discovered gRPC streaming. Now handling 100K concurrent streams at 50MB/s per stream.
The Four Patterns You Need to Know
gRPC offers four communication patterns, each optimized for different use cases. Understanding when to use each pattern is crucial for building efficient real-time systems.
Server streaming pushes continuous updates from server to client. Perfect for live dashboards, price feeds, and notifications. The server maintains state and pushes data as it becomes available.
Client streaming allows clients to send continuous data to the server. Ideal for file uploads, telemetry data, and bulk data ingestion. The server processes data as it arrives.
Bidirectional streaming enables full-duplex communication. Both client and server can send/receive independently. Essential for chat systems, collaborative editing, and real-time gaming.
Unary RPC is traditional request-response. Simple and reliable for standard API operations where streaming isn't needed.
service StreamingService {
rpc Subscribe(Topic) returns (stream Event); // Server streaming
rpc Upload(stream Chunk) returns (Summary); // Client streaming
rpc Chat(stream Message) returns (stream Message); // Bidirectional
rpc GetStatus(Empty) returns (Status); // Unary
}
Server Streaming: Push Updates to Clients
Perfect for live dashboards, price feeds, log tailing:
// server.go
type server struct {
pb.UnimplementedStreamingServiceServer
subscribers map[string][]chan *pb.Event
mu sync.RWMutex
}
func (s *server) Subscribe(req *pb.Topic, stream pb.StreamingService_SubscribeServer) error {
// Create channel for this subscriber
ch := make(chan *pb.Event, 100) // Buffer prevents slow consumer blocking
// Register subscriber
s.mu.Lock()
s.subscribers[req.Name] = append(s.subscribers[req.Name], ch)
s.mu.Unlock()
// Clean up on exit
defer func() {
s.mu.Lock()
subs := s.subscribers[req.Name]
for i, sub := range subs {
if sub == ch {
s.subscribers[req.Name] = append(subs[:i], subs[i+1:]...)
break
}
}
s.mu.Unlock()
close(ch)
}()
// Send events to client
for {
select {
case event := <-ch:
if err := stream.Send(event); err != nil {
return err // Client disconnected
}
case <-stream.Context().Done():
return nil // Client cancelled
}
}
}
// Publish events to all subscribers
func (s *server) PublishEvent(topic string, event *pb.Event) {
s.mu.RLock()
subscribers := s.subscribers[topic]
s.mu.RUnlock()
for _, ch := range subscribers {
select {
case ch <- event:
default:
// Channel full, drop message (or implement backpressure)
log.Printf("Dropping message for slow consumer")
}
}
}
Client with automatic reconnection:
// client.go
func SubscribeWithReconnect(client pb.StreamingServiceClient, topic string) {
backoff := time.Second
maxBackoff := time.Minute
for {
err := subscribe(client, topic)
if err != nil {
log.Printf("Stream error: %v, reconnecting in %v", err, backoff)
time.Sleep(backoff)
backoff *= 2
if backoff > maxBackoff {
backoff = maxBackoff
}
} else {
backoff = time.Second // Reset on successful connection
}
}
}
func subscribe(client pb.StreamingServiceClient, topic string) error {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
stream, err := client.Subscribe(ctx, &pb.Topic{Name: topic})
if err != nil {
return err
}
for {
event, err := stream.Recv()
if err == io.EOF {
return nil // Stream ended normally
}
if err != nil {
return err // Connection error
}
// Process event
handleEvent(event)
}
}
Client Streaming: Efficient Uploads
10x faster than REST for large files:
// server.go - Handle file upload
func (s *server) Upload(stream pb.StreamingService_UploadServer) error {
var (
fileID = uuid.New().String()
file *os.File
size int64
checksum = sha256.New()
)
// Create temp file
file, err := os.CreateTemp("", "upload-*.tmp")
if err != nil {
return status.Errorf(codes.Internal, "failed to create file: %v", err)
}
defer os.Remove(file.Name()) // Clean up on error
// Receive chunks
for {
chunk, err := stream.Recv()
if err == io.EOF {
// Upload complete
break
}
if err != nil {
return status.Errorf(codes.Internal, "failed to receive chunk: %v", err)
}
// Write chunk
n, err := file.Write(chunk.Data)
if err != nil {
return status.Errorf(codes.Internal, "failed to write: %v", err)
}
size += int64(n)
checksum.Write(chunk.Data)
// Enforce size limit
if size > 1<<30 { // 1GB
return status.Errorf(codes.InvalidArgument, "file too large")
}
}
// Move to final location
finalPath := fmt.Sprintf("/uploads/%s", fileID)
if err := os.Rename(file.Name(), finalPath); err != nil {
return status.Errorf(codes.Internal, "failed to save: %v", err)
}
// Send response
return stream.SendAndClose(&pb.Summary{
FileId: fileID,
Size: size,
Checksum: fmt.Sprintf("%x", checksum.Sum(nil)),
})
}
Client with progress tracking:
// client.go - Upload with progress
func UploadFile(client pb.StreamingServiceClient, filepath string) error {
file, err := os.Open(filepath)
if err != nil {
return err
}
defer file.Close()
// Get file size for progress
stat, _ := file.Stat()
totalSize := stat.Size()
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
stream, err := client.Upload(ctx)
if err != nil {
return err
}
// Send file in chunks
buf := make([]byte, 1<<20) // 1MB chunks
var sent int64
for {
n, err := file.Read(buf)
if err == io.EOF {
break
}
if err != nil {
return err
}
if err := stream.Send(&pb.Chunk{
Data: buf[:n],
}); err != nil {
return err
}
sent += int64(n)
progress := float64(sent) / float64(totalSize) * 100
fmt.Printf("\rUploading: %.1f%%", progress)
}
// Get response
summary, err := stream.CloseAndRecv()
if err != nil {
return err
}
fmt.Printf("\nUpload complete: %s (%d bytes)\n", summary.FileId, summary.Size)
return nil
}
Bidirectional Streaming: Real-Time Chat
The most complex but most powerful pattern:
// server.go - Chat room implementation
type chatServer struct {
rooms map[string]*Room
mu sync.RWMutex
}
type Room struct {
clients map[string]pb.StreamingService_ChatServer
mu sync.RWMutex
}
func (s *chatServer) Chat(stream pb.StreamingService_ChatServer) error {
// First message must be join request
msg, err := stream.Recv()
if err != nil {
return err
}
if msg.Type != pb.MessageType_JOIN {
return status.Errorf(codes.InvalidArgument, "first message must be JOIN")
}
roomID := msg.RoomId
userID := msg.UserId
// Get or create room
s.mu.Lock()
room, exists := s.rooms[roomID]
if !exists {
room = &Room{
clients: make(map[string]pb.StreamingService_ChatServer),
}
s.rooms[roomID] = room
}
s.mu.Unlock()
// Add client to room
room.mu.Lock()
room.clients[userID] = stream
room.mu.Unlock()
// Notify others
s.broadcast(room, &pb.Message{
Type: pb.MessageType_USER_JOINED,
UserId: userID,
Timestamp: time.Now().Unix(),
}, userID)
// Clean up on exit
defer func() {
room.mu.Lock()
delete(room.clients, userID)
room.mu.Unlock()
s.broadcast(room, &pb.Message{
Type: pb.MessageType_USER_LEFT,
UserId: userID,
Timestamp: time.Now().Unix(),
}, userID)
}()
// Handle messages
for {
msg, err := stream.Recv()
if err == io.EOF {
return nil
}
if err != nil {
return err
}
// Validate and enrich message
msg.Timestamp = time.Now().Unix()
msg.UserId = userID
// Broadcast to room
s.broadcast(room, msg, "")
}
}
func (s *chatServer) broadcast(room *Room, msg *pb.Message, exclude string) {
room.mu.RLock()
defer room.mu.RUnlock()
for userID, client := range room.clients {
if userID == exclude {
continue
}
// Send in goroutine to prevent blocking
go func(c pb.StreamingService_ChatServer, m *pb.Message) {
if err := c.Send(m); err != nil {
log.Printf("Failed to send to %s: %v", userID, err)
}
}(client, msg)
}
}
Client with separate send/receive:
// client.go - Chat client
func StartChat(client pb.StreamingServiceClient, roomID, userID string) error {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
stream, err := client.Chat(ctx)
if err != nil {
return err
}
// Join room
if err := stream.Send(&pb.Message{
Type: pb.MessageType_JOIN,
RoomId: roomID,
UserId: userID,
}); err != nil {
return err
}
// Start receiver goroutine
go func() {
for {
msg, err := stream.Recv()
if err != nil {
log.Printf("Receive error: %v", err)
cancel()
return
}
displayMessage(msg)
}
}()
// Send messages from stdin
scanner := bufio.NewScanner(os.Stdin)
for scanner.Scan() {
text := scanner.Text()
if text == "/quit" {
return nil
}
if err := stream.Send(&pb.Message{
Type: pb.MessageType_CHAT,
Content: text,
}); err != nil {
return err
}
}
return scanner.Err()
}
Flow Control and Backpressure
gRPC has built-in flow control, but you need to handle it:
// Server-side flow control
func (s *server) StreamWithFlowControl(req *pb.Request, stream pb.Service_StreamServer) error {
// Check if client is ready before sending
for i := 0; i < 1000000; i++ {
// This blocks if client's receive buffer is full
if err := stream.Send(&pb.Response{
Data: generateData(i),
}); err != nil {
if status.Code(err) == codes.Unavailable {
// Client can't keep up
log.Printf("Client overwhelmed at message %d", i)
}
return err
}
// Optional: Add rate limiting
if i%100 == 0 {
time.Sleep(10 * time.Millisecond)
}
}
return nil
}
// Client-side flow control
func ConsumeWithBackpressure(client pb.ServiceClient) error {
stream, err := client.Stream(context.Background(), &pb.Request{})
if err != nil {
return err
}
// Process with limited concurrency
sem := make(chan struct{}, 10) // Max 10 concurrent
for {
msg, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return err
}
sem <- struct{}{} // Acquire
go func(m *pb.Response) {
defer func() { <-sem }() // Release
// Slow processing
time.Sleep(100 * time.Millisecond)
process(m)
}(msg)
}
// Wait for all to finish
for i := 0; i < cap(sem); i++ {
sem <- struct{}{}
}
return nil
}
Metadata and Headers
// Send metadata with stream
func (s *server) AuthenticatedStream(req *pb.Request, stream pb.Service_StreamServer) error {
// Send initial metadata
header := metadata.Pairs(
"stream-id", uuid.New().String(),
"server-version", "1.0.0",
)
stream.SendHeader(header)
// Stream data
for i := 0; i < 100; i++ {
stream.Send(&pb.Response{Data: fmt.Sprintf("Message %d", i)})
}
// Send trailing metadata
trailer := metadata.Pairs(
"message-count", "100",
"checksum", "abc123",
)
stream.SetTrailer(trailer)
return nil
}
// Client reads metadata
func ReadWithMetadata(client pb.ServiceClient) error {
stream, err := client.Stream(context.Background(), &pb.Request{})
if err != nil {
return err
}
// Get initial metadata
header, err := stream.Header()
if err != nil {
return err
}
streamID := header.Get("stream-id")[0]
log.Printf("Stream ID: %s", streamID)
// Read messages
for {
_, err := stream.Recv()
if err == io.EOF {
break
}
if err != nil {
return err
}
}
// Get trailing metadata
trailer := stream.Trailer()
count := trailer.Get("message-count")[0]
log.Printf("Received %s messages", count)
return nil
}
Connection Management
// Robust client with keepalive and reconnection
func NewRobustClient(addr string) (pb.ServiceClient, error) {
conn, err := grpc.Dial(addr,
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 10 * time.Second, // Send ping every 10s
Timeout: 3 * time.Second, // Wait 3s for pong
PermitWithoutStream: true, // Send even without active streams
}),
grpc.WithDefaultCallOptions(
grpc.MaxCallRecvMsgSize(50*1024*1024), // 50MB max message
grpc.MaxCallSendMsgSize(50*1024*1024),
),
grpc.WithStreamInterceptor(streamRetryInterceptor()),
)
if err != nil {
return nil, err
}
return pb.NewServiceClient(conn), nil
}
// Retry interceptor for streams
func streamRetryInterceptor() grpc.StreamClientInterceptor {
return func(ctx context.Context, desc *grpc.StreamDesc, cc *grpc.ClientConn,
method string, streamer grpc.Streamer, opts ...grpc.CallOption) (grpc.ClientStream, error) {
var stream grpc.ClientStream
var err error
for i := 0; i < 3; i++ {
stream, err = streamer(ctx, desc, cc, method, opts...)
if err == nil {
return stream, nil
}
code := status.Code(err)
if code == codes.Unavailable || code == codes.ResourceExhausted {
time.Sleep(time.Duration(i+1) * time.Second)
continue
}
break // Non-retryable error
}
return nil, err
}
}
Load Balancing Streams
// Client-side load balancing for streams
conn, err := grpc.Dial("dns:///myservice.local:50051",
grpc.WithDefaultServiceConfig(`{
"loadBalancingPolicy": "round_robin",
"healthCheckConfig": {
"serviceName": "StreamingService"
}
}`),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
// Server registers with service discovery
func RegisterWithConsul(consulAddr, serviceID, grpcAddr string) error {
client, err := consul.NewClient(&consul.Config{
Address: consulAddr,
})
if err != nil {
return err
}
return client.Agent().ServiceRegister(&consul.AgentServiceRegistration{
ID: serviceID,
Name: "streaming-service",
Port: 50051,
Address: grpcAddr,
Check: &consul.AgentServiceCheck{
GRPC: grpcAddr,
Interval: "10s",
DeregisterCriticalServiceAfter: "1m",
},
})
}
Performance Benchmarks
| Method | Throughput | Latency (p99) | CPU Usage |
|---|---|---|---|
| REST (HTTP/1.1) | 10K req/s | 50ms | 80% |
| WebSocket | 50K msg/s | 10ms | 60% |
| gRPC Unary | 80K req/s | 5ms | 40% |
| gRPC Streaming | 500K msg/s | 1ms | 35% |
Common Pitfalls
1. Not handling reconnection: Streams will disconnect. Plan for it.
2. Ignoring flow control: Fast producer + slow consumer = OOM
3. No timeout on streams: Streams can run forever. Set deadlines.
4. Blocking in handlers: Use goroutines for heavy processing
5. Not cleaning up: Always defer cleanup in stream handlers
When to Use Each Pattern
- Server streaming: Live feeds, logs, monitoring
- Client streaming: File uploads, batch processing, telemetry
- Bidirectional: Chat, gaming, collaborative editing
- Unary: Simple request-response, CRUD operations
Pro tip: Start with unary. Add streaming only where it provides clear value. The complexity isn't always worth it.