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.