Files
cs-note/hhs/gRPC/3. 服务端实现/09-Streaming Handler.md
T
2026-05-24 11:42:38 +08:00

17 KiB
Raw Blame History

tags, create time
tags create time
gRPC
Go
Streaming
RPC
2026-05-11 16:30

Streaming Handler

概述

本文档讲解 gRPC 三种 streaming RPC handler 的编写模式——Server Streaming、Client Streaming 和 Bidirectional Streaming。核心主题是循环内的生命周期管理:如何正确感知 client 断开、如何在大数据量下避免 goroutine 泄漏、以及如何优雅地处理背压。

[!warning] Streaming handler 最大的陷阱:goroutine 泄漏 如果你在 loop 里 spawn 了子 goroutine 但没有正确的退出机制,一旦 client 断开连接,server 端的 goroutine 会永远无法回收。下面每一个例子都会标注如何避免这个问题。

Server Streaming Handler

Server streaming:client 发一次请求,server 持续返回多条响应。

func (s *server) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
	users := getAllUsers() // 从 DB / cache 获取

	for _, u := range users {
		if err := stream.Send(&pb.ListUsersResponse{User: u}); err != nil {
			// client 已经断开或主动取消
			s.logger.Warn("send failed", "error", err)
			return status.Error(codes.Internal, err.Error())
		}
	}
	return nil // 正常结束:数据全部发送完毕
}

核心要点

  • Server 通过 stream.Send() 逐条发送,每次 Send 都会编码 Protobuf 并通过 HTTP/2 frame 推送到 client
  • stream.Context() 是唯一可靠的退出信号 —— 当 client cancel 或网络中断时,context 会被取消
  • 正常结束(遍历完全部数据)返回 nil;断连时 context 报错
  • 不要在 loop 之外做 cleanup,否则即使 send 失败也会执行不必要的清理

大数据量场景下的超时检测

上面的简单示例只适用于结果集可控的场景。如果用户量大(比如百万级),需要在循环中主动检查 context 提前退出:

func (s *server) ListUsers(req *pb.ListUsersRequest, stream pb.UserService_ListUsersServer) error {
	ctx := stream.Context()
	results := hugeQuery(ctx) // 可能返回大量数据

	for i, u := range results {
		// 每发送 100 条检查一下 context
		if i >= 100 && len(results) >= 1000 {
			select {
			case <-ctx.Done():
				return ctx.Err() // client 已取消
			default:
			}
		}

		if err := stream.Send(&pb.ListUsersResponse{User: u}); err != nil {
			return status.Error(codes.Internal, err.Error())
		}
	}
	return nil
}

这里有个设计取舍:每次 Send 都检查 vs 每隔 N 次检查。全检查保证响应快但多了 select 开销;间隔检查减少 CPU 消耗但有延迟。根据业务对"取消感知的灵敏度"来决定。

数据流转时序

sequenceDiagram
    participant C as Client
    participant S as Server Handler
    participant DB as Data Source

    C->>S: Send(ListUsersRequest)
    S->>DB: Query All Users
    DB-->>S: []User

    loop For Each User
        alt Client Still Connected
            S->>S: Check Context
            S->>C: Send(User)
        else Context Cancelled
            S->>S: Return ctx.Err()
        end
    end

    alt Normal Completion
        C->>S: (no more data)
        S-->>C: Return nil (EOF)
    else Send Failed
        S-->>C: Return error Status
    end

Client Streaming Handler

Client streaming:client 持续发送多条请求,server 汇总后返回一次响应。

func (s *server) UploadData(stream pb.UserService_UploadDataServer) (*pb.UploadResult, error) {
	ctx := stream.Context()
	var totalBytes int64
	var processedItems []*Item

	for {
		chunk, err := stream.Recv()
		if err == io.EOF {
			break // client 结束了发送
		}
		if err != nil {
			// client 主动取消或网络异常 —— 不算业务错误,直接退出
			if ctx.Err() != nil {
				return nil, ctx.Err()
			}
			return nil, status.Errorf(codes.InvalidArgument, "recv failed: %v", err)
		}

		totalBytes += int64(len(chunk.Data))
		processedItems = append(processedItems, chunk.ToItem())
		// TODO: 批量写入 DB,不要每来一条就 flush
	}

	return &pb.UploadResult{TotalBytes: totalBytes, Count: int32(len(processedItems))}, nil
}

