Go实现基于WebSocket的弹幕服务

拉模式和推模式

拉模式

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}
点赞
收藏

评论区

加载中...

相关推荐

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

java将前端的json数组字符串转换为列表

记录下在前端通过ajax提交了一个json数组的字符串,在后端如何转换为列表。前端数据转化与请求varcontracts{id:'1',name:'yanggb合同1'},{id:'2',name:'yanggb合同2'},{id:'3',name:'yang