外观
Redis 缓存、Session 与 Pub/Sub
Phase 05 — Redis 涵盖:
Cache·Session·Pub/Sub
1. 学习目标
完成本知识点后,你应该能够:
- 理解 Cache-Aside 等常见缓存模式
- 使用 Redis 实现 Session 存储
- 使用 Redis Pub/Sub 实现消息发布与订阅
- 在 Go 中结合 go-redis 实现缓存、Session 和实时广播
- 为数字孪生系统设计缓存策略和实时数据管道
2. 为什么需要
Cache(缓存)
设备详情、仓库配置、报表数据——这些读多写少的数据,每次请求都查 PostgreSQL 浪费资源。Redis 缓存可以将响应时间从毫秒级降到亚毫秒级,并减轻数据库压力。
Session(会话)
数字孪生管理后台需要用户登录。JWT 无状态但无法主动失效;Redis Session 存储登录态,支持登出、踢人、过期——适合需要服务端控制的场景。
Pub/Sub(发布订阅)
AGV 状态变更、告警触发、设备上下线——多个服务或 WebSocket 网关需要实时感知这些事件。Pub/Sub 提供轻量级消息广播,是构建实时数据管道的基石。
3. 核心概念
3.1 Cache 模式
| 模式 | 流程 | 适用 |
|---|---|---|
| Cache-Aside | 读:先查缓存,miss 则查 DB 并回填;写:先写 DB,再删缓存 | 最常用 |
| Read-Through | 缓存层代理读,miss 时缓存层自己查 DB | 缓存中间件 |
| Write-Through | 写同时更新缓存和 DB | 强一致 |
| Write-Behind | 先写缓存,异步写 DB | 高性能写入 |
Cache-Aside 流程:
读:App → Redis GET → (miss) → PostgreSQL → Redis SET → 返回
写:App → PostgreSQL → Redis DEL3.2 缓存问题
| 问题 | 说明 | 应对 |
|---|---|---|
| 缓存穿透 | 查不存在的数据,每次都打到 DB | 缓存空值 / 布隆过滤器 |
| 缓存击穿 | 热点 key 过期瞬间大量请求打到 DB | 互斥锁 / 永不过期+异步更新 |
| 缓存雪崩 | 大量 key 同时过期 | TTL 加随机偏移 |
3.3 Session
登录成功 → 生成 sessionID → Redis SET session:{id} → 返回 cookie
后续请求 → 读取 cookie → Redis GET session:{id} → 验证身份
登出 → Redis DEL session:{id}Session 数据结构(Hash 或 JSON String):
json
{
"user_id": "admin-001",
"username": "admin",
"role": "operator",
"login_at": "2026-09-02T10:00:00Z"
}3.4 Pub/Sub
Publisher → channel: "device:events" → Subscribers (WebSocket Gateway, Logger, ...)| 概念 | 说明 |
|---|---|
| Publisher | 发布消息到 channel |
| Subscriber | 订阅 channel,接收消息 |
| Channel | 消息频道,如 device:AGV-001:telemetry |
| Pattern | PSUBSCRIBE device:* 模式订阅 |
特点:
- fire-and-forget:订阅者离线期间的消息会丢失
- 无持久化:不适合可靠消息队列(那是 Kafka 的场景)
- 低延迟:适合实时通知
4. 基础语法
4.1 Cache-Aside 实现
go
package cache
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type DeviceCache struct {
rdb *redis.Client
ttl time.Duration
}
func NewDeviceCache(rdb *redis.Client, ttl time.Duration) *DeviceCache {
return &DeviceCache{rdb: rdb, ttl: ttl}
}
func (c *DeviceCache) cacheKey(id string) string {
return fmt.Sprintf("cache:device:%s", id)
}
func (c *DeviceCache) Get(ctx context.Context, id string) (*Device, error) {
data, err := c.rdb.Get(ctx, c.cacheKey(id)).Bytes()
if err == redis.Nil {
return nil, nil // cache miss
}
if err != nil {
return nil, err
}
var d Device
if err := json.Unmarshal(data, &d); err != nil {
return nil, err
}
return &d, nil
}
func (c *DeviceCache) Set(ctx context.Context, d *Device) error {
data, err := json.Marshal(d)
if err != nil {
return err
}
// TTL 加随机偏移,避免雪崩
jitter := time.Duration(rand.Intn(60)) * time.Second
return c.rdb.Set(ctx, c.cacheKey(d.ID), data, c.ttl+jitter).Err()
}
func (c *DeviceCache) Delete(ctx context.Context, id string) error {
return c.rdb.Del(ctx, c.cacheKey(id)).Err()
}
type Device struct {
ID string `json:"id"`
Name string `json:"name"`
Status string `json:"status"`
}go
// Service 层 Cache-Aside
func (s *DeviceService) GetDevice(ctx context.Context, id string) (*Device, error) {
// 1. 查缓存
if cached, err := s.cache.Get(ctx, id); err != nil {
return nil, err
} else if cached != nil {
return cached, nil
}
// 2. 查数据库
device, err := s.repo.GetByID(ctx, id)
if err != nil {
return nil, err
}
// 3. 回填缓存
_ = s.cache.Set(ctx, device)
return device, nil
}
func (s *DeviceService) UpdateDevice(ctx context.Context, d *Device) error {
if err := s.repo.Update(ctx, d); err != nil {
return err
}
// 写 DB 后删缓存
return s.cache.Delete(ctx, d.ID)
}4.2 Session 实现
go
package session
import (
"context"
"crypto/rand"
"encoding/hex"
"encoding/json"
"time"
"github.com/redis/go-redis/v9"
)
type Store struct {
rdb *redis.Client
ttl time.Duration
}
type SessionData struct {
UserID string `json:"user_id"`
Username string `json:"username"`
Role string `json:"role"`
LoginAt time.Time `json:"login_at"`
}
func NewStore(rdb *redis.Client, ttl time.Duration) *Store {
return &Store{rdb: rdb, ttl: ttl}
}
func (s *Store) Create(ctx context.Context, data *SessionData) (string, error) {
id, err := generateSessionID()
if err != nil {
return "", err
}
key := "session:" + id
payload, err := json.Marshal(data)
if err != nil {
return "", err
}
if err := s.rdb.Set(ctx, key, payload, s.ttl).Err(); err != nil {
return "", err
}
return id, nil
}
func (s *Store) Get(ctx context.Context, id string) (*SessionData, error) {
data, err := s.rdb.Get(ctx, "session:"+id).Bytes()
if err == redis.Nil {
return nil, nil
}
if err != nil {
return nil, err
}
var session SessionData
if err := json.Unmarshal(data, &session); err != nil {
return nil, err
}
return &session, nil
}
func (s *Store) Delete(ctx context.Context, id string) error {
return s.rdb.Del(ctx, "session:"+id).Err()
}
func generateSessionID() (string, error) {
b := make([]byte, 32)
if _, err := rand.Read(b); err != nil {
return "", err
}
return hex.EncodeToString(b), nil
}go
// HTTP 中间件
func SessionMiddleware(store *session.Store) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
cookie, err := r.Cookie("session_id")
if err != nil {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
sess, err := store.Get(r.Context(), cookie.Value)
if err != nil || sess == nil {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
ctx := context.WithValue(r.Context(), "session", sess)
next.ServeHTTP(w, r.WithContext(ctx))
})
}
}4.3 Pub/Sub 实现
go
package pubsub
import (
"context"
"encoding/json"
"github.com/redis/go-redis/v9"
)
const DeviceEventsChannel = "device:events"
type Event struct {
Type string `json:"type"` // telemetry, alert, status_change
DeviceID string `json:"device_id"`
X float64 `json:"x,omitempty"`
Y float64 `json:"y,omitempty"`
Message string `json:"message,omitempty"`
}
type Publisher struct {
rdb *redis.Client
}
func NewPublisher(rdb *redis.Client) *Publisher {
return &Publisher{rdb: rdb}
}
func (p *Publisher) PublishEvent(ctx context.Context, event Event) error {
data, err := json.Marshal(event)
if err != nil {
return err
}
return p.rdb.Publish(ctx, DeviceEventsChannel, data).Err()
}
type Subscriber struct {
rdb *redis.Client
}
func NewSubscriber(rdb *redis.Client) *Subscriber {
return &Subscriber{rdb: rdb}
}
func (s *Subscriber) Subscribe(ctx context.Context, handler func(Event)) error {
pubsub := s.rdb.Subscribe(ctx, DeviceEventsChannel)
defer pubsub.Close()
ch := pubsub.Channel()
for {
select {
case <-ctx.Done():
return ctx.Err()
case msg := <-ch:
var event Event
if err := json.Unmarshal([]byte(msg.Payload), &event); err != nil {
continue
}
handler(event)
}
}
}go
// 发布 AGV 遥测事件
publisher.PublishEvent(ctx, pubsub.Event{
Type: "telemetry",
DeviceID: "AGV-001",
X: 12.5,
Y: 3.8,
})
// WebSocket 网关订阅并广播
go subscriber.Subscribe(ctx, func(event pubsub.Event) {
hub.Broadcast(event.DeviceID, event)
})5. 代码解析
go
if cached, err := s.cache.Get(ctx, id); cached != nil { return cached, nil }Cache-Aside 读路径:缓存命中直接返回,miss 则查 DB 并回填。cached == nil && err == nil 表示 miss。
go
return s.cache.Delete(ctx, d.ID)写路径:先写 DB,再删缓存(而非更新缓存)。避免并发写导致缓存与 DB 不一致。
go
jitter := time.Duration(rand.Intn(60)) * time.Second
c.rdb.Set(ctx, key, data, c.ttl+jitter)TTL 加随机偏移(如 5min ± 60s),避免大量 key 同时过期引发缓存雪崩。
go
pubsub := s.rdb.Subscribe(ctx, DeviceEventsChannel)
ch := pubsub.Channel()
for msg := range ch { ... }go-redis 的 Subscribe 是阻塞的,应在独立 goroutine 中运行,并通过 context 控制退出。
6. JavaScript / TypeScript 对比
| 概念 | Redis (Go) | 前端 |
|---|---|---|
| 缓存 | 服务端 Redis | 浏览器 HTTP Cache / Service Worker |
| Session | Redis + Cookie | Cookie / localStorage / JWT |
| Pub/Sub | Redis Pub/Sub | EventEmitter / BroadcastChannel |
| 实时推送 | Redis → Go → WebSocket → 前端 | 直接 WebSocket 消费 |
关键差异:
- 前端缓存(HTTP Cache)对用户不可控粒度;Redis 缓存是应用层精确控制
- JWT Session 无状态存在客户端;Redis Session 存在服务端,可主动失效
- 前端 EventEmitter 只在单页面内;Redis Pub/Sub 跨进程/跨服务
7. 常见错误
错误 1:先删缓存再写 DB
go
// ❌ 删缓存后、写 DB 前,另一个请求可能读到旧 DB 数据并回填旧缓存
s.cache.Delete(ctx, id)
s.repo.Update(ctx, d)
// ✅ 先写 DB,再删缓存
s.repo.Update(ctx, d)
s.cache.Delete(ctx, id)错误 2:缓存穿透未处理
go
// ❌ 恶意查询不存在的 device,每次都打到 DB
device, _ := s.repo.GetByID(ctx, "NONEXIST")
// ✅ 缓存空值,短 TTL
if device == nil {
s.cache.SetNull(ctx, id, 1*time.Minute)
}错误 3:Session ID 可预测
go
// ❌
sessionID := fmt.Sprintf("%d", time.Now().UnixNano())
// ✅ 加密安全随机
b := make([]byte, 32)
rand.Read(b)
sessionID := hex.EncodeToString(b)错误 4:Pub/Sub 当消息队列
go
// ❌ 期望离线消费者收到历史消息
// Pub/Sub 不持久化,订阅者离线则消息丢失
// ✅ 需要可靠投递时用 Redis Stream 或 Kafka错误 5:Subscribe 阻塞 HTTP Handler
go
// ❌ 在 Handler 里直接 Subscribe
func handler(w http.ResponseWriter, r *http.Request) {
subscriber.Subscribe(r.Context(), handler) // 阻塞
}
// ✅ 在 main 或独立 goroutine 启动订阅
go subscriber.Subscribe(ctx, eventHandler)8. 实际应用
数字孪生 — 完整数据流
┌─────────────┐
AGV 上报 ────────→│ Go API │
│ Handler │
└──────┬──────┘
│
┌────────────┼────────────┐
↓ ↓ ↓
Redis Hash PostgreSQL Redis Pub/Sub
(实时状态) (持久化) (device:events)
↓ ↓
GET /state API WebSocket Gateway
↓ ↓
前端轮询 Three.js 实时更新缓存策略
| 数据 | 缓存? | TTL | 原因 |
|---|---|---|---|
| 设备详情 | 是 | 5 min | 读多写少 |
| AGV 实时位置 | 是(Hash) | 60s 心跳 | 高频读写 |
| 历史轨迹 | 否 | — | 数据量大,查 DB |
| 用户 Session | 是 | 24h | 登录态 |
Pub/Sub 频道设计
| 频道 | 用途 |
|---|---|
device:events | 全局设备事件 |
device:AGV-001:telemetry | 单设备遥测(按需订阅) |
alerts:critical | 严重告警广播 |
9. 深入理解
9.1 Cache 与 Redis Hash 的分工
- Cache(String JSON):完整对象缓存,如设备详情,适合 Cache-Aside
- Hash:实时状态字段,部分更新,不需要「查 DB 回填」模式
9.2 Session vs JWT
| 维度 | Redis Session | JWT |
|---|---|---|
| 存储 | 服务端 | 客户端 |
| 失效 | 随时 DEL | 需等过期或黑名单 |
| 扩展 | 需共享 Redis | 无状态,易扩展 |
| 适用 | 管理后台、需踢人 | API、微服务 |
数字孪生项目可 JWT 做 API 认证,Redis Session 做 Web 管理后台。
9.3 Pub/Sub vs Redis Stream
| 维度 | Pub/Sub | Stream |
|---|---|---|
| 持久化 | 否 | 是 |
| 消费组 | 否 | 是 |
| 延迟 | 极低 | 低 |
| 适用 | 实时广播 | 可靠消息 |
数字孪生实时推送用 Pub/Sub 足够;告警持久化应写 DB 而非依赖 Pub/Sub。
9.4 优雅关闭 Subscribe
go
ctx, cancel := context.WithCancel(context.Background())
go subscriber.Subscribe(ctx, handler)
// 收到 SIGTERM 时
cancel() // Subscribe goroutine 退出10. 练习
请独立完成,不要查看答案。代码放在
workspace/phase-05/cache-pubsub/。
Level 1 — 基础
练习 1.1:实现 DeviceCache 的 Get、Set、Delete 三个方法。
练习 1.2:实现 Cache-Aside 读路径:缓存 miss 时返回 nil,由调用方查 DB。
练习 1.3:创建 Session,写入 Redis,读取并验证数据。
练习 1.4:发布一条 Pub/Sub 消息,订阅端打印收到的内容。
Level 2 — 应用
练习 2.1:实现完整 Cache-Aside:GetDevice(缓存 + DB)、UpdateDevice(DB + 删缓存)。
练习 2.2:实现 Session 中间件:无 cookie 或 session 无效返回 401。
练习 2.3:TTL 加随机偏移,避免缓存雪崩。
练习 2.4:缓存空值(设备不存在时缓存 "null" 标记,TTL 1 分钟)。
Level 3 — 综合
练习 3.1:AGV 上报 Handler:写 Redis Hash + 写 PostgreSQL + Publish 事件。
练习 3.2:独立 goroutine 订阅 device:events,收到遥测事件后打印 JSON。
练习 3.3:设计并实现登录 / 登出 API(Session 创建与删除)。
Level 4 — 项目实践
练习 4.1:搭建 mini 数字孪生数据管道:HTTP 接收遥测 → Redis 状态 + PG 持久化 + Pub/Sub 广播;独立订阅进程消费事件;设备详情 API 使用 Cache-Aside。用 curl 验证全流程。
11. 学习检查
完成练习后,确认你能回答:
- Cache-Aside 的读写路径分别是什么?
- 为什么写 DB 后「删缓存」而非「更新缓存」?
- 缓存穿透、击穿、雪崩分别是什么?如何应对?
- Redis Session 和 JWT 各适合什么场景?
- Pub/Sub 的消息会持久化吗?离线订阅者能收到吗?
- Pub/Sub 和 Redis Stream 如何选择?
12. 下一步
| 已完成 | 下一知识点 | 关系 |
|---|---|---|
| Cache, Session, Pub/Sub | WebSocket | Pub/Sub 事件推送到前端 |
| Connection | WebSocket 连接管理 | |
| Broadcast | 多客户端实时广播 |
Phase 05 完成后,Project 02 实时系统的 Redis 基础设施就绪。Phase 06 将用 WebSocket 把 Pub/Sub 事件推送给 Three.js 前端,完成数字孪生实时展示闭环。
前置知识提醒
- 已完成 Redis 数据类型文档,熟悉 go-redis 基本操作
- 已完成 Phase 04 PostgreSQL,理解 Cache-Aside 中 DB 的角色
- 已完成 Phase 03 JWT,可对比 Session 方案