Golang 网络编程

目录

  • TCP网络编程
  • UDP网络编程
  • Http网络编程
  • 理解函数是一等公民
  • HttpServer源码阅读
    • 注册路由
    • 启动服务
    • 处理请求
  • HttpClient源码阅读
    • DemoCode
    • 整理思路
    • 重要的struct
    • 流程
    • transport.dialConn
    • 发送请求

TCP网络编程

存在的问题:

  • 拆包:
    • 对发送端来说应用程序写入的数据远大于socket缓冲区大小,不能一次性将这些数据发送到server端就会出现拆包的情况。
    • 通过网络传输的数据包最大是1500字节,当TCP报文的长度 - TCP头部的长度 > MSS(最大报文长度时)将会发生拆包,MSS一般长(1460~1480)字节。
  • 粘包:
    • 对发送端来说:应用程序发送的数据很小,远小于socket的缓冲区的大小,导致一个数据包里面有很多不通请求的数据。
    • 对接收端来说:接收数据的方法不能及时的读取socket缓冲区中的数据,导致缓冲区中积压了不同请求的数据。

解决方法:

  • 使用带消息头的协议,在消息头中记录数据的长度。
  • 使用定长的协议,每次读取定长的内容,不够的使用空格补齐。
  • 使用消息边界,比如使用 \n 分隔 不同的消息。
  • 使用诸如 xml json protobuf这种复杂的协议。

实验:使用自定义协议

整体的流程:

客户端:发送端连接服务器,将要发送的数据通过编码器编码,发送。

服务端:启动、监听端口、接收连接、将连接放在协程中处理、通过解码器解码数据。

1 //########################### 2//###### Server端代码 ###### 3//########################### 4 5func main() { 6 // 1. 监听端口 2.accept连接 3.开goroutine处理连接 7 listen, err := net.Listen("tcp", "0.0.0.0:9090") 8 if err != nil { 9 fmt.Printf("error : %v", err) 10 return 11 } 12 for{ 13 conn, err := listen.Accept() 14 if err != nil { 15 fmt.Printf("Fail listen.Accept : %v", err) 16 continue 17 } 18 go ProcessConn(conn) 19 } 20} 21 22// 处理网络请求 23func ProcessConn(conn net.Conn) { 24 defer conn.Close() 25 for { 26 bt,err:=coder.Decode(conn) 27 if err != nil { 28 fmt.Printf("Fail to decode error [%v]", err) 29 return 30 } 31 s := string(bt) 32 fmt.Printf("Read from conn:[%v]\n",s) 33 } 34} 35 36//########################### 37//###### Clinet端代码 ###### 38//########################### 39func main() { 40 conn, err := net.Dial("tcp", ":9090") 41 defer conn.Close() 42 if err != nil { 43 fmt.Printf("error : %v", err) 44 return 45 } 46 47 // 将数据编码并发送出去 48 coder.Encode(conn,"hi server i am here"); 49} 50 51//########################### 52//###### 编解码器代码 ###### 53//########################### 54/** 55 * 解码: 56 */ 57func Decode(reader io.Reader) (bytes []byte, err error) { 58 // 先把消息头读出来 59 headerBuf := make([]byte, len(msgHeader)) 60 if _, err = io.ReadFull(reader, headerBuf); err != nil { 61 fmt.Printf("Fail to read header from conn error:[%v]", err) 62 return nil, err 63 } 64 // 检验消息头 65 if string(headerBuf) != msgHeader { 66 err = errors.New("msgHeader error") 67 return nil, err 68 } 69 // 读取实际内容的长度 70 lengthBuf := make([]byte, 4) 71 if _, err = io.ReadFull(reader, lengthBuf); err != nil { 72 return nil, err 73 } 74 contentLength := binary.BigEndian.Uint32(lengthBuf) 75 contentBuf := make([]byte, contentLength) 76 // 读出消息体 77 if _, err := io.ReadFull(reader, contentBuf); err != nil { 78 return nil, err 79 } 80 return contentBuf, err 81} 82 83/** 84 * 编码 85 * 定义消息的格式: msgHeader + contentLength + content 86 * conn 本身实现了 io.Writer 接口 87 */ 88func Encode(conn io.Writer, content string) (err error) { 89 // 写入消息头 90 if err = binary.Write(conn, binary.BigEndian, []byte(msgHeader)); err != nil { 91 fmt.Printf("Fail to write msgHeader to conn,err:[%v]", err) 92 } 93 // 写入消息体长度 94 contentLength := int32(len([]byte(content))) 95 if err = binary.Write(conn, binary.BigEndian, contentLength); err != nil { 96 fmt.Printf("Fail to write contentLength to conn,err:[%v]", err) 97 } 98 // 写入消息 99 if err = binary.Write(conn, binary.BigEndian, []byte(content)); err != nil { 100 fmt.Printf("Fail to write content to conn,err:[%v]", err) 101 } 102 return err 103

客户端的conn一直不被Close 有什么表现?

四次挥手各个状态的如下:

1主从关闭方 被动关闭方 2established established 3Fin-wait1 4 closeWait 5Fin-wait2 6Tiem-wait lastAck 7Closed Closed

如果客户端的连接手动的关闭,它和服务端的状态会一直保持established建立连接中的状态。

1MacBook-Pro% netstat -aln | grep 9090 2tcp4 0 0 127.0.0.1.9090 127.0.0.1.62348 ESTABLISHED 3tcp4 0 0 127.0.0.1.62348 127.0.0.1.9090 ESTABLISHED 4tcp46 0 0 *.9090 *.* LISTEN

服务端的conn一直不被关闭 有什么表现?

客户端的进程结束后,会发送fin数据包给服务端,向服务端请求断开连接。

