You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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. 错误处理

  • 向客户端推送消息失败时,自动从主题中移除该连接(通常因客户端已断开)。
  • 对无效命令返回明确的错误提示。

使用示例

  1. 启动服务后,客户端通过WebSocket连接到ws://x.x.x.x:8080/gs-guide-websocket。
  2. 发送subscribe:sports订阅体育主题。
  3. 另一个客户端同样订阅sports主题。
  4. 其中一个客户端发送publish:sports:世界杯小组赛结果出炉!,所有订阅该主题的客户端都会收到这条消息。

内容的提问来源于stack exchange,提问作者Rafael Souza

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 04:01:33