Go TCP、WebSocket 与 Redis 实战
# Go TCP、WebSocket 与 Redis 实战
从 TCP 字节流理解网络服务的底层边界,再掌握 WebSocket 长连接和 Redis 客户端在真实 Go 项目中的位置,能够看懂实时通知、在线状态和跨实例广播代码。
# 先记住一句话
TCP 提供可靠有序的字节流,但不提供业务消息边界;WebSocket 在长连接上定义消息帧;Redis 是独立网络服务,常用于共享缓存、短期状态和跨实例事件分发。
浏览器或客户端
├─ 普通 HTTP:一次请求对应一次响应
├─ WebSocket:建立长连接后双方随时发送消息
└─ 自定义 TCP:应用自己定义消息格式和边界
多个 Go 服务实例
└─ 通过 Redis 共享缓存、会话、限流状态或发布订阅事件
2
3
4
5
6
7
# 先分清它们所在的层次
| 技术 | 主要解决的问题 | 是否自带业务数据持久化 |
|---|---|---|
| TCP | 两个端点之间可靠、有序地传输字节 | 否 |
| HTTP | 在连接上定义请求与响应语义 | 否 |
| WebSocket | 在长连接上定义双向消息帧 | 否 |
| Redis | 在独立进程中保存数据并提供数据结构命令 | 取决于配置和使用方式 |
WebSocket 通常建立在 TCP 和 HTTP 握手之上,Redis 客户端也通常通过 TCP 访问 Redis Server。它们不是互相替代的方案。
# TCP 没有消息边界
发送方连续写入两次,不保证接收方也正好读取两次。一次 Read 可能只读到半条消息,也可能同时读到多条消息的一部分。因此应用必须自己定义分帧协议,例如:
- 每行一条消息,以换行符分隔;
- 固定长度消息;
- 长度前缀加消息正文;
- 使用 HTTP、WebSocket、gRPC 等已经定义好边界的上层协议。
下面是一个使用 “每行一条消息” 协议的完整 TCP Echo Server:
// 文件位置:cmd/tcp-echo/main.go
package main // main 包表示这段代码会被编译成可执行程序,而不是供其他包导入的库。
import (
"bufio" // buffered I/O(带缓冲的输入输出)工具;这里用 Scanner 按行读取连接中的数据。
"fmt" // 格式化输入输出;这里用 Fprintf 按指定格式向 TCP 连接写回响应。
"log" // 输出运行日志;这里记录监听、读写等操作产生的错误。
"net" // Go 标准库的网络包;提供 TCP 监听器和网络连接等 API。
"time" // 时间处理包;这里用于设置每次写响应的超时时间。
)
func main() {
// 在本机 127.0.0.1 的 9000 端口监听 TCP 连接:
// 第一个参数 "tcp" 表示使用 TCP,第二个参数表示监听地址。
// listener 用来接收客户端连接;err 表示监听端口时是否发生错误。
listener, err := net.Listen("tcp", "127.0.0.1:9000")
if err != nil {
// log.Fatal 是 log 包提供的函数:先把 err 输出到日志,再调用 os.Exit(1) 立即结束程序。
// 它适合处理端口被占用等 “服务无法继续启动” 的严重错误;程序退出时不会执行 defer。
log.Fatal(err)
}
defer listener.Close() // main 退出时释放监听端口。
// TCP Server 通常需要一直运行,因此使用无限循环不断接收新连接。
for {
// Accept 会阻塞等待客户端连接;连接建立后返回代表该连接的 net.Conn。
connection, err := listener.Accept()
if err != nil {
log.Printf("accept connection: %v", err)
continue // 本次接收失败时跳过后续代码,继续等待下一条连接。
}
// go 会启动一个新的 Goroutine 处理当前连接,主 Goroutine 不必等待它结束,
// 可以立刻回到 Accept 继续接收其他客户端,因此多个连接能够并发处理。
go handleConnection(connection)
}
}
// handleConnection 负责读取并响应一条客户端连接。
// connection 的类型是 net.Conn,表示已经建立的网络连接,可以从中读取数据,也可以向其中写入数据。
func handleConnection(connection net.Conn) {
defer connection.Close() // 连接处理结束后释放文件描述符。
// bufio.NewScanner 把网络连接包装成 Scanner。
// Scanner 默认以换行符分段,所以客户端每发送一行,这里就能读到一条完整消息;返回的文本不包含换行符。
scanner := bufio.NewScanner(connection)
// 调整 Scanner 读取单条消息时使用的缓冲区:
// make([]byte, 4096) 创建 4 KiB 的初始缓冲区;1024*1024 把每行上限设为 1 MiB。
// 消息超过上限时 Scan 会停止并产生错误,避免客户端发送超长内容无限占用服务器内存。
scanner.Buffer(make([]byte, 4096), 1024*1024)
// Scan 会阻塞等待下一行数据;成功读到一行时返回 true,连接关闭或读取失败时返回 false。
for scanner.Scan() {
// 每次写响应前设置 5 秒截止时间,避免客户端一直不接收数据,导致当前 Goroutine 永久卡在写操作上。
// time.Now() 是当前时间,Add(5 * time.Second) 得到当前时间 5 秒后的绝对时间点。
if err := connection.SetWriteDeadline(time.Now().Add(5 * time.Second)); err != nil {
log.Printf("set write deadline: %v", err)
return
}
// scanner.Text() 取得刚读到的一行文本。
// fmt.Fprintf 把格式化后的内容直接写入 connection:例如收到 hello,就写回 "echo: hello\n"。
// Fprintf 返回写入的字节数和错误;这里不需要字节数,因此用 _ 忽略,只检查 err。
if _, err := fmt.Fprintf(connection, "echo: %s\n", scanner.Text()); err != nil {
log.Printf("write response: %v", err)
return
}
}
// Scan 返回 false 可能只是客户端正常关闭了连接,也可能是读取失败。
// 正常结束时 Err 返回 nil;发生网络错误或消息超过 1 MiB 等情况时,Err 返回具体错误。
if err := scanner.Err(); err != nil {
log.Printf("read connection: %v", err)
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
# 启动服务后,在另一个终端建立 TCP 连接进行手工验证。
go run ./cmd/tcp-echo
nc 127.0.0.1 9000
2
3
这个例子为了突出消息边界,只设置写超时。生产服务还要设置读空闲超时、连接数上限、协议级认证、最大消息长度、优雅关闭和日志指标。
所谓 “粘包” 是什么
TCP 只交付连续字节流,不承诺保留每次 Write 的边界。所谓 “粘包” 并不是 TCP 把包处理错了,而是应用误以为一次 Read 就等于一条业务消息。正确做法是设计长度前缀、分隔符或使用已有上层协议。
# WebSocket 适合双向实时交互
WebSocket 建立后,客户端和服务端都能主动发送消息。聊天、协同编辑、实时控制和高频双向事件更适合 WebSocket;如果只需要服务端向浏览器持续推送进度或模型回答,SSE 往往更简单。
Go 标准库没有提供完整的 WebSocket Server API。下面使用 github.com/coder/websocket:
# 在项目根目录安装第三方 WebSocket 库。
go get github.com/coder/websocket
2
// 文件位置:internal/realtime/echo.go
package realtime // realtime 是当前业务包名,这个文件会被路由层导入使用,不是独立可执行程序。
import (
"context" // 控制一次读写操作的超时和取消,避免 Goroutine 永久阻塞。
"log" // 输出运行日志;这里用于记录写消息时出现的错误。
"net/http" // 提供 HTTP Handler 使用的 ResponseWriter 和 Request 类型。
"time" // 提供 Second 等时间单位,用来设置读写超时。
"github.com/coder/websocket" // 建立、管理和关闭 WebSocket 连接。
"github.com/coder/websocket/wsjson" // 在 WebSocket 连接中读写 JSON 消息。
)
// Message 定义客户端和服务端约定的 JSON 消息结构。
// 例如 Message{Type: "echo", Text: "hello"} 会被编码为 {"type":"echo","text":"hello"}。
type Message struct {
// 反引号中的 json:"..." 是 Struct Tag,用于指定字段在 JSON 中的名称。
Type string `json:"type"` // 消息类型,例如 echo,用于区分不同种类的消息。
Text string `json:"text"` // 消息正文,例如用户发送的 hello。
}
// Echo 是一个 HTTP Handler,用来把普通 HTTP 请求升级为 WebSocket 连接并回显消息。
// writer 用于向客户端写 HTTP 响应;request 保存客户端发来的请求头、地址等信息。
func Echo(writer http.ResponseWriter, request *http.Request) {
// Accept 检查 WebSocket 握手请求,并把当前 HTTP 连接升级为 WebSocket 连接。
// 第三个参数用于传入 AcceptOptions;这里传 nil,表示使用库的默认配置和 Origin 校验。
// connection 代表升级成功后的长连接;err 表示握手是否失败。
connection, err := websocket.Accept(writer, request, nil)
if err != nil {
// Accept 失败时已经向客户端写入相应的 HTTP 错误响应,因此直接结束本次请求。
return
}
// CloseNow 会立即关闭底层网络连接,不执行正常的 WebSocket 关闭握手。
// 放在 defer 中是为了兜底:Echo 因读写失败 return 时,连接仍然一定会被释放。
defer connection.CloseNow()
// WebSocket 是长连接,所以使用无限循环反复执行 “读取一条消息 → 写回一条消息”。
// 客户端断开、操作超时或读写失败时,循环中的 return 才会结束当前连接。
for {
// 创建最长 60 秒的读取 Context。
// 如果 60 秒内没有读到完整消息,readContext 会自动取消,使 Read 不会永久阻塞。
// WithTimeout 返回 Context 和取消函数;60*time.Second 表示 60 秒。
readContext, cancelRead := context.WithTimeout(context.Background(), 60*time.Second)
// incoming 用来接收客户端发来的 JSON 数据,初始值是 Message 的零值。
var incoming Message
// Read 等待下一条 WebSocket 消息,把 JSON 解码后写入 incoming。
// 必须传 &incoming(变量地址),wsjson.Read 才能修改这个变量的字段。
// 这里使用 = 而不是 :=,因为 err 已经由上面的 websocket.Accept 声明过。
err = wsjson.Read(readContext, connection, &incoming)
// 无论读取成功还是失败,本轮读取已经结束,应立即释放超时计时器占用的资源。
cancelRead()
if err != nil {
// 客户端正常关闭、网络断开、JSON 格式错误和读取超时都会结束当前连接。
// return 之后会执行上面的 defer connection.CloseNow()。
return
}
// 为本次写响应设置 5 秒超时,避免客户端不接收数据时一直阻塞当前 Goroutine。
writeContext, cancelWrite := context.WithTimeout(context.Background(), 5*time.Second)
// Write 会先把 Message 编码成 JSON,再作为一条 WebSocket 消息发送给客户端。
// Type 固定为 echo,Text 使用刚刚收到的内容,所以它实现的是 “收到什么就回显什么”。
err = wsjson.Write(writeContext, connection, Message{
Type: "echo",
Text: incoming.Text,
})
// 本轮写操作完成后立即释放写超时使用的计时器资源。
cancelWrite()
if err != nil {
log.Printf("write websocket message: %v", err)
return
}
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
在路由中把它注册为 http.HandleFunc("/ws", realtime.Echo)。默认的 Origin 校验适合同源访问;确实需要跨域时,应通过 websocket.AcceptOptions 显式配置可信 Origin,而不是关闭校验接受任意网站。
生产 WebSocket 还要处理:
- 身份认证和单个用户的连接上限;
- 读写超时、Ping/Pong 和失活连接清理;
- 每条连接的发送队列以及慢客户端背压;
- 单连接并发读写约束;
- 服务关闭时的连接摘流和正常关闭帧;
- 多实例部署时,消息怎样路由到持有目标连接的实例。
# Redis 在 Go 项目中的位置
Redis Server 部署在本机、容器或云托管环境中;Go 代码通过客户端连接它。go-redis 会管理连接池,因此通常为一个进程创建并复用一个 Client,不要每次请求都新建连接。
# 本地已有 Docker 时,可以启动一个仅供学习使用的 Redis 容器。
docker run --name go-notes-redis -p 6379:6379 redis:7-alpine
# 在 Go Module 中安装 Redis 客户端。
go get github.com/redis/go-redis/v9
2
3
4
5
不要把没有认证的 6379 端口暴露到公网。生产环境应使用私有网络、ACL、TLS、Secret 管理、超时和监控。
// 文件位置:internal/cache/note_cache.go
package cache // cache 包集中封装缓存操作,让 Service 不必直接依赖 Redis API。
import (
"context" // 向 Redis 操作传递超时、取消信号和请求生命周期。
"errors" // 判断一个错误是否属于指定错误,例如 redis.Nil。
"fmt" // 使用 Errorf 为底层错误补充当前操作的上下文。
"time" // 提供 Duration 类型,用来表示连接超时和缓存有效期。
"github.com/redis/go-redis/v9" // Redis 的 Go 客户端,提供连接池和各种 Redis 命令。
)
// NoteCache 是 Note 数据的 Redis 缓存组件。
// 它把 Redis Client 和统一的缓存有效期保存在一起,供 Get、Set 和 Close 方法使用。
type NoteCache struct {
client *redis.Client // 指向 Redis Client;一个 Client 内部维护可并发复用的连接池。
ttl time.Duration // 每条缓存数据的有效期,例如 10*time.Minute 表示 10 分钟。
}
// NewNoteCache 创建并检查一个可以使用的 NoteCache。
// redisURL 是 Redis 连接地址,例如 redis://:password@localhost:6379/0;ttl 是缓存有效期。
// 成功时返回 *NoteCache 和 nil;失败时返回 nil 和具体错误。
func NewNoteCache(redisURL string, ttl time.Duration) (*NoteCache, error) {
// ParseURL 把 URL 中的地址、数据库编号、用户名、密码和 TLS 等配置解析成 Options。
options, err := redis.ParseURL(redisURL)
if err != nil {
// %w 在补充错误背景的同时保留原始错误,调用方仍可用 errors.Is 或 errors.As 判断它。
return nil, fmt.Errorf("parse redis URL: %w", err)
}
// NewClient 根据 Options 创建 Client 和连接池管理器,但此时不代表 Redis 一定能够连接成功。
// Client 可以被多个 Goroutine 并发复用,通常一个应用进程创建一个即可。
client := redis.NewClient(options)
// 创建一个最长 3 秒的 Context,专门限制下面 Ping 检查的等待时间。
// Background 是根 Context;WithTimeout 返回带超时的 Context 和用于提前释放资源的 cancel。
checkContext, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel() // NewNoteCache 返回前释放计时器资源;重复调用 cancel 也是安全的。
// Ping 向 Redis 发送 PING 命令;成功收到响应才能确认地址、网络和认证配置可用。
// Ping 返回命令结果对象,调用 Err() 可以只取出这次命令产生的错误。
if err := client.Ping(checkContext).Err(); err != nil {
// 连通性检查失败时不再返回这个 Client,先关闭它持有的连接池资源。
// 这里用 _ 忽略 Close 的错误,避免它覆盖更关键的 Ping 错误。
_ = client.Close()
return nil, fmt.Errorf("connect to redis: %w", err)
}
// &NoteCache{...} 创建 NoteCache 并返回它的指针;nil 表示构造过程没有错误。
return &NoteCache{client: client, ttl: ttl}, nil
}
// Get 根据 Note ID 查询缓存。
// ctx 由上层请求传入,用于控制本次 Redis 查询的超时或取消。
// 三个返回值依次是:缓存内容、是否命中缓存、执行过程中是否发生系统错误。
func (c *NoteCache) Get(ctx context.Context, id string) (value string, found bool, err error) {
// "note:"+id 生成 Redis Key,例如 id 为 42 时得到 note:42。
// note: 前缀用于区分不同业务的数据,避免与其他类型的 Key 重名。
// Result() 同时取出 GET 命令返回的字符串和错误;这里把它们赋给命名返回值。
value, err = c.client.Get(ctx, "note:"+id).Result()
// Redis 找不到 Key 时,go-redis 返回特殊错误 redis.Nil。
// 缓存未命中是正常业务结果,不是系统故障,所以返回 found=false、err=nil。
if errors.Is(err, redis.Nil) {
return "", false, nil
}
// 连接中断、超时等其他错误才是真正的缓存故障,需要交给上层处理或决定是否降级。
if err != nil {
return "", false, fmt.Errorf("get note cache: %w", err)
}
// 没有错误说明成功取得缓存内容,因此返回 value、found=true 和 err=nil。
return value, true, nil
}
// Set 把一条 Note 数据写入 Redis 缓存。
// id 用于组成 Redis Key,value 是要缓存的内容,过期时间统一使用 c.ttl。
func (c *NoteCache) Set(ctx context.Context, id string, value string) error {
// Set 的四个参数依次是:Context、Key、Value 和过期时间。
// c.ttl 大于 0 时,Key 到期后由 Redis 自动删除;等于 0 时表示不过期。
// Err() 只取出 SET 命令执行过程中产生的错误。
if err := c.client.Set(ctx, "note:"+id, value, c.ttl).Err(); err != nil {
return fmt.Errorf("set note cache: %w", err)
}
return nil // nil 表示缓存写入成功。
}
// Close 关闭 Redis Client 管理的连接池。
// 它应在应用优雅关闭阶段调用一次,而不是每次请求结束后调用。
func (c *NoteCache) Close() error {
return c.client.Close()
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
上面实现的是 Cache-Aside 所需的缓存读写组件:Service 先读缓存,未命中再读数据库并回填。缓存不可用时是否降级访问数据库,要根据数据库容量和业务 SLO 决定,不能默认无限降级把数据库打垮。
# WebSocket 多实例为什么会用 Redis
单实例可以在内存中保存 userID -> connections。扩容后,用户 A 可能连接实例 1,业务事件却到达实例 2;实例 2 无法直接访问实例 1 的内存。这时可以通过 Redis Pub/Sub 把实时事件广播给所有实例,再由持有目标连接的实例推送。
业务事件
└─ 发布到 Redis Channel
├─ Go 实例 1 订阅 → 找到本机连接 → WebSocket 推送
├─ Go 实例 2 订阅 → 本机没有目标连接 → 忽略
└─ Go 实例 3 订阅 → 找到本机连接 → WebSocket 推送
2
3
4
5
Pub/Sub 不等于可靠消息队列
Redis Pub/Sub 不持久化历史消息,订阅者断线期间的消息通常不会补发。在线通知可以接受这种语义;订单、支付和必须重放的任务应使用 Redis Streams、专门消息队列或持久化事件表。
# 一个常见项目目录
go-realtime-service/
├─ cmd/
│ └─ api/
│ └─ main.go # 装配 HTTP、WebSocket、Redis 并处理关闭
├─ internal/
│ ├─ api/
│ │ └─ routes.go # 注册普通 HTTP 和 /ws 路由
│ ├─ realtime/
│ │ ├─ connection.go # 单条 WebSocket 连接的读写生命周期
│ │ └─ hub.go # 管理本实例的连接与发送队列
│ ├─ cache/
│ │ └─ note_cache.go # Redis 缓存读写
│ └─ events/
│ └─ redis_pubsub.go # 跨实例发布与订阅实时事件
├─ go.mod # 模块路径和第三方依赖
└─ go.sum # 依赖内容校验
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
# 高频面试题与回答
1. 为什么 TCP 一次 Read 不一定得到一条完整消息?参考答案
TCP 提供的是可靠有序的字节流,不保留应用每次 Write 的边界。接收端必须根据长度前缀、分隔符或上层协议恢复业务消息,不能把一次 Read 当作一条消息。
2. SSE 和 WebSocket 怎样选择?参考答案
SSE 主要是服务器到浏览器的单向事件流,基于普通 HTTP,适合模型流式回答、进度和日志;WebSocket 支持双方在同一长连接上主动发送消息,适合聊天、协同和实时控制,但连接状态、心跳、背压和扩容更复杂。
3. Go Redis Client 为什么通常做成单例复用?参考答案
go-redis Client 本身可以被多个 Goroutine 安全复用,并管理连接池。每次请求新建 Client 会反复建连、增加 Redis 和系统资源压力,也难以统一配置超时、指标和关闭流程。
4. Redis Pub/Sub 能不能当可靠队列?参考答案
通常不能。Pub/Sub 面向在线订阅者,不保存供离线消费者补读的消息,也没有业务级确认与重试。允许丢失的实时通知可以使用它;必须可靠处理的任务要选择 Redis Streams、消息队列或数据库事件表。
# 接下来学什么
下一篇学习 Go 测试、工程化与微服务边界,把网络组件放进可测试、可配置和可部署的项目结构中。