为什么选择gRPC流式通信
微服务架构中,服务间通信协议的选择直接影响系统吞吐量和延迟表现。HTTP/JSON在CRUD类场景够用,但遇到实时数据推送、大文件传输、长时间任务进度回报等场景,每次请求都建立新连接的开销就不可接受了。gRPC基于HTTP/2实现多路复用,单连接上并行传输多个请求帧,省掉了重复握手的开销。
后端开发中gRPC流式通信主要解决三类问题:服务端推送(实时数据下发)、客户端流式上传(批量数据导入)、双向流(聊天/协作场景)。本文用Go语言实现这三种模式,并给出连接池和生产级错误处理方案。
Protobuf定义与服务生成
syntax = "proto3";
package stream;
option go_package = "./proto";
service DataPushService {
rpc Subscribe(SubscribeRequest) returns (stream DataEvent);
}
service DataImportService {
rpc ImportBatch(stream ImportItem) returns (ImportResult);
}
service CollabService {
rpc Collaborate(stream CollabMessage) returns (stream CollabMessage);
}
message SubscribeRequest { string topic = 1; int32 buffer_size = 2; }
message DataEvent { string id = 1; string payload = 2; int64 timestamp = 3; }
message ImportItem { string key = 1; bytes data = 2; }
message ImportResult { int32 success_count = 1; int32 fail_count = 2; repeated string errors = 3; }
message CollabMessage { string user_id = 1; string content = 2; int64 version = 3; }
生成Go代码:
protoc --go_out=. --go-grpc_out=. proto/stream.proto
服务端流式推送实现
func (s *DataPushServer) Subscribe(
req *proto.SubscribeRequest,
stream proto.DataPushService_SubscribeServer,
) error {
ch := s.eventBus.Subscribe(req.Topic)
defer s.eventBus.Unsubscribe(req.Topic, ch)
bufferSize := int(req.BufferSize)
if bufferSize == 0 { bufferSize = 100 }
for {
select {
case <-stream.Context().Done():
return stream.Context().Err()
case event, ok := <-ch:
if !ok { return nil }
if err := stream.Send(&proto.DataEvent{
Id: event.ID, Payload: event.Payload, Timestamp: event.Timestamp,
}); err != nil {
return fmt.Errorf("send failed: %w", err)
}
}
}
}
客户端流式上传实现
func (s *DataImportServer) ImportBatch(
stream proto.DataImportService_ImportBatchServer,
) error {
var successCount, failCount int32
var errors []string
for {
item, err := stream.Recv()
if err == io.EOF {
return stream.SendAndClose(&proto.ImportResult{
SuccessCount: successCount, FailCount: failCount, Errors: errors,
})
}
if err != nil { return fmt.Errorf("recv failed: %w", err) }
if err := processItem(item); err != nil {
failCount++
errors = append(errors, fmt.Sprintf("key=%s: %v", item.Key, err))
} else { successCount++ }
}
}
双向流:实时协作场景
func (s *CollabServer) Collaborate(
stream proto.CollabService_CollaborateServer,
) error {
recvErr := make(chan error, 1)
go func() {
for {
msg, err := stream.Recv()
if err != nil { recvErr <- err; return }
s.hub.Broadcast(msg, stream)
}
}()
select {
case <-stream.Context().Done():
return stream.Context().Err()
case err := <-recvErr:
if err == io.EOF { return nil }
return err
}
}
连接池管理:避免连接泄漏
gRPC连接创建开销不小(HTTP/2握手、TLS协商),生产环境必须使用连接池:
type GrpcPool struct {
mu sync.Mutex
conns []*grpc.ClientConn
addr string
maxSize int
opts []grpc.DialOption
}
func NewGrpcPool(addr string, maxSize int, opts ...grpc.DialOption) *GrpcPool {
return &GrpcPool{
addr: addr, maxSize: maxSize,
opts: append([]grpc.DialOption{
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 30 * time.Second, Timeout: 10 * time.Second,
PermitWithoutStream: true,
}),
}, opts...),
}
}
func (p *GrpcPool) Get() (*grpc.ClientConn, error) {
p.mu.Lock()
defer p.mu.Unlock()
if len(p.conns) > 0 {
conn := p.conns[len(p.conns)-1]
p.conns = p.conns[:len(p.conns)-1]
if conn.GetState() == connectivity.Ready { return conn, nil }
conn.Close()
}
conn, err := grpc.Dial(p.addr, p.opts...)
if err != nil { return nil, fmt.Errorf("dial failed: %w", err) }
return conn, nil
}
func (p *GrpcPool) Put(conn *grpc.ClientConn) {
p.mu.Lock()
defer p.mu.Unlock()
if len(p.conns) >= p.maxSize || conn.GetState() != connectivity.Ready {
conn.Close(); return
}
p.conns = append(p.conns, conn)
}
服务治理:超时、重试与熔断
var retryOpts = []grpc.DialOption{
grpc.WithDefaultServiceConfig(`{
"methodConfig": [{
"name": [{"service": "stream.DataPushService"}],
"retryPolicy": {
"maxAttempts": 3, "initialBackoff": "0.5s",
"maxBackoff": "3s", "backoffMultiplier": 2.0,
"retryableStatusCodes": ["UNAVAILABLE", "DEADLINE_EXCEEDED"]
},
"timeout": "10s"
}]
}`),
grpc.WithKeepaliveParams(keepalive.ClientParameters{
Time: 30 * time.Second, Timeout: 10 * time.Second,
}),
}
消息中间件与服务治理配合时,gRPC流式连接的中断需要由上层做幂等性保障。服务端推送场景下,客户端断开重连后应携带上次收到的最后一条事件ID,服务端据此做断点续传。
高并发设计不是堆连接数,而是用少量长连接承载大量并发流。一个gRPC连接上可以同时运行数百个stream,只要HTTP/2帧的流量控制窗口够用。把连接池大小设置为CPU核心数的2-4倍,基本能覆盖大多数业务场景。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-wei-fu-wu-grpc-liu-shi-tong-xin-shi-zhan-shuang-xiang/