Skip to content

聊天室系统实战

前两篇我们打下了协议和 API 的基础。本篇进入真正的实战:用 gorilla/websocket 实现一个多人在线聊天室。我们会从需求分析入手,设计 Hub/Client 架构,定义 JSON 消息协议,用 channel 通信避免锁竞争,实现完整的加入、离开、广播、心跳逻辑,最后配套一个 HTML + JavaScript 前端页面。学完本篇,你将掌握 WebSocket 应用的经典架构模式。

一、需求分析:多人在线聊天室

我们要实现的聊天室具备以下能力:

1. 用户加入/离开

用户连接后需要「报到」:给自己起个昵称,告诉服务器「我来了」。服务器把「某某加入」的通知广播给所有人,并维护一份在线用户列表。用户断开连接时,服务器广播「某某离开」并更新列表。

2. 广播消息

任何一个用户发出的聊天消息,服务器都要转发给当前所有在线用户。这是聊天室的核心功能,也是「广播」模式的典型场景。

3. 在线用户列表

用户可以查询当前在线的所有人。新人加入或有人离开时,列表要实时更新到所有客户端。

4. 非功能性需求

  • 心跳保活:长时间不发消息的连接不能被中间设备断开。
  • 优雅关闭:用户离开要通知其他人,不能「悄悄消失」。
  • 并发安全:多个用户同时发消息不能出错。
  • 可扩展:架构要为后续私聊、群组等功能留好扩展点。

二、架构设计

聊天室的核心难点是「一个用户的消息要安全地发给其他所有用户」。直接让每个连接互相调用 WriteMessage 会陷入锁的地狱。gorilla 官方推荐的架构是用一个中心化的 Hub 管理所有连接,连接之间不直接通信,全部通过 Hub 中转。

1. Hub:管理所有连接和广播

Hub 是聊天室的「中枢」,负责:

  • 维护所有活跃 Client 的集合。
  • 处理 Client 的注册(register)和注销(unregister)。
  • 接收 Client 发来的消息,广播给所有 Client。

Hub 运行在单个 goroutine 中,所有状态变更通过 channel 传递,完全不需要锁。这是 Go 并发编程的精髓:「不要通过共享内存通信,而要通过通信共享内存」。

2. Client:单个客户端连接

Client 封装一个 WebSocket 连接及其相关信息:

  • 持有 *websocket.Conn
  • 持有发送缓冲 channel(send)。
  • 启动两个 goroutine:读 goroutine 持续读取用户消息,写 goroutine 持续从 send 取消息发出去。

Client 不直接写连接,而是把要发的消息塞进 send channel,由写 goroutine 统一写出。这样既保证了写操作的串行化(满足 gorilla 的并发约束),又解耦了「消息产生」和「消息发送」。

3. 消息类型定义

我们用 JSON 文本帧传输消息,定义统一的消息结构,包含类型字段区分不同业务:

  • chat:聊天消息
  • join:加入通知
  • leave:离开通知
  • system:系统消息(如在线列表更新)

4. 数据流图

[Client A] --读--> Hub --广播--> [Client A, B, C 的 send channel]
[Client B] --读--> Hub --广播--> ...
                              [Client C 写 goroutine] --写--> [Client C]

每个 Client 的读 goroutine 把消息发给 Hub 的 broadcast channel;Hub 把消息复制到每个 Client 的 send channel;每个 Client 的写 goroutine 从 send 取出并写出。整个链路单向、无锁、清晰。

三、Hub 实现

Hub 的核心是一个 select 循环,监听三类事件:注册、注销、广播。

go
package main

import (
	"encoding/json"
)

// Hub 管理所有连接和广播
type Hub struct {
	clients    map[*Client]bool // 所有活跃客户端
	broadcast  chan []byte      // 广播消息通道
	register   chan *Client     // 注册通道
	unregister chan *Client     // 注销通道
}

func newHub() *Hub {
	return &Hub{
		clients:    make(map[*Client]bool),
		broadcast:  make(chan []byte, 256),
		register:   make(chan *Client),
		unregister: make(chan *Client),
	}
}

