ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

grpc-go Resolver 与 Balancer源码走读

grpc-go Resolver 与 Balancer源码走读 完整的示例代码可以参考前面的文章go-grpc客户端调用.在调用grpc.NewClient方法的时候会进行地址的解析.initParsedTargetAndResolverBuilder方法:1.解析原始target:parseTarget方法:Target结构体和实现方法:getResolver方法:2.① parseTarget 解析失败② scheme 没有注册对应的 resolver .需要回退使用默认 scheme拼接成 canonicalTarget:resolver.Builder就已经找到了.stream.go的newClientStream方法:OnCallBegin方法:ExitIdleMode方法:func (m *Manager) ExitIdleMode() { // idleMu 保证和 tryEnterIdleMode进入空闲的逻辑互斥 m.idleMu.Lock() defer m.idleMu.Unlock() // 已经关闭 或者 本身就不在空闲状态直接返回幂等保护 if m.isClosed() || !m.actuallyIdle { // 注释里描述3种并发场景 // 1. 定时器触发进入idle还没拿到锁新RPC进来抢先调用ExitIdleMode // 2. 空闲状态多个RPC并发同时触发OnCallBegin抢锁第一个成功后面的进来发现已经退出idle直接return // 3. 通道本来就非空闲用户手动调用 cc.Connect() 触发ExitIdleMode return } // 调用 ClientConn 的 ExitIdleMode真正唤醒连接链路【重点】 m.cc.ExitIdleMode() // 撤销进入idle时对 activeCallsCount 的修改恢复计数状态 atomic.AddInt32(m.activeCallsCount, math.MaxInt32) m.actuallyIdle false // 重置空闲计时器重新开始计时从现在开始多久无RPC就再次进入idle m.resetIdleTimerLocked(m.timeout) }exitIdleMode方法:1.修改状态:2.构建resolverresolver 可能同步更新状态、上报错误:3.start方法:4.build方法:func (b *dnsBuilder) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (resolver.Resolver, error) { // 1. 从 target.Endpoint() 解析 host、port没有端口用 defaultPort host, port, err : parseTarget(target.Endpoint(), defaultPort) if err ! nil { return nil, err } // 2. 如果 host 本身就是IP地址IPv4/IPv6 if ipAddr, err : formatIP(host); err nil { addr : []resolver.Address{{Addr: ipAddr : port}} // 直接回调 UpdateState把IP地址返回给 ClientConn cc.UpdateState(resolver.State{ Addresses: addr, Endpoints: []resolver.Endpoint{{Addresses: addr}}, }) // 返回 deadResolver空实现不需要后台协程监听不需要ResolveNow return deadResolver{}, nil } // 3. 域名场景非IP需要做DNS查询 ctx, cancel : context.WithCancel(context.Background()) d : dnsResolver{ host: host, port: port, ctx: ctx, cancel: cancel, cc: cc, // resolver.ClientConn就是ccResolverWrapper rn: make(chan struct{}, 1), // ResolveNow信号通道缓冲1 enableServiceConfig: envconfig.EnableTXTServiceConfig !opts.DisableServiceConfig, } // 创建底层DNS解析器对象 d.resolver, err internal.NewNetResolver(target.URL.Host) if err ! nil { return nil, err } // 启动后台 watcher 协程做DNS轮询解析 d.wg.Add(1) go d.watcher() return d, nil }5.watch方法:func (d *dnsResolver) watcher() { defer d.wg.Done() backoffIndex : 1 for { // 执行一次完整DNS查询A/AAAA TXT service‑config state, err : d.lookup() if err ! nil { // DNS查询出错向上报告错误给ClientConn d.cc.ReportError(err) } else { // 查询成功把地址、service‑config回调上层 err d.cc.UpdateState(*state) } var nextResolutionTime time.Time if err nil { // ✅ 查询成功分支 backoffIndex 1 // 成功重置退避计数器 // MinResolutionInterval最小重解析间隔默认30s防止频繁刷屏解析 nextResolutionTime internal.TimeNowFunc().Add(MinResolutionInterval) select { case -d.ctx.Done(): return // resolver被Close退出协程 case -d.rn: // 等待外部 ResolveNow() 触发信号 } } else { // ❌ 查询失败分支DNS出错 或者 UpdateState返回错误 // 指数退避计算下次重试时间 nextResolutionTime internal.TimeNowFunc().Add(backoff.DefaultExponential.Backoff(backoffIndex)) backoffIndex } // 等待到下次解析时间或者resolver被关闭 select { case -d.ctx.Done(): return case -internal.TimeAfterFunc(internal.TimeUntilFunc(nextResolutionTime)): } } }6.lookup方法:func (d *dnsResolver) lookup() (*resolver.State, error) { // 设置单次DNS查询超时 ResolvingTimeout通常默认5s ctx, cancel : context.WithTimeout(d.ctx, ResolvingTimeout) defer cancel() // 1. 查询 DNS SRV 记录用于grpclb负载均衡 srv, srvErr : d.lookupSRV(ctx) // 2. 查询 A/AAAA 记录拿到后端业务IP地址 addrs, hostErr : d.lookupHost(ctx) // 错误判断A/AAAA查询失败并且 SRV也失败 / SRV结果为空 → 返回hostErr本次lookup整体失败 if hostErr ! nil (srvErr ! nil || len(srv) 0) { return nil, hostErr } // 把普通IP地址封装成 Endpoint 结构 eps : make([]resolver.Endpoint, 0, len(addrs)) for _, addr : range addrs { eps append(eps, resolver.Endpoint{Addresses: []resolver.Address{addr}}) } // 组装基础State业务后端地址列表 state : resolver.State{ Addresses: addrs, Endpoints: eps, } // 如果拿到SRV记录把grpclb的lb服务地址存入state自定义字段 if len(srv) 0 { state grpclbstate.Set(state, grpclbstate.State{BalancerAddresses: srv}) } // 开启TXT service‑config时查询TXT记录填入ServiceConfig if d.enableServiceConfig { state.ServiceConfig d.lookupTXT(ctx) } return state, nil }
返回列表