Appearance
聊天室系统实战
前两篇我们打下了协议和 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:持续发送消息
writePump 用 select 同时监听 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,服务端会补全from和time后广播给所有人(包括自己)。
运行方式
- 把服务端代码保存为
main.go。 - 把上面的 HTML 保存到
public/index.html。 - 执行
go run main.go。 - 浏览器打开
http://localhost:8080,输入昵称连接。 - 多开几个浏览器窗口(不同昵称),发消息就能看到实时广播效果,用户列表也会同步更新。
八、小结
本篇我们从零实现了一个多人在线聊天室。核心是 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 跨节点广播等生产级特性。