Go gRPC双向流无需重建更新Send请求授权头刷新Token方案咨询
问题核心原理说明
首先你要知道:gRPC基于HTTP/2协议实现,双向流属于单个HTTP/2流,请求头仅在流初始化阶段发送一次,流建立后所有Send调用都是在同一个流中传输数据帧,无法再更新初始请求头,这是HTTP/2协议的固有特性,所以你想通过更新请求头的方式在已建立的流上每次Send带新token,本身是走不通的。
现有实现的问题
- PerRPCCredentials调用时机不符合预期
对于普通一元RPC,GetRequestMetadata确实每次请求都会调用;但对于流RPC,该方法仅在调用service.CreateWatch()创建流的那一刻触发一次,后续所有stream.Send()操作不会再触发该方法,所以就算你刷新了tokenAuth里的token也不会生效。 - tokenAuth结构体用了值接收者
你GetRequestMetadata的接收者是(t tokenAuth)值类型,方法内部修改t.token、以及goroutine里修改的都是值副本,原结构体的token字段根本不会被更新,这也是你之前刷新逻辑无效的核心原因之一。 - GetRequestMetadata内启动goroutine会导致内存泄漏
每次调用GetRequestMetadata都启动一个25分钟生命周期的goroutine,只要程序不退出,这些goroutine会一直残留,长时间运行会有内存溢出风险。
可行解决方案
方案1:应用层透传token(最推荐,完全满足流不中断的需求)
不需要依赖gRPC的请求头传token,直接修改你双向流的Proto定义,在请求消息结构里增加authorization和tenant_id字段:
message Request { string self_link = 1; string authorization = 2; // 新增字段,存Bearer token string tenant_id = 3; // 新增字段 }
客户端每次Send前都获取最新的token塞到请求里:
func (w WatchClient) CreateWatch() error { stream, err := w.service.CreateWatch(context.Background()) if err != nil { return err } ticker := time.NewTicker(25 * time.Minute) defer ticker.Stop() for range ticker.C { // 每次Send前获取最新token token, err := w.tokenRequester.GetToken() // 把tokenRequester实例挂载到WatchClient即可 if err != nil { return err } topic := &proto.Request{ SelfLink: w.config.TopicSelfLink, Authorization: "Bearer " + token, TenantId: w.config.TenantID, } err = stream.Send(topic) if err != nil { return err } } return nil }
服务端每次收到流消息时,优先从消息体里取token校验即可,流全程不需要重建。
方案2:流自动优雅重建(不需要改Proto)
如果没办法修改Proto定义,就采用流断连自动重建的方案,对上层业务逻辑影响很小:
- 首先修正你的
PerRPCCredentials实现,用指针接收者,并且单独做token的定时刷新,不要在GetRequestMetadata里启动goroutine:
type tokenAuth struct { tenantID string tokenRequester auth.PlatformTokenGetter mu sync.RWMutex // 加锁保证并发读写安全 token string } func NewTokenAuth(tenantID string, requester auth.PlatformTokenGetter) *tokenAuth { ta := &tokenAuth{ tenantID: tenantID, tokenRequester: requester, } // 启动单独的goroutine定时刷新token go func() { ticker := time.NewTicker(25 * time.Minute) defer ticker.Stop() for range ticker.C { ta.mu.Lock() ta.token, _ = requester.GetToken() ta.mu.Unlock() } }() return ta } func (t *tokenAuth) RequireTransportSecurity() bool { return false } func (t *tokenAuth) GetRequestMetadata(_ context.Context, _ ...string) (map[string]string, error) { t.mu.RLock() defer t.mu.RUnlock() // 首次调用如果token为空就主动拉取一次 if t.token == "" { token, err := t.tokenRequester.GetToken() if err != nil { return nil, err } t.token = token } return map[string]string{ "tenant-id": t.tenantID, "authorization": "Bearer " + t.token, }, nil }
- 流创建逻辑加重试重建机制:
func (w WatchClient) CreateWatch() error { topic := &proto.Request{SelfLink: w.config.TopicSelfLink} // 无限重试重建流,遇到错误直接重建,对上层业务透明 for { stream, err := w.service.CreateWatch(context.Background()) if err != nil { w.logger.Errorf("create watch stream failed, retry after 1s: %v", err) time.Sleep(time.Second) continue } sendTicker := time.NewTicker(25 * time.Minute) for range sendTicker.C { err = stream.Send(topic) if err != nil { w.logger.Warnf("send failed, will rebuild stream: %v", err) sendTicker.Stop() break } } } }
服务端检测到token过期断开流后,客户端会自动用最新的token重建流,重建时间很短,业务几乎无感知。
内容的提问来源于stack exchange,提问作者Trevor Johnson
相关产品推荐
相关产品推荐

