You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go gRPC双向流无需重建更新Send请求授权头刷新Token方案咨询

问题核心原理说明

首先你要知道:gRPC基于HTTP/2协议实现,双向流属于单个HTTP/2流,请求头仅在流初始化阶段发送一次,流建立后所有Send调用都是在同一个流中传输数据帧,无法再更新初始请求头,这是HTTP/2协议的固有特性,所以你想通过更新请求头的方式在已建立的流上每次Send带新token,本身是走不通的。

现有实现的问题
  1. PerRPCCredentials调用时机不符合预期
    对于普通一元RPC,GetRequestMetadata确实每次请求都会调用;但对于流RPC,该方法仅在调用service.CreateWatch()创建流的那一刻触发一次,后续所有stream.Send()操作不会再触发该方法,所以就算你刷新了tokenAuth里的token也不会生效。
  2. tokenAuth结构体用了值接收者
    你GetRequestMetadata的接收者是(t tokenAuth)值类型,方法内部修改t.token、以及goroutine里修改的都是值副本,原结构体的token字段根本不会被更新,这也是你之前刷新逻辑无效的核心原因之一。
  3. 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定义,就采用流断连自动重建的方案,对上层业务逻辑影响很小:

  1. 首先修正你的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
}
  1. 流创建逻辑加重试重建机制:
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 04:24:03