如何在Golang中基于Gorilla/WebSocket实现主题式消息收发
基于Gorilla/WebSocket实现主题订阅与推送功能
要实现主题(Topic)功能,核心是维护主题与客户端连接的映射关系,并提供订阅、取消订阅和消息广播的逻辑。以下是完整的实现方案:
核心实现代码
package main import ( "fmt" "log" "net/http" "strings" "sync" "github.com/gorilla/websocket" ) var upgrader = websocket.Upgrader{ ReadBufferSize: 1024, WriteBufferSize: 1024, // 允许跨域(根据实际需求调整) CheckOrigin: func(r *http.Request) bool { return true }, } // TopicManager 管理所有主题和订阅的客户端 type TopicManager struct { mu sync.Mutex // 保证并发安全的互斥锁 topics map[string]map[*websocket.Conn]bool // 主题 -> 客户端连接集合 } // NewTopicManager 创建一个新的主题管理器 func NewTopicManager() *TopicManager { return &TopicManager{ topics: make(map[string]map[*websocket.Conn]bool), } } // Subscribe 让客户端订阅指定主题 func (tm *TopicManager) Subscribe(topic string, conn *websocket.Conn) { tm.mu.Lock() defer tm.mu.Unlock() // 如果主题不存在,创建新的客户端集合 if _, ok := tm.topics[topic]; !ok { tm.topics[topic] = make(map[*websocket.Conn]bool) } tm.topics[topic][conn] = true log.Printf("客户端订阅主题: %s,当前订阅数: %d", topic, len(tm.topics[topic])) } // Unsubscribe 让客户端取消订阅指定主题 func (tm *TopicManager) Unsubscribe(topic string, conn *websocket.Conn) { tm.mu.Lock() defer tm.mu.Unlock() if clients, ok := tm.topics[topic]; ok { delete(clients, conn) log.Printf("客户端取消订阅主题: %s,当前订阅数: %d", topic, len(clients)) // 如果主题没有订阅者了,删除该主题以节省资源 if len(clients) == 0 { delete(tm.topics, topic) } } } // Broadcast 向指定主题的所有订阅者推送消息 func (tm *TopicManager) Broadcast(topic string, message []byte) { tm.mu.Lock() defer tm.mu.Unlock() if clients, ok := tm.topics[topic]; ok { // 遍历所有订阅客户端发送消息 for conn := range clients { err := conn.WriteMessage(websocket.TextMessage, message) if err != nil { log.Printf("向客户端推送消息失败: %v", err) // 发送失败时,自动移除该客户端(可能已断开连接) delete(clients, conn) } } } } // RemoveClient 移除客户端的所有订阅 func (tm *TopicManager) RemoveClient(conn *websocket.Conn) { tm.mu.Lock() defer tm.mu.Unlock() // 遍历所有主题,移除该客户端 for topic, clients := range tm.topics { delete(clients, conn) // 如果主题没有订阅者,删除主题 if len(clients) == 0 { delete(tm.topics, topic) } } log.Println("客户端已断开,移除所有订阅") } var topicManager = NewTopicManager() func main() { http.HandleFunc("/gs-guide-websocket", wsHandler) log.Fatal(http.ListenAndServe("x.x.x.x:8080", nil)) } func wsHandler(w http.ResponseWriter, r *http.Request) { conn, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Println(err) return } defer func() { // 客户端断开时,移除所有订阅并关闭连接 topicManager.RemoveClient(conn) conn.Close() }() for { messageType, message, err := conn.ReadMessage() if err != nil { log.Println(err) return } msgStr := string(message) log.Printf("收到客户端消息: %s", msgStr) // 解析客户端命令 switch { // 订阅命令格式: subscribe:topicName case strings.HasPrefix(msgStr, "subscribe:"): topic := strings.TrimPrefix(msgStr, "subscribe:") topicManager.Subscribe(topic, conn) // 向客户端发送订阅成功通知 conn.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf("已订阅主题: %s", topic))) // 取消订阅命令格式: unsubscribe:topicName case strings.HasPrefix(msgStr, "unsubscribe:"): topic := strings.TrimPrefix(msgStr, "unsubscribe:") topicManager.Unsubscribe(topic, conn) conn.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf("已取消订阅主题: %s", topic))) // 发布消息命令格式: publish:topicName:messageContent case strings.HasPrefix(msgStr, "publish:"): parts := strings.SplitN(msgStr, ":", 3) if len(parts) != 3 { conn.WriteMessage(websocket.TextMessage, []byte("无效的发布命令,格式应为: publish:topicName:message")) continue } topic := parts[1] content := parts[2] topicManager.Broadcast(topic, []byte(content)) conn.WriteMessage(websocket.TextMessage, []byte(fmt.Sprintf("已向主题 %s 推送消息", topic))) default: // 未知命令,返回提示 conn.WriteMessage(websocket.TextMessage, []byte("未知命令,支持的命令有: subscribe:topic, unsubscribe:topic, publish:topic:message")) } } }
关键功能说明
1. 主题管理器(TopicManager)
- 使用
sync.Mutex保证并发安全:多个WebSocket连接会同时操作主题映射,加锁可避免竞态条件。 topics字段是双层映射:第一层为主题名称,第二层为该主题下的所有客户端连接集合。- 封装
Subscribe、Unsubscribe、Broadcast、RemoveClient四个核心方法,统一处理主题操作逻辑。
2. 客户端命令解析
客户端通过发送特定格式的文本消息执行操作:
- 订阅主题:
subscribe:news - 取消订阅:
unsubscribe:news - 发布消息:
publish:news:今天有新的技术文章发布!
3. 连接清理
当客户端断开连接时,通过defer调用RemoveClient方法,自动移除该客户端在所有主题中的订阅,避免无效连接占用资源。
4. 错误处理
- 向客户端推送消息失败时,自动从主题中移除该连接(通常因客户端已断开)。
- 对无效命令返回明确的错误提示。
使用示例
- 启动服务后,客户端通过WebSocket连接到
ws://x.x.x.x:8080/gs-guide-websocket。 - 发送
subscribe:sports订阅体育主题。 - 另一个客户端同样订阅
sports主题。 - 其中一个客户端发送
publish:sports:世界杯小组赛结果出炉!,所有订阅该主题的客户端都会收到这条消息。
内容的提问来源于stack exchange,提问作者Rafael Souza
相关产品推荐
相关产品推荐