服务端的conn不关闭的话,服务端就会停留在四次挥手的close_wait阶段(我们不手动Close,服务端就任务还有数据/任务没处理完,因此它不关闭)。

客户端停留在 fin_wait2的阶段(在这个阶段等着服务端告诉自己可以真正断开连接的消息)。

1DXMdeMacBook-Pro% netstat -aln | grep 9090 2tcp4 0 0 127.0.0.1.9090 127.0.0.1.62888 CLOSE_WAIT 3tcp4 0 0 127.0.0.1.62888 127.0.0.1.9090 FIN_WAIT_2 4tcp46 0 0 *.9090 *.* LISTEN

什么是binary.BigEndian?什么是binary.LittleEndian?

对计算机来说一切都是二进制的数据,BigEndian和LittleEndian描述的就是二进制数据的字节顺序。计算机内部,小端序被广泛应用于现代性 CPU 内部存储数据;大端序常用于网络传输和文件存储。

比如:

1一个数的二进制表示为 0x12345678 2BigEndian 表示为: 0x12 0x34 0x56 0x78 3LittleEndian表示为: 0x78 0x56 0x34 0x12

UDP网络编程

思路:

UDP服务器:1、监听 2、循环读取消息 3、回复数据。

UDP客户端:1、连接服务器 2、发送消息 3、接收消息。

1// ################################ 2// ######## UDPServer ######### 3// ################################ 4func main() { 5 // 1. 监听端口 2.accept连接 3.开goroutine处理连接 6 listen, err := net.Listen("tcp", "0.0.0.0:9090") 7 if err != nil { 8 fmt.Printf("error : %v", err) 9 return 10 } 11 for{ 12 conn, err := listen.Accept() 13 if err != nil { 14 fmt.Printf("Fail listen.Accept : %v", err) 15 continue 16 } 17 go ProcessConn(conn) 18 } 19} 20 21// 处理网络请求 22func ProcessConn(conn net.Conn) { 23 defer conn.Close() 24 for { 25 bt,err:= coder.Decode(conn) 26 if err != nil { 27 fmt.Printf("Fail to decode error [%v]", err) 28 return 29 } 30 s := string(bt) 31 fmt.Printf("Read from conn:[%v]\n",s) 32 } 33} 34 35// ################################ 36// ######## UDPClient ######### 37// ################################ 38func main() { 39 40 udpConn, err := net.DialUDP("udp", nil, &net.UDPAddr{ 41 IP: net.IPv4(127, 0, 0, 1), 42 Port: 9091, 43 }) 44 if err != nil { 45 fmt.Printf("error : %v", err) 46 return 47 } 48 49 _, err = udpConn.Write([]byte("i am udp client")) 50 if err != nil { 51 fmt.Printf("error : %v", err) 52 return 53 } 54 bytes:=make([]byte,1024) 55 num, addr, err := udpConn.ReadFromUDP(bytes) 56 if err != nil { 57 fmt.Printf("Fail to read from udp error: [%v]", err) 58 return 59 } 60 fmt.Printf("Recieve from udp address:[%v], bytes:[%v], content:[%v]",addr,num,string(bytes)) 61}

Http网络编程

思路整理:

HttpServer:1、创建路由器。2、为路由器绑定路由规则。3、创建服务器、监听端口。 4启动读服务。

HttpClient: 1、创建连接池。2、创建客户端,绑定连接池。3、发送请求。4、读取响应。

1func main() { 2 mux := http.NewServeMux() 3 mux.HandleFunc("/login", doLogin) 4 server := &http.Server{ 5 Addr: ":8081", 6 WriteTimeout: time.Second * 2, 7 Handler: mux, 8 } 9 log.Fatal(server.ListenAndServe()) 10} 11 12func doLogin(writer http.ResponseWriter,req *http.Request){ 13 _, err := writer.Write([]byte("do login")) 14 if err != nil { 15 fmt.Printf("error : %v", err) 16 return 17 } 18}

HttpClient端