核心要点

  • 用 io.EOF 判断 client 结束发送(注意:是 io.EOF,不是其他 error)
  • Recv 的 error 可能来自两种情况:client 正常关闭(io.EOF)或 client 取消/断连。后者不应视为业务错误,而应直接通过 ctx.Err() 退出
  • Recv 循环必须在 EOF 之前 return 或 break,否则会无限阻塞
  • 所有临时状态放在 loop 内,返回前一次性聚合结果

内存安全注意事项

当 client 发送的数据量不可控时,必须限制内存使用:

const maxBufferSize = 100 * 1024 * 1024 // 100 MB limit

func (s *server) UploadData(stream pb.UserService_UploadDataServer) (*pb.UploadResult, error) {
	var totalBytes int64
	var items []*Item // ⚠️ 生产环境应写流到磁盘而非放内存

	for {
		chunk, err := stream.Recv()
		if err == io.EOF {
			break
		}
		if err != nil {
			return nil, status.Errorf(codes.InvalidArgument, "recv failed: %v", err)
		}

		totalBytes += int64(len(chunk.Data))
		if totalBytes > maxBufferSize {
			return nil, status.Errorf(codes.ResourceExhausted,
				"exceeds max upload size: %d bytes", maxBufferSize)
		}

		items = append(items, chunk.ToItem())
	}

	return aggregateResult(items), nil
}

[!tip] 何时该用临时文件? 如果单条消息超过 10MB 或者总量难以预估,直接将 chunk.Data 写到一个 os.File 中比存在内存切片里安全得多。handler 返回前一次性处理临时文件即可。

Bidirectional Streaming Handler

Bidirectional streaming 是最复杂的模式:收和发两个方向完全独立,各自有自己的生命周期。

func (s *server) ChatRoom(stream pb.UserService_ChatRoomServer) error {
	ctx := stream.Context()
	inboundCh := make(chan *pb.ChatMessage, 100)
	errCh := make(chan error, 1) // 错误传播通道

	// ---- 接收方向(后台 goroutine)----
	go func() {
		defer close(inboundCh)
		for {
			msg, err := stream.Recv()
			if err != nil {
				errCh <- err // EOF / cancel / network error
				return
			}
			select {
			case inboundCh <- msg:
			default:
				// channel full, drop message
			}
		}
	}()

	// ---- 发送方向(主循环)----
	for {
		select {
		case <-ctx.Done():
			return ctx.Err() // client 断了
		case err := <-errCh:
			return err // recv 出错,主循环退出
		case envelope := <-inboundCh:
			envelope := transform(envelope) // 业务转换
			if err := stream.Send(envelope); err != nil {
				return err
			}
		}
	}
}

核心架构

flowchart TB
    subgraph "Handler Goroutine"
        S["Send Loop<br/>for-select"]
    end

    subgraph "Recv Goroutine"
        R["Recv Loop"] -->|"inboundCh buffered"| B["Business Logic"]
    end

    R --> stream
    S --> stream
    B -->|"outboundCh buffered"| S

    stream["HTTP/2 Stream"]

    style S fill:#4FC08D,color:#fff
    style R fill:#00B6BC,color:#fff
    style B fill:#FF9F43,color:#000

关键设计原则

  1. Context.Done() 是唯一可靠的退出信号 — 不管哪个方向出错,最终都要通过 context 来协调
  2. 收和发是两个独立的生命周期 — recv 在一个 goroutine,send 在主循环
  3. 务必给 channel 设置 buffer,否则 sender 堵住时 receiver 也会饿死

背压策略选择

channel full 时如何处理,决定了系统的韧性:

策略 实现方式 适用场景 风险
Drop select + default(如上文) 实时性优先,丢消息可接受 丢失关键消息
Block 普通 ch <- msg 吞吐量重要,允许等待 receiver 被拖慢
Backpressure select + ctx.Done() 需要优雅降级 代码略复杂
Reject 统计丢弃数,超标后返回 error 有 SLA 要求的系统 直接中断连接

推荐大多数场景用 Backpressure 模式——超时未消费则丢弃,连续丢弃过多时通过 errCh 通知主循环终止:

