拉模式和推模式
拉模式
1、数据更新频率低,则大多数请求是无效的 2、在线用户量多,则服务端的查询负载高 3、定时轮询拉取,实时性低
推模式
1、仅在数据更新时才需要推送 2、需要维护大量的在线长连接 3、数据更新后可以立即推送
基于webSocket推送
1、浏览器支持的socket编程,轻松维持服务端长连接 2、基于TCP可靠传输之上的协议,无需开发者关心通讯细节 3、提供了高度抽象的编程接口,业务开发成本较低
webSocket协议与交互
通讯流程
客户端->upgrade->服务端 客户端<-switching<-服务端 客户端->message->服务端 客户端<-message<-服务端
实现http服务端
1、webSocket是http协议upgrade而来 2、使用http标准库快速实现空接口:/ws
webSocket握手
1、使用webSocket.Upgrader完成协议握手,得到webSocket长连接 2、操作webSocket api,读取客户端消息,然后原样发送回去
封装webSocket
缺乏工程化设计
1、其他代码模块,无法直接操作webSocket连接 2、webSocket连接非线程安全,并发读/写需要同步手段
隐藏细节,封装api
1、封装Connection结构,隐藏webSocket底层连接 2、封装Connection的api,提供Send/Read/Close等线程安全接口
api原理(channel是线程安全的)
1、SendMessage将消息投递到out channel 2、ReadMessage从in channel读取消息
内部原理
1、启动读协程,循环读取webSocket,将消息投递到in channel 2、启动写协程,循环读取out channel,将消息写给webSocket
1// server.go 2package main 3 4import ( 5 "net/http" 6 "github.com/gorilla/websocket" 7 "./impl" 8 "time" 9) 10 11var ( 12 upgrader = websocket.Upgrader{ 13 //允许跨域 14 CheckOrigin: func(r *http.Request) bool { 15 return true 16 }, 17 } 18) 19 20func wsHandler(w http.ResponseWriter, r *http.Request) { 21 var ( 22 wsConn *websocket.Conn 23 err error 24 conn *impl.Connection 25 data []byte 26 ) 27 28 //Upgrade:websocket 29 if wsConn, err = upgrader.Upgrade(w, r, nil); err != nil { 30 return 31 } 32 if conn, err = impl.InitConnection(wsConn); err != nil { 33 goto ERR 34 } 35 36 go func() { 37 var ( 38 err error 39 ) 40 for { 41 if err =conn.WriteMessage([]byte("heartbeat")); err != nil { 42 return 43 } 44 time.Sleep(1 * time.Second) 45 } 46 }() 47 48 for { 49 if data, err = conn.ReadMessage(); err != nil { 50 goto ERR 51 } 52 if err = conn.WriteMessage(data); err != nil { 53 goto ERR 54 } 55 56 } 57 58 ERR: 59 //关闭连接 60 conn.Close() 61} 62 63func main() { 64 //http:localhost:7777/ws 65 http.HandleFunc("/ws", wsHandler) 66 http.ListenAndServe("0.0.0.0:7777", nil) 67} 68 69 70// connection.go 71package impl 72 73import ( 74 "github.com/gorilla/websocket" 75 "sync" 76 "github.com/influxdata/platform/kit/errors" 77) 78 79var once sync.Once 80 81type Connection struct { 82 wsConn *websocket.Conn 83 inChan chan []byte 84 outChan chan []byte 85 closeChan chan byte 86 isClosed bool 87 mutex sync.Mutex 88} 89 90func InitConnection(wsConn *websocket.Conn) (conn *Connection, err error) { 91 conn = &Connection{ 92 wsConn:wsConn, 93 inChan:make(chan []byte, 1000), 94 outChan:make(chan []byte, 1000), 95 closeChan:make(chan byte, 1), 96 } 97 98 //启动读协程 99 go conn.readLoop() 100 101 //启动写协程 102 go conn.writeLoop() 103 104 return 105} 106 107//API 108func (conn *Connection) ReadMessage() (data []byte, err error) { 109 select { 110 case data = <- conn.inChan: 111 case <- conn.closeChan: 112 err = errors.New("connection is closed") 113 } 114 return 115} 116 117func (conn *Connection) WriteMessage(data []byte) (err error) { 118 select { 119 case conn.outChan <- data: 120 case <- conn.closeChan: 121 err = errors.New("connection is closed") 122 } 123 return 124} 125 126func (conn *Connection) Close() { 127 // 线程安全的close,可重入 128 conn.wsConn.Close() 129 conn.mutex.Lock() 130 if !conn.isClosed { 131 close(conn.closeChan) 132 conn.isClosed = true 133 } 134 conn.mutex.Unlock() 135} 136 137//内部实现 138func (conn *Connection) readLoop() { 139 var ( 140 data []byte 141 err error 142 ) 143 for { 144 if _, data, err = conn.wsConn.ReadMessage(); err != nil { 145 goto ERR 146 } 147 148 //阻塞在这里,等待inChan有空位置 149 //但是如果writeLoop连接关闭了,这边无法得知 150 //conn.inChan <- data 151 152 select { 153 case conn.inChan <- data: 154 case <-conn.closeChan: 155 //closeChan关闭的时候,会进入此分支 156 goto ERR 157 } 158 } 159 ERR: 160 conn.Close() 161} 162 163func (conn *Connection) writeLoop() { 164 var ( 165 data []byte 166 err error 167 ) 168 for { 169 select { 170 case data = <- conn.outChan: 171 case <- conn.closeChan: 172 goto ERR 173 174 } 175 176 if err = conn.wsConn.WriteMessage(websocket.TextMessage, data); err != nil { 177 goto ERR 178 } 179 conn.outChan <- data 180 } 181 ERR: 182 conn.Close() 183}