// Hub 的主循环,运行在单个 goroutine 中
func (h *Hub) run() {
	for {
		select {
		case client := <-h.register:
			// 注册新客户端
			h.clients[client] = true
		case client := <-h.unregister:
			// 注销客户端
			if _, ok := h.clients[client]; ok {
				delete(h.clients, client)
				close(client.send)
			}
		case message := <-h.broadcast:
			// 广播消息给所有客户端
			for client := range h.clients {
				select {
				case client.send <- message:
				default:
					// 发送缓冲满,认为客户端卡死,踢掉
					delete(h.clients, client)
					close(client.send)
				}
			}
		}
	}
}

// 广播一条 JSON 消息
func (h *Hub) broadcastMessage(msg map[string]interface{}) {
	data, _ := json.Marshal(msg)
	h.broadcast <- data
}

这段代码体现了几个关键设计:

  • clients 这个 map 只在 Hub 的主 goroutine 中访问,不需要锁
  • 注销时 close(client.send),写 goroutine 感知到 channel 关闭后退出。
  • 广播时用 select + default 实现「非阻塞发送」:如果某个客户端的 send 缓冲满了(说明它消费太慢),直接踢掉,避免拖慢整个 Hub。

四、Client 实现

Client 的核心是读、写两个 goroutine 加上心跳。

go
package main

import (
	"time"
	"github.com/gorilla/websocket"
)

const (
	writeWait      = 10 * time.Second
	pongWait       = 60 * time.Second
	pingPeriod     = 50 * time.Second
	maxMessageSize = 4096
)

// Client 封装一个 WebSocket 连接
type Client struct {
	hub  *Hub
	conn *websocket.Conn
	send chan []byte
	name string
}

// 读 goroutine:持续读取消息并发给 Hub
func (c *Client) readPump() {
	defer func() {
		c.hub.unregister <- c
		c.conn.Close()
	}()

	c.conn.SetReadLimit(maxMessageSize)
	c.conn.SetReadDeadline(time.Now().Add(pongWait))
	c.conn.SetPongHandler(func(string) error {
		c.conn.SetReadDeadline(time.Now().Add(pongWait))
		return nil
	})

	for {
		_, message, err := c.conn.ReadMessage()
		if err != nil {
			break
		}
		// 把读到的消息广播出去
		c.hub.broadcast <- message
	}
}

// 写 goroutine:持续从 send 取消息写出,并发 Ping
func (c *Client) writePump() {
	ticker := time.NewTicker(pingPeriod)
	defer func() {
		ticker.Stop()
		c.conn.Close()
	}()

	for {
		select {
		case message, ok := <-c.send:
			c.conn.SetWriteDeadline(time.Now().Add(writeWait))
			if !ok {
				// channel 已关闭,发送关闭帧
				c.conn.WriteMessage(websocket.CloseMessage, []byte{})
				return
			}
			if err := c.conn.WriteMessage(websocket.TextMessage, message); err != nil {
				return
			}
		case <-ticker.C:
			// 定期发 Ping
			c.conn.SetWriteDeadline(time.Now().Add(writeWait))
			if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
				return
			}
		}
	}
}

1. 读 goroutine:持续读取消息

readPump 是一个死循环,不断 ReadMessage。读到的原始消息直接塞进 hub.broadcast。退出时(连接出错或关闭)通过 hub.unregister 通知 Hub 把自己移除,并关闭底层连接。

注意 SetPongHandler:客户端回 Pong 时刷新读超时,实现心跳续期。

2. 写 goroutine:持续发送消息

writePumpselect 同时监听 send channel 和心跳 ticker。从 send 取到消息就写出,到点就发 Ping。如果 send 被关闭(Hub 注销时),发送一个关闭帧后退出。

3. 乒乓心跳保活

心跳机制分两端:

  • 服务端:每 pingPeriod(50 秒)发一个 Ping,期望客户端在 pongWait(60 秒)内回 Pong。如果 60 秒没收到任何帧(含 Pong),读超时触发,readPump 退出。
  • 客户端:浏览器原生 WebSocket 会自动回 Pong,无需手动处理。

这样即使客户端长时间不发业务消息,连接也能靠 Ping/Pong 保持活跃。

五、消息协议设计

为了让前后端能对话,需要约定消息格式。我们用 JSON 文本帧,统一结构如下:

1. JSON 消息格式

json
{
  "type": "chat",
  "from": "张三",
  "text": "大家好",
  "time": "2025-01-01 12:00:00",
  "users": ["张三", "李四"]
}
  • type:消息类型,决定前端如何渲染。
  • from:消息发送者。
  • text:消息内容。
  • time:时间戳。
  • users:在线用户列表(仅 system 消息携带)。