func (s *server) ChatRoom(stream pb.UserService_ChatRoomServer) error {
	ctx := stream.Context()
	inboundCh := make(chan *pb.ChatMessage, 100)
	errCh := make(chan error, 1)

	go func() {
		defer close(inboundCh)
		drops := 0
		for {
			msg, err := stream.Recv()
			if err != nil {
				errCh <- err
				return
			}
			select {
			case inboundCh <- msg:
			case <-time.After(10 * time.Millisecond):
				drops++
				if drops > 1000 { // 连续丢弃过多,通知主循环终止
					errCh <- status.Errorf(codes.ResourceExhausted,
						"too many messages dropped")
					return
				}
			}
		}
	}()

	// 主循环同上,新增对 errCh 的监听即可
	for {
		select {
		case <-ctx.Done():
			return ctx.Err()
		case err := <-errCh:
			return err
		case envelope := <-inboundCh:
			if err := stream.Send(envelope); err != nil {
				return err
			}
		}
	}
}

[!tip] bidirectional streaming 的设计直觉 想象成一个管道:一端进水(Recv)、一端出水(Send)。中间的业务逻辑可以是任意的转换、过滤、聚合。但管壁(context)不能漏。

错误处理对照表

streaming handler 中的错误来源复杂,不同错误的处理方式也不同:

错误来源 stream.Recv 返回值 处理方式
Client 正常关闭连接 io.EOF break 或 return nil
Client 主动取消请求 context.Canceled(包装为普通 error) 记录日志,break
Network 中断 非 nil error(含 canceled) 记录日志,break
数据格式错误 具体 error 原样返回给 client
业务校验失败 status.Errorf(codes.InvalidArgument, ...) 返回具体错误码,终止循环
数据量超限 status.Errorf(codes.ResourceExhausted, ...) 返回错误码,终止循环

[!note] 关于 Recv 返回 error 的细节 gRPC-go 会把 context.Canceled 包装成普通 error 返回,不会返回 io.EOF。所以判断 client 断开应该同时检查两种 case。

flowchart TD
    A["Recv() returns error?"] -->|"No"| B["继续处理消息"]
    A -->|"Yes"| C{"err == io.EOF?"}
    C -->|"Yes"| D["client 正常关闭<br/>return nil ✅"]
    C -->|"No"| E{"codes.FromError is OK?"}
    E -->|"Yes"| F["canceled / network issue<br/>return ctx.Err() ⚠️"]
    E -->|"No"| G["真实错误<br/>return status.Errorf ❌"]

Context 超时传播

Streaming handler 中,context 超时由 client 侧的 WithTimeout 控制,server 必须主动感知并退出:

func (s *server) LongRunningOp(req *pb.LongRequest, stream pb.Service_LongRunningOpServer) error {
	ctx := stream.Context()

	for i := 0; i < 100; i++ {
		// 每一步都检查 context
		select {
		case <-ctx.Done():
			return ctx.Err() // client 取消了
		default:
		}

		result := doStep(i)
		if err := stream.Send(result); err != nil {
			return err // send 失败也立即退出
		}
		time.Sleep(100 * time.Millisecond)
	}
	return nil
}

这个模式的精髓是 "每一次 IO 操作前都检查 context"。只要有一次遗漏,就可能变成孤儿 goroutine。

关闭与连接检测

在 streaming 场景中,理解"如何感知对端断开"比"如何发消息"更重要。

Client 关闭的信号路径

flowchart LR
    A["Client 正常关闭发送端"] -->|"Recv() → io.EOF"| B["handler return nil ✅"]
    C["Client Cancel context"] -->|"Recv() → ctx.Err()"| D["handler return ctx.Err() ⚠️"]
    E["网络断开"] -->|"Recv() → connection error"| F["handler return error ⚠️"]
    G["不读 response"] -->|"Send() → stream error"| H["handler return error ❌"]
Client 行为 Server 感知现象 handler 处理
正常关闭发送端 Recv() 返回 io.EOF return nil ✅
Cancel context Recv() 返回 context canceled return ctx.Err() ⚠️
网络断开 Recv() 返回连接相关 error return err ⚠️
不读 response Send() 返回 stream error return err ❌

核心原则:无论哪种情况,handler 的正确做法都是尽快 return,让 gRPC 内部清理资源。不需要手动 close channel 或做其他清理。

为什么不存在 CloseSend 通知 channel

有些开发者误以为 gRPC 提供了类似 CloseSend() chan struct{} 的通知通道。实际上 gRPC-Go 并没有这样的 API。Server 只能通过以下两种方式检测 client 停止发送:

  1. Recv() io.EOF —— client 正常调用 CloseSend() 后
  2. Recv() error —— context 取消或网络异常

任何依赖"某个 channel 被关闭"来判断的模式都是不可靠的。

