Skip to content

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 DEL

3.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
PatternPSUBSCRIBE 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
SessionRedis + CookieCookie / localStorage / JWT
Pub/SubRedis Pub/SubEventEmitter / BroadcastChannel
实时推送Redis → Go → WebSocket → 前端直接 WebSocket 消费

关键差异

  1. 前端缓存(HTTP Cache)对用户不可控粒度;Redis 缓存是应用层精确控制
  2. JWT Session 无状态存在客户端;Redis Session 存在服务端,可主动失效
  3. 前端 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
用户 Session24h登录态

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 SessionJWT
存储服务端客户端
失效随时 DEL需等过期或黑名单
扩展需共享 Redis无状态,易扩展
适用管理后台、需踢人API、微服务

数字孪生项目可 JWT 做 API 认证,Redis Session 做 Web 管理后台。

9.3 Pub/Sub vs Redis Stream

维度Pub/SubStream
持久化
消费组
延迟极低
适用实时广播可靠消息

数字孪生实时推送用 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:实现 DeviceCacheGetSetDelete 三个方法。

练习 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. 学习检查

完成练习后,确认你能回答:

  1. Cache-Aside 的读写路径分别是什么?
  2. 为什么写 DB 后「删缓存」而非「更新缓存」?
  3. 缓存穿透、击穿、雪崩分别是什么?如何应对?
  4. Redis Session 和 JWT 各适合什么场景?
  5. Pub/Sub 的消息会持久化吗?离线订阅者能收到吗?
  6. Pub/Sub 和 Redis Stream 如何选择?

12. 下一步

已完成下一知识点关系
Cache, Session, Pub/SubWebSocketPub/Sub 事件推送到前端
ConnectionWebSocket 连接管理
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 方案

学习导航

上一篇:Redis 数据类型 · 对应练习 · 下一篇:WebSocket 基础