2. 消息类型:chat、join、leave、system

type触发场景携带内容
chat用户发送聊天消息from、text、time
join新用户加入from(新人)、time
leave用户离开from(离开者)、time
system系统更新(如用户列表)users、time

下面是 Go 端的消息结构定义:

go
package main

import "time"

// 消息类型常量
const (
	TypeChat   = "chat"
	TypeJoin   = "join"
	TypeLeave  = "leave"
	TypeSystem = "system"
)

// Message 统一消息结构
type Message struct {
	Type string   `json:"type"`
	From string   `json:"from"`
	Text string   `json:"text"`
	Time string   `json:"time"`
	Users []string `json:"users,omitempty"`
}

// 构造一条消息
func newMessage(msgType, from, text string) Message {
	return Message{
		Type: msgType,
		From: from,
		Text: text,
		Time: time.Now().Format("2006-01-02 15:04:05"),
	}
}

六、完整代码实现:聊天室服务端

把 Hub、Client、消息协议组合起来,加上 HTTP 服务和连接处理,就是一个完整的聊天室服务端。

go
package main

import (
	"encoding/json"
	"fmt"
	"log"
	"net/http"
	"time"
	"github.com/gorilla/websocket"
)

const (
	writeWait      = 10 * time.Second
	pongWait       = 60 * time.Second
	pingPeriod     = 50 * time.Second
	maxMessageSize = 4096
)

var upgrader = websocket.Upgrader{
	ReadBufferSize:  1024,
	WriteBufferSize: 1024,
	CheckOrigin:     func(r *http.Request) bool { return true },
}

// ---- 消息协议 ----
type Message struct {
	Type  string   `json:"type"`
	From  string   `json:"from"`
	Text  string   `json:"text"`
	Time  string   `json:"time"`
	Users []string `json:"users,omitempty"`
}

func newMessage(t, from, text string) Message {
	return Message{Type: t, From: from, Text: text,
		Time: time.Now().Format("2006-01-02 15:04:05")}
}

// ---- Hub ----
type Hub struct {
	clients    map[*Client]bool
	broadcast  chan []byte
	register   chan *Client
	unregister chan *Client
}

func newHub() *Hub {
	return &Hub{
		clients:    make(map[*Client]bool),
		broadcast:  make(chan []byte, 256),
		register:   make(chan *Client),
		unregister: make(chan *Client),
	}
}

func (h *Hub) run() {
	for {
		select {
		case c := <-h.register:
			h.clients[c] = true
			h.notifyUserList()
		case c := <-h.unregister:
			if _, ok := h.clients[c]; ok {
				delete(h.clients, c)
				close(c.send)
				h.notifyUserList()
			}
		case msg := <-h.broadcast:
			for c := range h.clients {
				select {
				case c.send <- msg:
				default:
					delete(h.clients, c)
					close(c.send)
				}
			}
		}
	}
}

// 广播当前在线用户列表
func (h *Hub) notifyUserList() {
	users := make([]string, 0, len(h.clients))
	for c := range h.clients {
		users = append(users, c.name)
	}
	m := Message{Type: "system", Users: users, Time: time.Now().Format("2006-01-02 15:04:05")}
	data, _ := json.Marshal(m)
	for c := range h.clients {
		select {
		case c.send <- data:
		default:
		}
	}
}

// 广播一条消息
func (h *Hub) broadcastMsg(m Message) {
	data, _ := json.Marshal(m)
	h.broadcast <- data
}

// ---- Client ----
type Client struct {
	hub  *Hub
	conn *websocket.Conn
	send chan []byte
	name string
}

func (c *Client) readPump() {
	defer func() {
		// 离开时广播 leave
		c.hub.broadcastMsg(newMessage("leave", c.name, "离开了聊天室"))
		c.hub.unregister <- c
		c.conn.Close()
	}()

	c.conn.SetReadLimit(maxMessageSize)
	c.conn.SetReadDeadline(time.Now().Add(pongWait))
	c.conn.SetPongHandler(func(string) error {
		c.conn.SetReadDeadline(time.Now().Add(pongWait))
		return nil
	})

	for {
		_, data, err := c.conn.ReadMessage()
		if err != nil {
			break
		}
		// 解析客户端发来的消息,补全 from 和 time
		var m Message
		if err := json.Unmarshal(data, &m); err != nil {
			continue
		}
		m.From = c.name
		m.Type = "chat"
		m.Time = time.Now().Format("2006-01-02 15:04:05")
		c.hub.broadcastMsg(m)
	}
}