资源清理模式

Streaming handler 中经常需要订阅外部事件源(DB watcher、消息队列等),资源清理是关键:

func (s *server) WatchEvents(req *pb.WatchRequest, stream pb.Service_WatchServer) error {
	ctx := stream.Context()

	// 订阅事件流 —— 必须传 context
	ch, cancel := s.eventBroker.Subscribe(ctx, req.GetFilter())
	defer cancel() // handler 返回时自动释放订阅

	for {
		select {
		case <-ctx.Done():
			return nil // client 断了
		case event, ok := <-ch:
			if !ok {
				return nil // broker 关闭了 channel
			}
			if err := stream.Send(event); err != nil {
				return err // send 失败
			}
		}
	}
}

要点总结:

  1. Subscribe(ctx, ...) 一定要传 context,不能用 context.Background()
  2. defer cancel() 确保 handler 退出时释放订阅
  3. select 中始终优先检测 ctx.Done()
sequenceDiagram
    participant C as Client
    participant H as Handler
    participant EB as Event Broker
    participant S as Stream

    H->>EB: Subscribe(ctx, filter)
    EB-->>H: ch, cancel

    loop Until Disconnected
        alt Client Still Connected
            H->>H: select { ctx.Done(), <-ch }
            EB->>H: event
            H->>S: Send(event)
        else Client Disconnected
            H->>H: ctx.Done() triggers
            H->>H: defer cancel()
            H->>EB: Unsubscribe via cancel
            H-->>C: Return nil
        end
    end

生产实践要点

写完 handler 只是第一步,生产环境还需要考虑以下方面:

背压与速率限制

  • Server Streaming 中如果 client 读取速度慢于 server 发送速度,gRPC 会在 TCP/HTTP2 层积压数据。建议:监控 Send 耗时,如果 P99 持续升高则降低发送频率或加入 rate limiter
  • Client Streaming 中可通过 maxReceiveMessageSize dial option 限制最大消息大小

监控指标

在 streaming handler 中埋点三个关键指标:

// 1. handler 存活时长
duration := time.Since(start)
histogram.Record("grpc.stream.duration", duration.Milliseconds())

// 2. 发送/接收消息计数
counter.Increment("grpc.stream.messages_sent")

// 3. 非正常终止次数
counter.Increment("grpc.stream.errors", errorCode)

优雅关闭(Graceful Stop)

当 server 要 shutdown 时,不能直接 kill 正在执行的 streaming handler。正确流程:

  1. 调用 grpc.Server.GracefulStop() —— 拒绝新连接
  2. 等待已有 handler 自然返回(可以设 deadline)
  3. 超时后强制 shutdown
srv.GracefulStop()          // 停止接收新请求
// 或带超时的优雅关闭:
stopCh := make(chan struct{})
go func() {
	srv.GracefulStop()
	close(stopCh)
}()
select {
case <-stopCh:
    log.Println("all handlers finished")
case <-time.After(30 * time.Second):
    srv.Stop() // 强制杀
    log.Println("force stopped after timeout")
}

常见问题 Checklist

写完 streaming handler 后逐项核对:

  • Recv 错误分清了 io.EOF 和 ctx.Err() 吗?(正常结束 vs client 取消的处理不同)
  • Context 取消了吗?(ctx.Done() 在所有关键路径被检测)
  • Send 失败正确处理了吗?(不会因为 send 错误继续无效循环)
  • goroutine 会泄漏吗?(所有 spawn 的 goroutine 都有明确退出条件,bidirectional 中 recv goroutine 的错误能传播到主循环)
  • Channel 有 buffer 吗?(bidirectional 模式下 sender/receiver 互相不饿死)
  • 内存有上限吗?(client streaming 中对累积的消息数量做了限制)
  • 消息大小有限制吗?(gRPC 默认 4MB 限制,超出会直接报错)
  • Subscribe 传了 context 吗?(外部事件源随 handler 一起退出)
  • 大数据集做了间隔检测吗?(百万级数据不会无意义地每个元素都 select)
  • 背压策略选对了吗?(drop/block/backpressure/reject 符合业务需求)

关联笔记

  • hhs/gRPC/3. 服务端实现/08-Server 搭建与注册
  • hhs/gRPC/3. 服务端实现/10-健康检查与反射
  • hhs/gRPC/4. 客户端开发/13-Streaming Client
  • hhs/gRPC/6. 工程实践篇/20-性能优化与压测