Golang中WebSocket实时数据内存存储的并发问题解决
问题分析与解决方案
问题场景
我通过WebSocket对接交易平台API获取100个标的数据,每次onMessage回调会收到[]map[string]interface{}格式的部分标的数据。我维护了livestore数组存储实时数据,用sync.RWMutex控制并发读写,但遇到以下问题:
- WebSocket数据接收过快,写锁长期无法释放,导致
getsymdata函数无法读取新数据 - 尝试用无缓冲channel优化后问题依旧
- 每秒调用100次
getsymdata获取不同标的价格,始终拿不到新数据
相关代码:
// websocket on message函数 onMessage(message []map[string]interface{}){ lock.Lock(); defer lock.Unlock(); storedata(message []map[string]interface{}) } var livestore []map[string]interface{} func storedata(livedata []map[string]interface{}){ for _,m:= range livedata{ for_,l := range livestore{ if(m["ts"]==l["ts"]){ m["name"]==l["name"] m["price"]==l["price"] m["open"]==l["open"] ....... } } } } func getsymdata(name string) map[string]interface{}{ lock.Rlock(); defer lock.RUnlock(); d:= make(map[string]interface{}) for_,l := range livestore{ if(l["name"]==name) { d["name"]=l["name"] d["price"]=l["price"] d["open"]=l["open"] ....... } return d } }
核心问题分析
- 存储结构效率极低:用数组
livestore存标的,每次更新/查询都要遍历全数组,时间复杂度O(n)。消息频繁时,写操作会长期占用锁,直接堵死读请求。 storedata逻辑全错:- 赋值用了
==而非=,根本没更新数据 - 逻辑搞反:应该用新数据更新旧数据,不是把旧数据赋值给新数据
- 嵌套循环导致每次更新耗时随标的数量线性增长,写锁持有时间拉满
- 赋值用了
getsymdata逻辑错误:循环里的return d位置不对,第一次循环就直接返回空map,根本没遍历找目标标的- 锁的使用不合理:写锁持有期间执行大量耗时遍历,读请求完全抢不到锁
解决方案
1. 替换存储结构:用map代替数组
把livestore改成map[string]map[string]interface{},key用标的唯一标识(比如name或ts),更新和查询的时间复杂度直接降到O(1),大幅缩短锁持有时间。
2. 修正核心逻辑
- 修复
storedata的赋值符号和更新逻辑 - 调整
getsymdata的返回位置,确保找到目标标的再返回
3. 优化锁的使用
写锁只在更新map的瞬间持有,避免在锁内做额外耗时操作;读锁正常使用,保证并发读安全。
优化后的代码示例
import "sync" // 用map替代数组,key为标的唯一标识(这里用name) var livestore = make(map[string]map[string]interface{}) var lock sync.RWMutex // WebSocket消息回调 func onMessage(message []map[string]interface{}) { lock.Lock() defer lock.Unlock() storeData(message) } func storeData(livedata []map[string]interface{}) { for _, newData := range livedata { // 提取标的唯一标识,这里假设name是唯一的 name, ok := newData["name"].(string) if !ok { continue // 数据格式错误,直接跳过 } // 取出已有数据,不存在则初始化 existingData, exists := livestore[name] if !exists { existingData = make(map[string]interface{}) livestore[name] = existingData } // 用新数据更新字段(按需覆盖或合并) existingData["name"] = newData["name"] existingData["price"] = newData["price"] existingData["open"] = newData["open"] existingData["ts"] = newData["ts"] // 其他字段同理补充 } } func getSymData(name string) map[string]interface{} { lock.RLock() defer lock.RUnlock() // 直接从map中取目标数据,O(1)操作 data, exists := livestore[name] if !exists { return make(map[string]interface{}) // 没找到返回空map } // 返回数据副本,避免外部修改内部存储的数据 copyData := make(map[string]interface{}) for k, v := range data { copyData[k] = v } return copyData }
额外优化建议
- 如果WebSocket消息量极大,可以加带缓冲的channel做异步处理:把收到的消息先丢进channel,用单独goroutine消费并更新livestore,让
onMessage快速返回,避免阻塞WebSocket接收。 - 高并发场景下,可考虑用
sync.Map替代map+RWMutex,读多写少的场景性能更优。
内容的提问来源于stack exchange,提问作者Suyog Patil
相关产品推荐
相关产品推荐