func (c *Client) writePump() {
	ticker := time.NewTicker(pingPeriod)
	defer func() {
		ticker.Stop()
		c.conn.Close()
	}()
	for {
		select {
		case msg, ok := <-c.send:
			c.conn.SetWriteDeadline(time.Now().Add(writeWait))
			if !ok {
				c.conn.WriteMessage(websocket.CloseMessage, []byte{})
				return
			}
			if err := c.conn.WriteMessage(websocket.TextMessage, msg); err != nil {
				return
			}
		case <-ticker.C:
			c.conn.SetWriteDeadline(time.Now().Add(writeWait))
			if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
				return
			}
		}
	}
}

// ---- HTTP 服务 ----
func serveWs(hub *Hub, w http.ResponseWriter, r *http.Request) {
	conn, err := upgrader.Upgrade(w, r, nil)
	if err != nil {
		log.Println("升级失败:", err)
		return
	}

	name := r.URL.Query().Get("name")
	if name == "" {
		name = "匿名用户"
	}

	client := &Client{
		hub:  hub,
		conn: conn,
		send: make(chan []byte, 256),
		name: name,
	}
	hub.register <- client

	// 广播加入消息
	hub.broadcastMsg(newMessage("join", name, "加入了聊天室"))

	// 启动读写 goroutine
	go client.writePump()
	go client.readPump()
}

func main() {
	hub := newHub()
	go hub.run()

	http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
		serveWs(hub, w, r)
	})
	http.Handle("/", http.FileServer(http.Dir("./public")))

	fmt.Println("聊天室服务监听 http://localhost:8080")
	log.Fatal(http.ListenAndServe(":8080", nil))
}

代码要点解读:

  • serveWs 在注册 client 后立即广播 join 消息,让所有人看到「某某加入」。
  • readPump 退出时(连接断开)广播 leave 消息并注销自己。
  • notifyUserList 在每次注册/注销后把最新用户列表广播给所有人,前端据此更新侧边栏。
  • 静态文件服务把 ./public 目录作为前端页面根目录。

七、前端页面:HTML + JavaScript 客户端

前端用一个纯 HTML + JavaScript 页面,不依赖任何框架。它负责建立连接、收发消息、渲染聊天记录和用户列表。

下面是 public/index.html 的内容,保存到与服务端同级的 public 目录即可运行:

html
<!DOCTYPE html>
<html lang="zh-CN">
<head>
  <meta charset="UTF-8">
  <title>Go 聊天室</title>
  <style>
    body { font-family: sans-serif; margin: 0; display: flex; height: 100vh; }
    #chat { flex: 1; display: flex; flex-direction: column; }
    #messages { flex: 1; overflow-y: auto; padding: 16px; background: #f5f5f5; }
    .msg { margin-bottom: 8px; padding: 6px 10px; background: #fff; border-radius: 4px; }
    .msg .meta { font-size: 12px; color: #888; margin-right: 8px; }
    .msg .from { font-weight: bold; color: #1976d2; }
    .sys { color: #888; font-style: italic; }
    #input { display: flex; padding: 8px; border-top: 1px solid #ddd; }
    #text { flex: 1; padding: 8px; font-size: 14px; }
    #send { padding: 8px 16px; margin-left: 8px; }
    #users { width: 200px; border-left: 1px solid #ddd; padding: 16px; }
    #users h3 { margin-top: 0; }
    #users li { list-style: none; padding: 4px 0; }
    #bar { padding: 8px; background: #1976d2; color: #fff; }
    #name { padding: 4px; }
  </style>
</head>
<body>
  <div id="chat">
    <div id="bar">
      Go 聊天室 昵称: <input id="name" placeholder="输入昵称">
      <button id="connect">连接</button>
      <button id="disconnect" disabled>断开</button>
    </div>
    <div id="messages"></div>
    <div id="input">
      <input id="text" placeholder="输入消息回车发送" disabled>
      <button id="send" disabled>发送</button>
    </div>
  </div>
  <div id="users">
    <h3>在线用户</h3>
    <ul id="userList"></ul>
  </div>

<script>
let ws = null;

function appendMsg(m) {
  const box = document.getElementById('messages');
  const div = document.createElement('div');
  div.className = 'msg';
  if (m.type === 'system') {
    div.className = 'msg sys';
    div.textContent = m.time + ' 当前在线: ' + (m.users || []).join(', ');
  } else if (m.type === 'join' || m.type === 'leave') {
    div.className = 'msg sys';
    div.textContent = m.time + ' ' + m.from + ' ' + m.text;
  } else {
    div.innerHTML = '<span class="meta">' + m.time + '</span>' +
      '<span class="from">' + m.from + ':</span> ' + m.text;
  }
  box.appendChild(div);
  box.scrollTop = box.scrollHeight;
}

function updateUserList(users) {
  const ul = document.getElementById('userList');
  ul.innerHTML = '';
  (users || []).forEach(function(u) {
    const li = document.createElement('li');
    li.textContent = u;
    ul.appendChild(li);
  });
}

function connect() {
  const name = document.getElementById('name').value || '匿名用户';
  ws = new WebSocket('ws://' + location.host + '/ws?name=' + encodeURIComponent(name));

  ws.onopen = function() {
    document.getElementById('text').disabled = false;
    document.getElementById('send').disabled = false;
    document.getElementById('connect').disabled = true;
    document.getElementById('disconnect').disabled = false;
  };

  ws.onmessage = function(e) {
    const m = JSON.parse(e.data);
    appendMsg(m);
    if (m.type === 'system') {
      updateUserList(m.users);
    }
  };

  ws.onclose = function() {
    document.getElementById('text').disabled = true;
    document.getElementById('send').disabled = true;
    document.getElementById('connect').disabled = false;
    document.getElementById('disconnect').disabled = true;
  };
}

function send() {
  const input = document.getElementById('text');
  const text = input.value.trim();
  if (!text || !ws) return;
  ws.send(JSON.stringify({ type: 'chat', text: text }));
  input.value = '';
}

document.getElementById('connect').onclick = connect;
document.getElementById('disconnect').onclick = function() { if (ws) ws.close(); };
document.getElementById('send').onclick = send;
document.getElementById('text').onkeydown = function(e) {
  if (e.key === 'Enter') send();
};
</script>
</body>
</html>

前端逻辑说明:

  • 用户输入昵称点「连接」,建立 WebSocket 连接,把昵称通过 query 参数传给服务端。
  • onmessage 收到消息后按 type 分发渲染:system 更新用户列表,join/leave 显示系统提示,chat 显示聊天气泡。
  • 输入框回车或点「发送」把消息 JSON 化后 ws.send,服务端会补全 fromtime 后广播给所有人(包括自己)。

运行方式

  1. 把服务端代码保存为 main.go
  2. 把上面的 HTML 保存到 public/index.html
  3. 执行 go run main.go
  4. 浏览器打开 http://localhost:8080,输入昵称连接。
  5. 多开几个浏览器窗口(不同昵称),发消息就能看到实时广播效果,用户列表也会同步更新。

八、小结

本篇我们从零实现了一个多人在线聊天室。核心是 Hub/Client 架构:Hub 作为中枢,在单个 goroutine 中用 select 处理注册、注销、广播三类事件,通过 channel 通信完全避免了锁;Client 用「读 goroutine + 写 goroutine + send channel」的经典模式,满足 gorilla/websocket 的并发约束。消息层用 JSON 文本帧 + type 字段区分业务,前端用原生 WebSocket API 即可对接。

关键要点回顾:

  • Hub 集中管理:所有连接注册到 Hub,广播由 Hub 完成,连接之间不直接通信。
  • channel 代替锁:Hub 的状态只在自身 goroutine 中修改,无需 mutex。
  • 非阻塞发送:广播时用 select + default,消费过慢的客户端直接踢掉,保护 Hub。
  • 读写分离:每个 Client 一个读 goroutine、一个写 goroutine,写操作串行化。
  • 心跳保活:服务端定期 Ping,读超时续期,避免长连接被中间设备断开。
  • 协议设计:JSON + type 字段是前后端协作的通用模式。

下一篇我们将把这套架构迁移到 Gin 框架,加入认证、限流、连接管理器、私聊、群组、Redis 跨节点广播等生产级特性。