【go语言grpc之client端源码分析三】
go语言grpc之server端源码分析三
- newClientStream
- newAttemptLocked
上一篇在介绍了grpc.Dial之后,然后再介绍一下后面的
//创建RPC客户端client := pb.NewGreetsClient(conn)//设置超时时间_, cancel := context.WithTimeout(context.Background(), time.Second)defer cancel()reply, err := client.SayHello(context.Background(), &pb.HelloRequest{Name: "小超", Message: "回来吃饭吗"})if err != nil {log.Fatalf("couldn not greet: %v", err)return}log.Println(reply.Name, reply.Message)
然后看一下pb.NewGreetsClient还有SayHello的方法。
func NewGreetsClient(cc grpc.ClientConnInterface) GreetsClient {return &greetsClient{cc}
}func (c *greetsClient) SayHello(ctx context.Context, in *HelloRequest, opts ...grpc.CallOption) (*HelloReply, error) {out := new(HelloReply)err := c.cc.Invoke(ctx, "/proto.Greets/SayHello", in, out, opts...)if err != nil {return nil, err}return out, nil
}
可以看出来核心就是调用ClientConn的Invoke方法。
来看一下具体的实现
func Invoke(ctx context.Context, method string, args, reply interface{}, cc *ClientConn, opts ...CallOption) error {return cc.Invoke(ctx, method, args, reply, opts...)
}var unaryStreamDesc = &StreamDesc{ServerStreams: false, ClientStreams: false}func invoke(ctx context.Context, method string, req, reply interface{}, cc *ClientConn, opts ...CallOption) error {cs, err := newClientStream(ctx, unaryStreamDesc, cc, method, opts...)if err != nil {return err}if err := cs.SendMsg(req); err != nil {return err}return cs.RecvMsg(reply)
}
所以这里就是上面的三个方法,
newClientStream
func newClientStream(ctx context.Context, desc *StreamDesc, cc *ClientConn, method string, opts ...CallOption) (_ ClientStream, err error) {if channelz.IsOn() {cc.incrCallsStarted()defer func() {if err != nil {cc.incrCallsFailed()}}()}// Provide an opportunity for the first RPC to see the first service config// provided by the resolver.if err := cc.waitForResolvedAddrs(ctx); err != nil {return nil, err}var mc serviceconfig.MethodConfigvar onCommit func()var newStream = func(ctx context.Context, done func()) (iresolver.ClientStream, error) {return newClientStreamWithParams(ctx, desc, cc, method, mc, onCommit, done, opts...)}rpcInfo := iresolver.RPCInfo{Context: ctx, Method: method}rpcConfig, err := cc.safeConfigSelector.SelectConfig(rpcInfo)if err != nil {return nil, toRPCErr(err)}return newStream(ctx, func() {})
}
可以看出来这个方法就是newClientStreamWithParams。
func newClientStreamWithParams(ctx context.Context, desc *StreamDesc, cc *ClientConn, method string, mc serviceconfig.MethodConfig, onCommit, doneFunc func(), opts ...CallOption) (_ iresolver.ClientStream, err error) {c := defaultCallInfo()if mc.WaitForReady != nil {c.failFast = !*mc.WaitForReady}// Possible context leak:// The cancel function for the child context we create will only be called// when RecvMsg returns a non-nil error, if the ClientConn is closed, or if// an error is generated by SendMsg.// https://github.com/grpc/grpc-go/issues/1818.var cancel context.CancelFuncif mc.Timeout != nil && *mc.Timeout >= 0 {ctx, cancel = context.WithTimeout(ctx, *mc.Timeout)} else {ctx, cancel = context.WithCancel(ctx)}defer func() {if err != nil {cancel()}}()for _, o := range opts {if err := o.before(c); err != nil {return nil, toRPCErr(err)}}c.maxSendMessageSize = getMaxSize(mc.MaxReqSize, c.maxSendMessageSize, defaultClientMaxSendMessageSize)c.maxReceiveMessageSize = getMaxSize(mc.MaxRespSize, c.maxReceiveMessageSize, defaultClientMaxReceiveMessageSize)if err := setCallInfoCodec(c); err != nil {return nil, err}callHdr := &transport.CallHdr{Host: cc.authority,Method: method,ContentSubtype: c.contentSubtype,DoneFunc: doneFunc,}// Set our outgoing compression according to the UseCompressor CallOption, if// set. In that case, also find the compressor from the encoding package.// Otherwise, use the compressor configured by the WithCompressor DialOption,// if set.var cp Compressorvar comp encoding.Compressorif ct := c.compressorType; ct != "" {callHdr.SendCompress = ctif ct != encoding.Identity {comp = encoding.GetCompressor(ct)if comp == nil {return nil, status.Errorf(codes.Internal, "grpc: Compressor is not installed for requested grpc-encoding %q", ct)}}} else if cc.dopts.cp != nil {callHdr.SendCompress = cc.dopts.cp.Type()cp = cc.dopts.cp}if c.creds != nil {callHdr.Creds = c.creds}cs := &clientStream{callHdr: callHdr,ctx: ctx,methodConfig: &mc,opts: opts,callInfo: c,cc: cc,desc: desc,codec: c.codec,cp: cp,comp: comp,cancel: cancel,firstAttempt: true,onCommit: onCommit,}if !cc.dopts.disableRetry {cs.retryThrottler = cc.retryThrottler.Load().(*retryThrottler)}cs.binlog = binarylog.GetMethodLogger(method)if err := cs.newAttemptLocked(false /* isTransparent */); err != nil {cs.finish(err)return nil, err}op := func(a *csAttempt) error { return a.newStream() }if err := cs.withRetry(op, func() { cs.bufferForRetryLocked(0, op) }); err != nil {cs.finish(err)return nil, err}if cs.binlog != nil {md, _ := metadata.FromOutgoingContext(ctx)logEntry := &binarylog.ClientHeader{OnClientSide: true,Header: md,MethodName: method,Authority: cs.cc.authority,}if deadline, ok := ctx.Deadline(); ok {logEntry.Timeout = time.Until(deadline)if logEntry.Timeout < 0 {logEntry.Timeout = 0}}cs.binlog.Log(logEntry)}if desc != unaryStreamDesc {// Listen on cc and stream contexts to cleanup when the user closes the// ClientConn or cancels the stream context. In all other cases, an error// should already be injected into the recv buffer by the transport, which// the client will eventually receive, and then we will cancel the stream's// context in clientStream.finish.go func() {select {case <-cc.ctx.Done():cs.finish(ErrClientConnClosing)case <-ctx.Done():cs.finish(toRPCErr(ctx.Err()))}}()}return cs, nil
}
可以看出来这里是初始化了clientStream这个结构体,然后是调用了
newAttemptLocked方法
newAttemptLocked
// newAttemptLocked creates a new attempt with a transport.
// If it succeeds, then it replaces clientStream's attempt with this new attempt.
func (cs *clientStream) newAttemptLocked(isTransparent bool) (retErr error) {ctx := newContextWithRPCInfo(cs.ctx, cs.callInfo.failFast, cs.callInfo.codec, cs.cp, cs.comp)method := cs.callHdr.Methodsh := cs.cc.dopts.copts.StatsHandlervar beginTime time.Timeif sh != nil {ctx = sh.TagRPC(ctx, &stats.RPCTagInfo{FullMethodName: method, FailFast: cs.callInfo.failFast})beginTime = time.Now()begin := &stats.Begin{Client: true,BeginTime: beginTime,FailFast: cs.callInfo.failFast,IsClientStream: cs.desc.ClientStreams,IsServerStream: cs.desc.ServerStreams,IsTransparentRetryAttempt: isTransparent,}sh.HandleRPC(ctx, begin)}var trInfo *traceInfoif EnableTracing {trInfo = &traceInfo{tr: trace.New("grpc.Sent."+methodFamily(method), method),firstLine: firstLine{client: true,},}if deadline, ok := ctx.Deadline(); ok {trInfo.firstLine.deadline = time.Until(deadline)}trInfo.tr.LazyLog(&trInfo.firstLine, false)ctx = trace.NewContext(ctx, trInfo.tr)}newAttempt := &csAttempt{ctx: ctx,beginTime: beginTime,cs: cs,dc: cs.cc.dopts.dc,statsHandler: sh,trInfo: trInfo,}defer func() {if retErr != nil {// This attempt is not set in the clientStream, so it's finish won't// be called. Call it here for stats and trace in case they are not// nil.newAttempt.finish(retErr)}}()if err := ctx.Err(); err != nil {return toRPCErr(err)}if cs.cc.parsedTarget.Scheme == "xds" {// Add extra metadata (metadata that will be added by transport) to context// so the balancer can see them.ctx = grpcutil.WithExtraMetadata(ctx, metadata.Pairs("content-type", grpcutil.ContentType(cs.callHdr.ContentSubtype),))}t, done, err := cs.cc.getTransport(ctx, cs.callInfo.failFast, cs.callHdr.Method)if err != nil {return err}if trInfo != nil {trInfo.firstLine.SetRemoteAddr(t.RemoteAddr())}newAttempt.t = tnewAttempt.done = donecs.attempt = newAttemptreturn nil
}
看一下这里的getTransport方法。
func (cc *ClientConn) getTransport(ctx context.Context, failfast bool, method string) (transport.ClientTransport, func(balancer.DoneInfo), error) {t, done, err := cc.blockingpicker.pick(ctx, failfast, balancer.PickInfo{Ctx: ctx,FullMethodName: method,})if err != nil {return nil, nil, toRPCErr(err)}return t, done, nil
}
注意一下这里的cc.blockingpicker.pick。是不是很熟悉,其实就是前面的。
然后看一下pick这个方法
// pick returns the transport that will be used for the RPC.
// It may block in the following cases:
// - there's no picker
// - the current picker returns ErrNoSubConnAvailable
// - the current picker returns other errors and failfast is false.
// - the subConn returned by the current picker is not READY
// When one of these situations happens, pick blocks until the picker gets updated.
func (pw *pickerWrapper) pick(ctx context.Context, failfast bool, info balancer.PickInfo) (transport.ClientTransport, func(balancer.DoneInfo), error) {var ch chan struct{}var lastPickErr errorfor {pw.mu.Lock()if pw.done {pw.mu.Unlock()return nil, nil, ErrClientConnClosing}if pw.picker == nil {ch = pw.blockingCh}//fmt.Printf(" pw.picker:%v nil:%v ch == pw.blockingCh:%v\n", pw.picker, pw.picker == nil, ch == pw.blockingCh)if ch == pw.blockingCh {// This could happen when either:// - pw.picker is nil (the previous if condition), or// - has called pick on the current picker.pw.mu.Unlock()select {case <-ctx.Done():var errStr stringif lastPickErr != nil {errStr = "latest balancer error: " + lastPickErr.Error()} else {errStr = ctx.Err().Error()}switch ctx.Err() {case context.DeadlineExceeded:return nil, nil, status.Error(codes.DeadlineExceeded, errStr)case context.Canceled:return nil, nil, status.Error(codes.Canceled, errStr)}case <-ch:}continue}ch = pw.blockingChp := pw.pickerpw.mu.Unlock()pickResult, err := p.Pick(info)if err != nil {if err == balancer.ErrNoSubConnAvailable {continue}if _, ok := status.FromError(err); ok {// Status error: end the RPC unconditionally with this status.return nil, nil, err}// For all other errors, wait for ready RPCs should block and other// RPCs should fail with unavailable.if !failfast {lastPickErr = errcontinue}return nil, nil, status.Error(codes.Unavailable, err.Error())}acw, ok := pickResult.SubConn.(*acBalancerWrapper)if !ok {logger.Errorf("subconn returned from pick is type %T, not *acBalancerWrapper", pickResult.SubConn)continue}if t := acw.getAddrConn().getReadyTransport(); t != nil {if channelz.IsOn() {return t, doneChannelzWrapper(acw, pickResult.Done), nil}return t, pickResult.Done, nil}if pickResult.Done != nil {// Calling done with nil error, no bytes sent and no bytes received.// DoneInfo with default value works.pickResult.Done(balancer.DoneInfo{})}logger.Infof("blockingPicker: the picked transport is not ready, loop back to repick")// If ok == false, ac.state is not READY.// A valid picker always returns READY subConn. This means the state of ac// just changed, and picker will be updated shortly.// continue back to the beginning of the for loop to repick.}
}
注意这里的p.Pick,这个就是在pickfirst中进行更新后调用的,如下
func (b *pickfirstBalancer) UpdateSubConnState(sc balancer.SubConn, s balancer.SubConnState) {if logger.V(2) {logger.Infof("pickfirstBalancer: UpdateSubConnState: %p, %v", sc, s)}if b.sc != sc {if logger.V(2) {logger.Infof("pickfirstBalancer: ignored state change because sc is not recognized")}return}b.state = s.ConnectivityStateif s.ConnectivityState == connectivity.Shutdown {b.sc = nilreturn}switch s.ConnectivityState {case connectivity.Ready:b.cc.UpdateState(balancer.State{ConnectivityState: s.ConnectivityState, Picker: &picker{result: balancer.PickResult{SubConn: sc}}})case connectivity.Connecting:b.cc.UpdateState(balancer.State{ConnectivityState: s.ConnectivityState, Picker: &picker{err: balancer.ErrNoSubConnAvailable}})case connectivity.Idle:b.cc.UpdateState(balancer.State{ConnectivityState: s.ConnectivityState, Picker: &idlePicker{sc: sc}})case connectivity.TransientFailure:b.cc.UpdateState(balancer.State{ConnectivityState: s.ConnectivityState,Picker: &picker{err: s.ConnectionError},})}
}
如果成功了就是返回ready下面的balancer.PickResult。返回了SubConn,然后失败了就是返回了balancer.PickResult但是其中的err是错误的,这样在cc.blockingpicker.pick的时候,就可以返回具体的成功或者失败。
这样完成了整个逻辑的闭环。
下面的sendMsg和ReckMsg和之前的没有什么区别,就是在http2的基础上加上hpack头部压缩和proto的序列化,就不进行赘述了。
相关文章:
【go语言grpc之client端源码分析三】
go语言grpc之server端源码分析三newClientStreamnewAttemptLocked上一篇在介绍了grpc.Dial之后,然后再介绍一下后面的 //创建RPC客户端client : pb.NewGreetsClient(conn)//设置超时时间_, cancel : context.WithTimeout(context.Background(), time.Second)defer c…...

Android 基础知识4-2.6LinearLayout(线性布局)
一、LinearLayout的概述 线性布局(LinearLayout)主要以水平或垂直方式来排列界面中的控件。并将控件排列到一条直线上。在线性布局中,如果水平排列,垂直方向上只能放一个控件,如果垂直排列,水平方向上也只能…...

补充前端面试题(三)
图片懒加载<!DOCTYPE html> <html lang"en"><head><meta charset"UTF-8"><meta http-equiv"X-UA-Compatible" content"IEedge"><meta name"viewport" content"widthdevice-width, in…...

.net开发安卓入门-自动升级(配合.net6 webapi 作为服务端)
文章目录思路客户端权限清单(AndroidManifest.xml)权限列表(完整内容看 权限清单(AndroidManifest.xml))打开外部应用的权限(完整内容看 权限清单(AndroidManifest.xml))添加文件如下…...

分享111个HTML艺术时尚模板,总有一款适合您
分享111个HTML艺术时尚模板,总有一款适合您 111个HTML艺术时尚模板下载链接:https://pan.baidu.com/s/1sYo2IPma4rzeku3yCG7jGw?pwdk8dx 提取码:k8dx Python采集代码下载链接:采集代码.zip - 蓝奏云 时尚理发沙龙服务网站模…...

spring之Spring AOP基于注解
文章目录前言一、Spring AOP基于注解的所有通知类型1、前置通知2、后置通知3、环绕通知4、最终通知5、异常通知二、Spring AOP基于注解之切面顺序三、Spring AOP基于注解之通用切点三、Spring AOP基于注解之连接点四、Spring AOP基于注解之全注解开发前言 通知类型包括&#x…...

LeetCode题目笔记——6362. 合并两个二维数组 - 求和法
文章目录题目描述题目链接题目难度——简单方法一:常规双指针遍历代码/Python方法二:字典\哈希表代码/Python总结题目描述 给你两个 二维 整数数组 nums1 和 nums2. nums1[i] [idi, vali] 表示编号为 idi 的数字对应的值等于 vali 。nums2[i] [idi, …...

【C#基础】C# 常用语句讲解
序号系列文章3【C#基础】C# 数据类型总结4【C#基础】C# 变量和常量的使用5【C#基础】C# 运算符总结文章目录前言语句概念1,迭代语句1.1 for 语句1.2 foreach 语句1.3 while 语句1.4 do 语句2,选择语句2.1,if 语句2.2,else 语句2.3…...

腾讯云——负载均衡CLB
负载均衡 CLB 提供四层(TCP 协议/UDP 协议/TCP SSL 协议)和七层(HTTP 协议/HTTPS 协议)负载均衡。您可以通过 CLB 将业务流量分发到多个后端服务器上,消除单点故障并保障业务可用性。CLB 自身采用集群部署,…...

6.关于系统服务的思考—— native vs java
文章目录native服务 以sensor service为例Native 服务java 服务, 以vibrate为例java 服务 以一个demo为例native服务 以sensor service为例 service启动 SystemServer.startBootstrapServices---->>>mSystemServiceManager.startService—>>>Sen…...

SQL语句创建视图:
前言 🎈个人主页:🎈 :✨✨✨初阶牛✨✨✨ 🐻推荐专栏: 🍔🍟🌯 c语言初阶 🔑个人信条: 🌵知行合一 🍉本篇简介:>:介绍数据库中有关视图的知识,参考学校作业. 金句分享:…...

使用BP神经网络和Elman Net预测航班价格(Matlab代码实现)
👨🎓个人主页:研学社的博客💥💥💞💞欢迎来到本博客❤️❤️💥💥🏆博主优势:🌞🌞🌞博客内容尽量做到思维缜密…...

JavaWeb9-volatile解决内存可见性和指令重排序问题
目录 1.解决内存可见性问题 2.解决指令重排序问题 3.volatile缺点 4.特使使用场景 volatile(易变的,易挥发的,不稳定的)可以解决内存可见性和指令重排序的问题。 1.解决内存可见性问题 代码在写入 volatile 修饰的变量时&am…...

Docker - 镜像操作命令
镜像名称一般分为两部分组成:[repository]:[tag]在没有指定tag时,默认是latest,代表最新版本的镜像1.下载docker镜像 docker pull repository:tag2.查看本地所有镜像 docker images3.创建镜像别名 docker tag repository:tag repository111:tag4.查看镜像…...

全栈之路-前端篇 | 第三讲.基础前置知识【前端标准与研发工具】学习笔记
欢迎关注「全栈工程师修炼指南」公众号 点击 👇 下方卡片 即可关注我哟! 设为「星标⭐」每天带你 基础入门 到 进阶实践 再到 放弃学习! 涉及 企业运维、网络安全、应用开发、物联网、人工智能、大数据 学习知识 “ 花开堪折直须折,莫待无花…...

Tomcat 线上调优记录
原始Tomcat配置 启动参数Plaintext-Xms256m -Xmx512m -XX:MaxPermSize128m Tomcat 参数配置XML<Executor name"tomcatThreadPool" namePrefix"catalina-exec-" maxThreads"1500" minSpareThreads"50" maxIdleTime"600000&q…...
学习 Python 之 Pygame 开发坦克大战(四)
学习 Python 之 Pygame 开发坦克大战(四)坦克大战添加音效1. 初始化音效2. 加入游戏开始音效和坦克移动音效3. 添加坦克开火音效4. 添加装甲削减音效5. 添加坦克爆炸音效6. 添加子弹击中边界音效坦克大战添加音效 我的素材放到了百度网盘里,…...

New和Malloc的使用及其差异
1,new的使用关于new的定义:new其实就是告诉计算机开辟一段新的空间,但是和一般的声明不同的是,new开辟的空间在堆上,而一般声明的变量存放在栈上。通常来说,当在局部函数中new出一段新的空间,该…...
2023年细胞生物学复习汇总
细胞分化 1.什么是细胞分化?细胞分化的特点是什么? 答:(1)细胞分化(cell differentiation)是指同一来源的细胞逐渐产生出形态结构、功能特征各不相同的细胞类群的过程,其结果是在空间…...

光伏VSG-基于虚拟同步发电机的光伏并网逆变器系统MATLAB仿真
采用MATLAB2021b仿真!!!仿真模型1光伏电池模块(采用MATLAB自带光伏模块)、MPPT控制模块、升压模块、VSG控制模块、电流滞环控制模块。2s时改变光照强度 !!!VSG输出有功功率、无功功率…...
Golang 面试经典题:map 的 key 可以是什么类型?哪些不可以?
Golang 面试经典题:map 的 key 可以是什么类型?哪些不可以? 在 Golang 的面试中,map 类型的使用是一个常见的考点,其中对 key 类型的合法性 是一道常被提及的基础却很容易被忽视的问题。本文将带你深入理解 Golang 中…...

理解 MCP 工作流:使用 Ollama 和 LangChain 构建本地 MCP 客户端
🌟 什么是 MCP? 模型控制协议 (MCP) 是一种创新的协议,旨在无缝连接 AI 模型与应用程序。 MCP 是一个开源协议,它标准化了我们的 LLM 应用程序连接所需工具和数据源并与之协作的方式。 可以把它想象成你的 AI 模型 和想要使用它…...
Qwen3-Embedding-0.6B深度解析:多语言语义检索的轻量级利器
第一章 引言:语义表示的新时代挑战与Qwen3的破局之路 1.1 文本嵌入的核心价值与技术演进 在人工智能领域,文本嵌入技术如同连接自然语言与机器理解的“神经突触”——它将人类语言转化为计算机可计算的语义向量,支撑着搜索引擎、推荐系统、…...

苍穹外卖--缓存菜品
1.问题说明 用户端小程序展示的菜品数据都是通过查询数据库获得,如果用户端访问量比较大,数据库访问压力随之增大 2.实现思路 通过Redis来缓存菜品数据,减少数据库查询操作。 缓存逻辑分析: ①每个分类下的菜品保持一份缓存数据…...

04-初识css
一、css样式引入 1.1.内部样式 <div style"width: 100px;"></div>1.2.外部样式 1.2.1.外部样式1 <style>.aa {width: 100px;} </style> <div class"aa"></div>1.2.2.外部样式2 <!-- rel内表面引入的是style样…...

uniapp微信小程序视频实时流+pc端预览方案
方案类型技术实现是否免费优点缺点适用场景延迟范围开发复杂度WebSocket图片帧定时拍照Base64传输✅ 完全免费无需服务器 纯前端实现高延迟高流量 帧率极低个人demo测试 超低频监控500ms-2s⭐⭐RTMP推流TRTC/即构SDK推流❌ 付费方案 (部分有免费额度&#x…...
Rust 异步编程
Rust 异步编程 引言 Rust 是一种系统编程语言,以其高性能、安全性以及零成本抽象而著称。在多核处理器成为主流的今天,异步编程成为了一种提高应用性能、优化资源利用的有效手段。本文将深入探讨 Rust 异步编程的核心概念、常用库以及最佳实践。 异步编程基础 什么是异步…...
ip子接口配置及删除
配置永久生效的子接口,2个IP 都可以登录你这一台服务器。重启不失效。 永久的 [应用] vi /etc/sysconfig/network-scripts/ifcfg-eth0修改文件内内容 TYPE"Ethernet" BOOTPROTO"none" NAME"eth0" DEVICE"eth0" ONBOOT&q…...

Yolov8 目标检测蒸馏学习记录
yolov8系列模型蒸馏基本流程,代码下载:这里本人提交了一个demo:djdll/Yolov8_Distillation: Yolov8轻量化_蒸馏代码实现 在轻量化模型设计中,**知识蒸馏(Knowledge Distillation)**被广泛应用,作为提升模型…...

JVM虚拟机:内存结构、垃圾回收、性能优化
1、JVM虚拟机的简介 Java 虚拟机(Java Virtual Machine 简称:JVM)是运行所有 Java 程序的抽象计算机,是 Java 语言的运行环境,实现了 Java 程序的跨平台特性。JVM 屏蔽了与具体操作系统平台相关的信息,使得 Java 程序只需生成在 JVM 上运行的目标代码(字节码),就可以…...