1func main() { 2 transport := &http.Transport{ 3 // 拨号的上下文 4 DialContext: (&net.Dialer{ 5 Timeout: 30 * time.Second, // 拨号建立连接时的超时时间 6 KeepAlive: 30 * time.Second, // 长连接存活的时间 7 }).DialContext, 8 // 最大空闲连接数 9 MaxIdleConns: 100, 10 // 超过最大的空闲连接数的连接,经过 IdleConnTimeout时间后会失效 11 IdleConnTimeout: 10 * time.Second, 12 // https使用了SSL安全证书,TSL是SSL的升级版 13 // 当我们使用https时,这行配置生效 14 TLSHandshakeTimeout: 10 * time.Second, 15 ExpectContinueTimeout: 1 * time.Second, // 100-continue 状态码超时时间 16 } 17 18 // 创建客户端 19 client := &http.Client{ 20 Timeout: time.Second * 10, //请求超时时间 21 Transport: transport, 22 } 23 24 // 请求数据 25 res, err := client.Get("http://localhost:8081/login") 26 if err != nil { 27 fmt.Printf("error : %v", err) 28 return 29 } 30 defer res.Body.Close() 31 32 bytes, err := ioutil.ReadAll(res.Body) 33 if err != nil { 34 fmt.Printf("error : %v", err) 35 return 36 } 37 fmt.Printf("Read from http server res:[%v]", string(bytes)) 38}

理解函数是一等公民

点击查看在github中函数相关的笔记

在golang中函数是一等公民,我们可以把一个函数当作普通变量一样使用。

比如我们有个函数HelloHandle,我们可以直接使用它。

1func HelloHandle(name string, age int) { 2 fmt.Printf("name:[%v] age:[%v]", name, age) 3} 4 5func main() { 6 HelloHandle("tom",12) 7}

闭包

如何理解闭包:闭包本质上是一个函数,而且这个函数会引用它外部的变量,如下例子中的f3中的匿名函数本身就是一个闭包。 通常我们使用闭包起到一个适配的作用。

例1:

1// f2是一个普通函数,有两个入参数 2func f2() { 3 fmt.Printf("f2222") 4} 5 6// f1函数的入参是一个f2类型的函数 7func f1(f2 func()) { 8 f2() 9} 10 11func main() { 12 // 由于golang中函数是一等公民,所以我们可以把f2同普通变量一般传递给f1 13 f1(f2) 14}

例2: 在上例中更进一步。f2有了自己的参数, 这时就不能直接把f2传递给f1了。

总不能傻傻的这样吧f1(f2(1,2)) ???

而闭包就能解决这个问题。

1// f2是一个普通函数,有两个入参数 2func f2(x int, y int) { 3 fmt.Println("this is f2 start") 4 fmt.Printf("x: %d y: %d \n", x, y) 5 fmt.Println("this is f2 end") 6} 7 8// f1函数的入参是一个f2类型的函数 9func f1(f2 func()) { 10 fmt.Println("this is f1 will call f2") 11 f2() 12 fmt.Println("this is f1 finished call f2") 13} 14 15// 接受一个两个参数的函数, 返回一个包装函数 16func f3(f func(int,int) ,x,y int) func() { 17 fun := func() { 18 f(x,y) 19 } 20 return fun 21} 22 23func main() { 24 // 目标是实现如下的传递与调用 25 f1(f3(f2,6,6)) 26}

实现方法的回调:

下面的例子中实现这样的功能:就好像是我设计了一个框架,定好了整个框架运转的流程(或者说是提供了一个编程模版),框架具体做事的函数你根据自己的需求自己实现,我的框架只是负责帮你回调你具体的方法。

1// 自定义类型,handler本质上是一个函数 2type HandlerFunc func(string, int) 3 4// 闭包 5func (f HandlerFunc) Serve(name string, age int) { 6 f(name, age) 7} 8 9// 具体的处理函数 10func HelloHandle(name string, age int) { 11 fmt.Printf("name:[%v] age:[%v]", name, age) 12} 13 14func main() { 15 // 把HelloHandle转换进自定义的func中 16 handlerFunc := HandlerFunc(HelloHandle) 17 // 本质上会去回调HelloHandle方法 18 handlerFunc.Serve("tom", 12) 19 20 // 上面两行效果 == 下面这行 21 // 只不过上面的代码是我在帮你回调,下面的是你自己主动调用 22 HelloHandle("tom",12) 23}

HttpServer源码阅读

注册路由

直观上看注册路由这一步,就是它要做的就是将在路由器url pattern和开发者提供的func关联起来。 很容易想到,它里面很可能是通过map实现的。

1func main() { 2 // 创建路由器 3 // 为路由器绑定路由规则 4 mux := http.NewServeMux() 5 mux.HandleFunc("/login", doLogin) 6 ... 7} 8 9func doLogin(writer http.ResponseWriter,req *http.Request){ 10 _, err := writer.Write([]byte("do login")) 11 if err != nil { 12 fmt.Printf("error : %v", err) 13 return 14 } 15}

姑且将ServeMux当作是路由器。我们使用http包下的 NewServerMux 函数创建一个新的路由器对象,进而使用它的HandleFunc(pattern,func)函数完成路由的注册。

跟进NewServerMux函数,可以看到,它通过new函数返回给我们一个ServeMux结构体。

1func NewServeMux() *ServeMux { 2 return new(ServeMux) 3}

这个ServeMux结构体长下面这样:在这个ServeMux结构体中我们就看到了这个维护pattern和func的map

1type ServeMux struct { 2 mu sync.RWMutex 3 m map[string]muxEntry 4 hosts bool // whether any patterns contain hostnames 5}

这个muxEntry长下面这样:

1type muxEntry struct { 2 h Handler 3 pattern string 4} 5 6type Handler interface { 7 ServeHTTP(ResponseWriter, *Request) 8}

image-20200627161641153

看到这里问题就来了,上面我们手动注册进路由器中的仅仅是一个有规定参数的方法,到这里怎么成了一个Handle了?我们也没有说去手动的实现Handler这个接口,也没有重写ServeHTTP函数啊, 在golang中实现一个接口不得像下面这样搞吗?**

1type Handle interface { 2 Serve(string, int, string) 3} 4 5type HandleImpl struct { 6 7} 8 9func (h HandleImpl)Serve(string, int, string){ 10 11}

带着这个疑问看下面的方法:

1 // 由于函数是一等公民,故我们将doLogin函数同普通变量一样当做入参传递进去。 2 mux.HandleFunc("/login", doLogin) 3 4 func doLogin(writer http.ResponseWriter,req *http.Request){ 5 ... 6 }

跟进去看 HandleFunc 函数的实现:

首先:HandleFunc函数的第二个参数是接收的函数的类型和doLogin函数的类型是一致的,所以doLogin能正常的传递进HandleFunc中。

其次:我们的关注点应该是下面的HandlerFunc(handler)

1// HandleFunc registers the handler function for the given pattern. 2func (mux *ServeMux) HandleFunc(pattern string, handler func(ResponseWriter, *Request)) { 3 if handler == nil { 4 panic("http: nil handler") 5 } 6 mux.Handle(pattern, HandlerFunc(handler)) 7}

跟进这个HandlerFunc(handler) 看到下图,真相就大白于天下了。golang以一种优雅的方式悄无声息的为我们完成了一次适配。这么看来上面的HandlerFunc(handler)并不是函数的调用,而是doLogin转换成自定义类型。这个自定义类型去实现了Handle接口(因为它重写了ServeHTTP函数)以闭包的形式完美的将我们的doLogin适配成了Handle类型。

image-20200625171922500

在往下看Handle方法:

第一:将pattern和handler注册进map中

第二:为了保证整个过程的并发安全,使用锁保护整个过程。

1// Handle registers the handler for the given pattern. 2// If a handler already exists for pattern, Handle panics. 3func (mux *ServeMux) Handle(pattern string, handler Handler) { 4 mux.mu.Lock() 5 defer mux.mu.Unlock() 6 7 if pattern == "" { 8 panic("http: invalid pattern") 9 } 10 if handler == nil { 11 panic("http: nil handler") 12 } 13 if _, exist := mux.m[pattern]; exist { 14 panic("http: multiple registrations for " + pattern) 15 } 16 17 if mux.m == nil { 18 mux.m = make(map[string]muxEntry) 19 } 20 mux.m[pattern] = muxEntry{h: handler, pattern: pattern} 21 22 if pattern[0] != '/' { 23 mux.hosts = true 24 } 25

启动服务

概览图:

image-20200627163736422

和java对比着看,在java一组复杂的逻辑会被封装成一个class。在golang中对应的就是一组复杂的逻辑会被封装成一个结构体。

对应HttpServer肯定也是这样,http服务器在golang的实现中有自己的结构体。它就是http包下的Server。

它有一系列描述性属性。如监听的地址、写超时时间、路由器。

1 server := &http.Server{ 2 Addr: ":8081", 3 WriteTimeout: time.Second * 2, 4 Handler: mux, 5 } 6 log.Fatal(server.ListenAndServe())

我们看它启动服务的函数:server.ListenAndServe()

实现的逻辑是使用net包下的Listen函数,获取给定地址上的tcp连接。

再将这个tcp连接封装进 tcpKeepAliveListenner 结构体中。

在将这个tcpKeepAliveListenner丢进Server的Serve函数中处理

1// ListenAndServe 会监听开发者给定网络地址上的tcp连接,当有请求到来时,会调用Serve函数去处理这个连接。 2// 它接收到所有连接都使用 TCP keep-alives相关的配置 3// 4// 如果构造Server时没有指定Addr,他就会使用默认值: “:http” 5// 6// 当Server ShutDown或者是Close,ListenAndServe总是会返回一个非nil的error。 7// 返回的这个Error是 ErrServerClosed 8func (srv *Server) ListenAndServe() error { 9 if srv.shuttingDown() { 10 return ErrServerClosed 11 } 12 addr := srv.Addr 13 if addr == "" { 14 addr = ":http" 15 } 16 // 底层借助于tcp实现 17 ln, err := net.Listen("tcp", addr) 18 if err != nil { 19 return err 20 } 21 return srv.Serve(tcpKeepAliveListener{ln.(*net.TCPListener)}) 22} 23 24// tcpKeepAliveListener会为TCP设置一个keep-alive 超时时长。 25// 它通常被 ListenAndServe 和 ListenAndServeTLS使用。 26// 它保证了已经dead的TCP最终都会消失。 27type tcpKeepAliveListener struct { 28 *net.TCPListener 29}

接着去看看Serve方法,上一个函数中获取到了一个基于tcp的Listener,从这个Listener中可以不断的获取出新的连接,下面的方法中使用无限for循环完成这件事。conn获取到后将连接封装进httpConn,为了保证不阻塞下一个连接到到来,开启新的goroutine处理这个http连接。

1func (srv *Server) Serve(l net.Listener) error { 2 // 如果有一个包裹了 srv 和 listener 的钩子函数,就执行它 3 if fn := testHookServerServe; fn != nil { 4 fn(srv, l) // call hook with unwrapped listener 5 } 6 7 // 将tcp的Listener封装进onceCloseListener,保证连接不会被关闭多次。 8 l = &onceCloseListener{Listener: l} 9 defer l.Close() 10 11 // http2相关的配置 12 if err := srv.setupHTTP2_Serve(); err != nil { 13 return err 14 } 15 16 if !srv.trackListener(&l, true) { 17 return ErrServerClosed 18 } 19 defer srv.trackListener(&l, false) 20 21 // 如果没有接收到请求睡眠多久 22 var tempDelay time.Duration // how long to sleep on accept failure 23 baseCtx := context.Background() // base is always background, per Issue 16220 24 ctx := context.WithValue(baseCtx, ServerContextKey, srv) 25 // 开启无限循环,尝试从Listenner中获取连接。 26 for { 27 rw, e := l.Accept() 28 // accpet过程中发生错屋 29 if e != nil { 30 select { 31 // 如果从server的doneChan中可以获取内容,返回Server关闭了 32 case <-srv.getDoneChan(): 33 return ErrServerClosed 34 default: 35 } 36 // 如果发生了 net.Error 并且是临时的错误就睡5毫秒,再发生错误睡眠的时间*2,上线是1s 37 if ne, ok := e.(net.Error); ok && ne.Temporary() { 38 if tempDelay == 0 { 39 tempDelay = 5 * time.Millisecond 40 } else { 41 tempDelay *= 2 42 } 43 if max := 1 * time.Second; tempDelay > max { 44 tempDelay = max 45 } 46 srv.logf("http: Accept error: %v; retrying in %v", e, tempDelay) 47 time.Sleep(tempDelay) 48 continue 49 } 50 return e 51 } 52 // 如果没有发生错误,清空睡眠的时间 53 tempDelay = 0 54 // 将接收到连接封装进httpConn 55 c := srv.newConn(rw) 56 c.setState(c.rwc, StateNew) // before Serve can return 57 // 开启一条新的协程处理这个连接 58 go c.serve(ctx) 59 } 60}

处理请求

c.serve(ctx)中就会去解析http相关的报文信息~,将http报文解析进Request结构体中。

部分代码如下:

1 // 将 server 包裹为 serverHandler 的实例,执行它的 ServeHTTP 方法,处理请求,返回响应。 2 // serverHandler 委托给 server 的 Handler 或者 DefaultServeMux(默认路由器) 3 // 来处理 "OPTIONS *" 请求。 4 serverHandler{c.server}.ServeHTTP(w, w.req) 5 6 7// serverHandler delegates to either the server's Handler or 8// DefaultServeMux and also handles "OPTIONS *" requests. 9type serverHandler struct { 10 srv *Server 11} 12 13func (sh serverHandler) ServeHTTP(rw ResponseWriter, req *Request) { 14 // 如果没有定义Handler就使用默认的 15 handler := sh.srv.Handler 16 if handler == nil { 17 handler = DefaultServeMux 18 } 19 if req.RequestURI == "*" && req.Method == "OPTIONS" { 20 handler = globalOptionsHandler{} 21 } 22 // 处理请求,返回响应。 23 handler.ServeHTTP(rw, req) 24}

image-20200625183225261

可以看到,req中包含了我们前面说的pattern,叫做RequestUri,有了它下一步就知道该回调ServeMux中的哪一个函数。

HttpClient源码阅读

DemoCode

1func main() { 2 // 创建连接池 3 // 创建客户端,绑定连接池 4 // 发送请求 5 // 读取响应 6 transport := &http.Transport{ 7 DialContext: (&net.Dialer{ 8 Timeout: 30 * time.Second, // 连接超时 9 KeepAlive: 30 * time.Second, // 长连接存活的时间 10 }).DialContext, 11 // 最大空闲连接数 12 MaxIdleConns: 100, 13 // 超过最大空闲连接数的连接会在IdleConnTimeout后被销毁 14 IdleConnTimeout: 10 * time.Second, 15 TLSHandshakeTimeout: 10 * time.Second, // tls握手超时时间 16 ExpectContinueTimeout: 1 * time.Second, // 100-continue 状态码超时时间 17 } 18 19 // 创建客户端 20 client := &http.Client{ 21 Timeout: time.Second * 10, //请求超时时间 22 Transport: transport, 23 } 24 25 // 请求数据,获得响应 26 res, err := client.Get("http://localhost:8081/login") 27 if err != nil { 28 fmt.Printf("error : %v", err) 29 return 30 } 31 defer res.Body.Close() 32 // 处理数据 33 bytes, err := ioutil.ReadAll(res.Body) 34 if err != nil { 35 fmt.Printf("error : %v", err) 36 return 37 } 38 fmt.Printf("Read from http server res:[%v]", string(bytes)) 39}

整理思路

http.Client的代码其实是很多的,全部很细的过一遍肯定也会难度,下面可能也是只能提及其中的一部分。

首先明白一件事,我们编写的HttpClient是在干什么?(虽然这个问题很傻,但是总得问一下)是在发送Http请求。

一般我们在开发的时候,更多的编写的是HttpServer的代码。是在处理Http请求, 而不是去发送Http请求,Http请求都是是前端通过ajax经由浏览器发送到后端的。

其次,Http请求实际上是建立在tcp连接之上的,所以如果我们去看http.Client肯定能找到net.Dial("tcp",adds)相关的代码。

那也就是说,我们要看看,http.Client是如何在和服务端建立连接、发送数据、接收数据的。

重要的struct

http.Client中有机几个比较重要的struct,如下

http.Client结构体中封装了和http请求相关的属性,诸如 cookie,timeout,redirect以及Transport。

1type Client struct { 2 Transport RoundTripper 3 CheckRedirect func(req *Request, via []*Request) error 4 Jar CookieJar 5 Timeout time.Duration 6}

Tranport实现了RoundTrpper接口:

1 type RoundTripper interface { 2 // 1、RoundTrip会去执行一个简单的 Http Trancation,并为requestt返回一个响应 3 // 2、RoundTrip不会尝试去解析response 4 // 3、注意:只要返回了Reponse,无论response的状态码是多少,RoundTrip返回的结果:err == nil 5 // 4、RoundTrip将请求发送出去后,如果他没有获取到response,他会返回一个非空的err。 6 // 5、同样,RoundTrip不会尝试去解析诸如重定向、认证、cookie这种更高级的协议。 7 // 6、除了消费和关闭请求体之外,RoundTrip不会修改request的其他字段 8 // 7、RoundTrip可以在一个单独的gorountine中读取request的部分字段。一直到ResponseBody关闭之前,调用者都不能取消,或者重用这个request 9 // 8、RoundTrip始终会保证关闭Body(包含在发生err时)。根据实现的不同,在RoundTrip关闭前,关闭Body这件事可能会在一个单独的goroutine中去做。这就意味着,如果调用者想将请求体用于后续的请求,必须等待知道发生Close 10 // 9、请求的URL和Header字段必须是被初始化的。 11 RoundTrip(*Request) (*Response, error) 12}

看上面RoundTrpper接口,它里面只有一个方法RoundTrip,方法的作用就是执行一次Http请求,发送Request然后获取Response。

RoundTrpper被设计成了一个支持并发的结构体。

Transport结构体如下:

1type Transport struct { 2 idleMu sync.Mutex 3 // user has requested to close all idle conns 4 wantIdle bool 5 // Transport的作用就是用来建立一个连接,这个idleConn就是Transport维护的空闲连接池。 6 idleConn map[connectMethodKey][]*persistConn // most recently used at end 7 idleConnCh map[connectMethodKey]chan *persistConn 8}

其中的connectMethodKey也是结构体:

1type connectMethodKey struct { 2 // proxy 代理的URL,当他不为空时,就会一直使用这个key 3 // scheme 协议的类型, http https 4 // addr 代理的url,也就是下游的url 5 proxy, scheme, addr string 6}

persistConn是一个具体的连接实例,包含连接的上下文。

1type persistConn struct { 2 // alt可选地指定TLS NextProto RoundTripper。 3 // 这用于今天的HTTP / 2和以后的将来的协议。 如果非零,则其余字段未使用。 4 alt RoundTripper 5 t *Transport 6 cacheKey connectMethodKey 7 conn net.Conn 8 tlsState *tls.ConnectionState 9 // 用于从conn中读取内容 10 br *bufio.Reader // from conn 11 // 用于往conn中写内容 12 bw *bufio.Writer // to conn 13 nwrite int64 // bytes written 14 // 他是个chan,roundTrip会将readLoop中的内容写入到reqch中 15 reqch chan requestAndChan 16 // 他是个chan,roundTrip会将writeLoop中的内容写到writech中 17 writech chan writeRequest 18 closech chan struct{} // closed when conn closed

另外补充一个结构体:Request,他用来描述一次http请求的实例,它定义于http包request.go, 里面封装了对Http请求相关的属性

1type Request struct { 2 Method string 3 URL *url.URL 4 Proto string // "HTTP/1.0" 5 ProtoMajor int // 1 6 ProtoMinor int // 0 7 Header Header 8 Body io.ReadCloser 9 GetBody func() (io.ReadCloser, error) 10 ContentLength int64 11 TransferEncoding []string 12 Close bool 13 Host string 14 Form url.Values 15 PostForm url.Values 16 MultipartForm *multipart.Form 17 Trailer Header 18 RemoteAddr string 19 RequestURI string 20 TLS *tls.ConnectionState 21 Cancel <-chan struct{} 22 Response *Response 23 ctx context.Context 24}

这几个结构体共同完成如下图所示http.Client的工作流程

image-20200627131720251

流程

我们想发送一次Http请求。首先我们需要构造一个Request,Request本质上是对Http协议的描述(因为大家使用的都是Http协议,所以将这个Request发送到HttpServer后,HttpServer能识别并解析它)。

1// 从这行代码开始往下看 2 res, err := client.Get("http://localhost:8081/login") 3 4// 跟进Get 5 req, err := NewRequest("GET", url, nil) 6 if err != nil { 7 return nil, err 8 } 9 return c.Do(req) 10 11// 跟进Do 12 func (c *Client) Do(req *Request) (*Response, error) { 13 return c.do(req) 14 } 15 16// 跟进do,do函数中有下面的逻辑,可以看到执行完send后已经拿到返回值了。所以我们得继续跟进send方法 17 if resp, didTimeout, err = c.send(req, deadline); err != nil 18 19// 跟进send方法,可以看到send中还有一send方法,入参分别是:request,tranpost,deadline 20// 到现在为止,我们没有看到有任何和服务端建立连接的动作发生,但是构造的req和拥有连接池的tranport已经见面了~ 21 resp, didTimeout, err = send(req, c.transport(), deadline) 22 23// 继续跟进这个send方法,看到了调用了rt的RoundTrip方法。 24// 这个rt就是我们编写HttpClient代码时创建的,绑定在http.Client上的tranport实例。 25// 这个RoundTrip方法的作用我们在上面已经说过了,最直接的作用就是:发送request 并获取response。 26 resp, err = rt.RoundTrip(req) 27

但是RoundTrip他是个定义在RoundTripper接口中的抽象方法,我们看代码肯定是要去看具体的实现嘛
这里可以使用断点调试法:在上面最后一行上打上断点,会进入到他的具体实现中。从图中可以看到具体的实现在roundtrip中。

image-20200627103402751

RoundTrip中调用的函数是我们自定义的transport的roundTrip函数, 跟进去如下:

紧接着我们需要一个conn,这个conn我们通过Transport可以获取到。conn的类型为persistConn。

1// roundTrip函数中又一个无限for循环 2for { 3 // 检查请求的上下文是否关闭了 4 select { 5 case <-ctx.Done(): 6 req.closeBody() 7 return nil, ctx.Err() 8 default: 9 } 10 11 // 对传递进来的req进行了有一层的封装,封装后的这个treq可以被roundTrip修改,所以每次重试都会新建 12 treq := &transportRequest{Request: req, trace: trace} 13 cm, err := t.connectMethodForRequest(treq) 14 if err != nil { 15 req.closeBody() 16 return nil, err 17 } 18 19 // 到这里真的执行从tranport中获取和对应主机的连接,这个连接可能是http、https、http代理、http代理的高速缓存, 但是无论如何我们都已经准备好了向这个连接发送treq 20 // 这里获取出来的连接就是我们在上文中提及的persistConn 21 pconn, err := t.getConn(treq, cm) 22 if err != nil { 23 t.setReqCanceler(req, nil) 24 req.closeBody() 25 return nil, err 26 } 27 28 var resp *Response 29 if pconn.alt != nil { 30 // HTTP/2 path. 31 t.decHostConnCount(cm.key()) // don't count cached http2 conns toward conns per host 32 t.setReqCanceler(req, nil) // not cancelable with CancelRequest 33 resp, err = pconn.alt.RoundTrip(req) 34 } else { 35 36 // 调用persistConn的roundTrip方法,发送treq并获取响应。 37 resp, err = pconn.roundTrip(treq) 38 } 39 if err == nil { 40 return resp, nil 41 } 42 if !pconn.shouldRetryRequest(req, err) { 43 // Issue 16465: return underlying net.Conn.Read error from peek, 44 // as we've historically done. 45 if e, ok := err.(transportReadFromServerError); ok { 46 err = e.err 47 } 48 return nil, err 49 } 50 testHookRoundTripRetried() 51 52 // Rewind the body if we're able to. (HTTP/2 does this itself so we only 53 // need to do it for HTTP/1.1 connections.) 54 if req.GetBody != nil && pconn.alt == nil { 55 newReq := *req 56 var err error 57 newReq.Body, err = req.GetBody() 58 if err != nil { 59 return nil, err 60 } 61 req = &newReq 62 } 63 }

整理思路:然后看上面代码中获取conn和roundTrip的实现细节。

我们需要一个conn,这个conn可以通过Transport获取到。conn的类型为persistConn。但是不管怎么样,都得先获取出 persistConn,才能进一步完成发送请求再得到服务端到响应。

然后关于这个persistConn结构体其实上面已经提及过了。重新贴在下面

1type persistConn struct { 2 // alt可选地指定TLS NextProto RoundTripper。 3 // 这用于今天的HTTP / 2和以后的将来的协议。 如果非零,则其余字段未使用。 4 alt RoundTripper 5 6 conn net.Conn 7 t *Transport 8 br *bufio.Reader // 用于从conn中读取内容 9 bw *bufio.Writer // 用于往conn中写内容 10 // 他是个chan,roundTrip会将readLoop中的内容写入到reqch中 11 reqch chan requestAndChan 12 // 他是个chan,roundTrip会将writeLoop中的内容写到writech中 13 14 nwrite int64 // bytes written 15 cacheKey connectMethodKey 16 tlsState *tls.ConnectionState 17 writech chan writeRequest 18 closech chan struct{} // closed when conn closed

跟进 t.getConn(treq, cm)代码如下:

1 // 先尝试从空闲缓冲池中取得连接 2 // 所谓的空闲缓冲池就是Tranport结构体中的: idleConn map[connectMethodKey][]*persistConn 3 // 入参位置的cm如下: 4 /* type connectMethod struct { 5 // 代理的url,如果没有代理的话,这个值为nil 6 proxyURL *url.URL 7 8 // 连接所使用的协议 http、https 9 targetScheme string 10 11 // 如果proxyURL指定了http代理或者是https代理,并且使用的协议是http而不是https。 12 // 那么下面的targetAddr就会不包含在connect method key中。 13 // 因为socket可以复用不同的targetAddr值 14 targetAddr string 15 }*/ 16 t.getIdleConn(cm); 17 18 // 空闲缓冲池有的空闲连接的话返回conn,否则进行如下的select 19 select { 20 // todo 这里我还不确定是在干什么,目前猜测是这样的:每个服务器能打开的socket句柄是有限的 21 // 每次来获取链接的时候,我们就计数+1。当整体的句柄在Host允许范围内时我们不做任何干涉~ 22 case <-t.incHostConnCount(cmKey): 23 // count below conn per host limit; proceed 24 25 // 重新尝试从空闲连接池中获取连接,因为可能有的连接使用完后被放回连接池了 26 case pc := <-t.getIdleConnCh(cm): 27 if trace != nil && trace.GotConn != nil { 28 trace.GotConn(httptrace.GotConnInfo{Conn: pc.conn, Reused: pc.isReused()}) 29 } 30 return pc, nil 31 // 请求是否被取消了 32 case <-req.Cancel: 33 return nil, errRequestCanceledConn 34 // 请求的上下文是否Done掉了 35 case <-req.Context().Done(): 36 return nil, req.Context().Err() 37 case err := <-cancelc: 38 if err == errRequestCanceled { 39 err = errRequestCanceledConn 40 } 41 return nil, err 42 } 43 44 // 开启新的gorountine新建连接一个连接 45 go func() { 46 /** 47 * 新建连接,方法底层封装了tcp client dial相关的逻辑 48 * conn, err := t.dial(ctx, "tcp", cm.addr()) 49 * 以及根据不同的targetScheme构建不同的request的逻辑。 50 */ 51 // 获取到persistConn 52 pc, err := t.dialConn(ctx, cm) 53 // 将persistConn写到chan中 54 dialc <- dialRes{pc, err} 55 }() 56 57 // 再尝试从空闲连接池中获取 58 idleConnCh := t.getIdleConnCh(cm) 59 select { 60 // 如果上面的go协程拨号成功了,这里就能取出值来 61 case v := <-dialc: 62 // Our dial finished. 63 if v.pc != nil { 64 if trace != nil && trace.GotConn != nil && v.pc.alt == nil { 65 trace.GotConn(httptrace.GotConnInfo{Conn: v.pc.conn}) 66 } 67 return v.pc, nil 68 } 69 // Our dial failed. See why to return a nicer error 70 // value. 71 // 将Host的连接-1 72 t.decHostConnCount(cmKey) 73 select { 74 ... 75

transport.dialConn

下面代码中的cm长这样

image-20200627121729925

1// dialConn是Transprot的方法 2// 入参:context上下文, connectMethod 3// 出参:persisnConn 4func (t *Transport) dialConn(ctx context.Context, cm connectMethod) (*persistConn, error) { 5 // 构建将要返回的 persistConn 6 pconn := &persistConn{ 7 t: t, 8 cacheKey: cm.key(), 9 reqch: make(chan requestAndChan, 1), 10 writech: make(chan writeRequest, 1), 11 closech: make(chan struct{}), 12 writeErrCh: make(chan error, 1), 13 writeLoopDone: make(chan struct{}), 14 } 15 trace := httptrace.ContextClientTrace(ctx) 16 wrapErr := func(err error) error { 17 if cm.proxyURL != nil { 18 // Return a typed error, per Issue 16997 19 return &net.OpError{Op: "proxyconnect", Net: "tcp", Err: err} 20 } 21 return err 22 } 23 24 // 判断cm中使用的协议是否是https 25 if cm.scheme() == "https" && t.DialTLS != nil { 26 var err error 27 pconn.conn, err = t.DialTLS("tcp", cm.addr()) 28 if err != nil { 29 return nil, wrapErr(err) 30 } 31 if pconn.conn == nil { 32 return nil, wrapErr(errors.New("net/http: Transport.DialTLS returned (nil, nil)")) 33 } 34 if tc, ok := pconn.conn.(*tls.Conn); ok { 35 // Handshake here, in case DialTLS didn't. TLSNextProto below 36 // depends on it for knowing the connection state. 37 if trace != nil && trace.TLSHandshakeStart != nil { 38 trace.TLSHandshakeStart() 39 } 40 if err := tc.Handshake(); err != nil { 41 go pconn.conn.Close() 42 if trace != nil && trace.TLSHandshakeDone != nil { 43 trace.TLSHandshakeDone(tls.ConnectionState{}, err) 44 } 45 return nil, err 46 } 47 cs := tc.ConnectionState() 48 if trace != nil && trace.TLSHandshakeDone != nil { 49 trace.TLSHandshakeDone(cs, nil) 50 } 51 pconn.tlsState = &cs 52 } 53 } else { 54 // 如果不是https协议就来到这里,使用tcp向httpserver拨号,获取一个tcp连接。 55 conn, err := t.dial(ctx, "tcp", cm.addr()) 56 if err != nil { 57 return nil, wrapErr(err) 58 } 59 // 将获取到tcp连接交给我们的persistConn维护 60 pconn.conn = conn 61 62 // 处理https相关逻辑 63 if cm.scheme() == "https" { 64 var firstTLSHost string 65 if firstTLSHost, _, err = net.SplitHostPort(cm.addr()); err != nil { 66 return nil, wrapErr(err) 67 } 68 if err = pconn.addTLS(firstTLSHost, trace); err != nil { 69 return nil, wrapErr(err) 70 } 71 } 72 } 73 74 // Proxy setup. 75 switch { 76 // 如果代理URL为空,不做任何处理 77 case cm.proxyURL == nil: 78 // Do nothing. Not using a proxy. 79 // 80 case cm.proxyURL.Scheme == "socks5": 81 conn := pconn.conn 82 d := socksNewDialer("tcp", conn.RemoteAddr().String()) 83 if u := cm.proxyURL.User; u != nil { 84 auth := &socksUsernamePassword{ 85 Username: u.Username(), 86 } 87 auth.Password, _ = u.Password() 88 d.AuthMethods = []socksAuthMethod{ 89 socksAuthMethodNotRequired, 90 socksAuthMethodUsernamePassword, 91 } 92 d.Authenticate = auth.Authenticate 93 } 94 if _, err := d.DialWithConn(ctx, conn, "tcp", cm.targetAddr); err != nil { 95 conn.Close() 96 return nil, err 97 } 98 case cm.targetScheme == "http": 99 pconn.isProxy = true 100 if pa := cm.proxyAuth(); pa != "" { 101 pconn.mutateHeaderFunc = func(h Header) { 102 h.Set("Proxy-Authorization", pa) 103 } 104 } 105 case cm.targetScheme == "https": 106 conn := pconn.conn 107 hdr := t.ProxyConnectHeader 108 if hdr == nil { 109 hdr = make(Header) 110 } 111 connectReq := &Request{ 112 Method: "CONNECT", 113 URL: &url.URL{Opaque: cm.targetAddr}, 114 Host: cm.targetAddr, 115 Header: hdr, 116 } 117 if pa := cm.proxyAuth(); pa != "" { 118 connectReq.Header.Set("Proxy-Authorization", pa) 119 } 120 connectReq.Write(conn) 121 122 // Read response. 123 // Okay to use and discard buffered reader here, because 124 // TLS server will not speak until spoken to. 125 br := bufio.NewReader(conn) 126 resp, err := ReadResponse(br, connectReq) 127 if err != nil { 128 conn.Close() 129 return nil, err 130 } 131 if resp.StatusCode != 200 { 132 f := strings.SplitN(resp.Status, " ", 2) 133 conn.Close() 134 if len(f) < 2 { 135 return nil, errors.New("unknown status code") 136 } 137 return nil, errors.New(f[1]) 138 } 139 } 140 141 if cm.proxyURL != nil && cm.targetScheme == "https" { 142 if err := pconn.addTLS(cm.tlsHost(), trace); err != nil { 143 return nil, err 144 } 145 } 146 147 if s := pconn.tlsState; s != nil && s.NegotiatedProtocolIsMutual && s.NegotiatedProtocol != "" { 148 if next, ok := t.TLSNextProto[s.NegotiatedProtocol]; ok { 149 return &persistConn{alt: next(cm.targetAddr, pconn.conn.(*tls.Conn))}, nil 150 } 151 } 152 153 if t.MaxConnsPerHost > 0 { 154 pconn.conn = &connCloseListener{Conn: pconn.conn, t: t, cmKey: pconn.cacheKey} 155 } 156 157 // 初始化persistConn的bufferReader和bufferWriter 158 pconn.br = bufio.NewReader(pconn) // 可以从上面给pconn维护的tcpConn中读数据 159 pconn.bw = bufio.NewWriter(persistConnWriter{pconn})// 可以往上面pconn维护的tcpConn中写数据 160 161 // 新开启两条和persistConn相关的go协程。 162 go pconn.readLoop() 163 go pconn.writeLoop() 164 return pconn, nil 165}

上面的两条goroutine 和 br bw共同完成如下图的流程

image-20200627131859112

发送请求

发送req的逻辑在http包的下的tranport包中的func (t *Transport) roundTrip(req *Request) (*Response, error) {}函数中。

如下:

1 // 发送treq 2 resp, err = pconn.roundTrip(treq) 3 4 // 跟进roundTrip 5 // 可以看到他将一个writeRequest结构体类型的实例写入了writech中 6 // 而这个writech会被上图中的writeLoop消费,借助bufferWriter写入tcp连接中,完成往服务端数据的发送。 7 pc.writech <- writeRequest{req, writeErrCh, continueCh}
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )