From 8f1008793938e701883979020185eb115b57df2e Mon Sep 17 00:00:00 2001 From: hhs <386998068@qq.com> Date: Thu, 28 May 2026 12:41:37 +0800 Subject: [PATCH] vault backup: 2026-05-28 12:41:37 --- hhs/Redis/09-高级特性.md | 351 +++++++++++++++++++++++------- hhs/Redis/17-Go项目集成Redis.md | 363 ++++++++++++++++++++++++++++++++ hhs/Redis/README.md | 2 +- 3 files changed, 641 insertions(+), 75 deletions(-) create mode 100644 hhs/Redis/17-Go项目集成Redis.md diff --git a/hhs/Redis/09-高级特性.md b/hhs/Redis/09-高级特性.md index 110cc67..c575d72 100644 --- a/hhs/Redis/09-高级特性.md +++ b/hhs/Redis/09-高级特性.md @@ -7,10 +7,36 @@ create time: 2026-05-15 18:15 ## 概述 -掌握事务、Lua 脚本和 Pipeline,是 Redis 从「基础缓存」迈向「生产级系统」的关键。本节逐一剖析其原理、场景和陷阱。 +当你用 Redis 做的不只是「存一个值、取一个值」时——比如转账时要同时扣 A 加 B、限流时要「先判断再计数」——单条命令就不够用了。你需要**把多条命令打包**,确保它们要么一起成功,要么逻辑上不出错。 + +本节围绕这个核心问题,逐一介绍三种「打包」工具: +- **事务(MULTI/EXEC)**:保证命令按顺序执行、不被插队,但不能根据前一条结果决定下一步 +- **Lua 脚本**:真正的「逻辑原子性」——读、判断、写一气呵成 +- **Pipeline**:纯性能优化——批量发送命令,减少网络往返 + +此外还会介绍 Pub/Sub 和 Stream 这两种消息传递机制,以及 Keyspace Notifications 键事件监听。 ## 一、事务(MULTI / EXEC) +### 先想一个问题 + +假设你要给两个计数器各加 1。如果直接逐条发送 `INCR counter`、`INCR other_counter`,在高并发下可能出现这种情况: + +```mermaid +sequenceDiagram + participant A as "Client A" + participant B as "Client B" + participant S as "Redis" + A->>S: "INCR counter 得到 1" + B->>S: "INCR counter 得到 2 (插队了!)" + B->>S: "INCR other_counter 得到 1" + A->>S: "INCR other_counter 得到 2" +``` + +两个命令之间被别的客户端**插队**了。虽然最终结果可能没问题,但在更复杂的场景(比如「读余额→判断→扣款」)中,插队会导致严重 bug。 + +**事务就是用来解决这个问题的:把一组命令打包,保证它们连续执行、中间不被打断。** + ### 基本用法 ```bash @@ -27,15 +53,18 @@ pipe.Set(ctx, "key", "value", 0) results, _ := pipe.Exec(ctx) ``` -> [!TIP] Go 的 Pipeline vs TxPipeline -> `rdb.Pipeline()` 只批量发送命令,不加事务包裹;`rdb.TxPipeline()` 会在批量命令前后自动加上 `MULTI` / `EXEC`,等价于上文的事务语义。在 Cluster 模式下,go-redis 会按 slot 分组,对每个节点单独发送 MULTI/EXEC(因此同一 pipeline 内的 key 必须在相同 hash slot)。两者区别: +> [!TIP] Go 的 Pipeline vs TxPipeline —— 简单理解 +> - `rdb.Pipeline()` = 批量发送,**不加锁**,谁都能插队 → 纯粹为了快 +> - `rdb.TxPipeline()` = 批量发送,**前后包上 MULTI/EXEC** → 保证连续执行 +> +> 在 Cluster 模式下,go-redis 会按 slot 分组,对每个节点单独发 MULTI/EXEC(所以同一 pipeline 内的 key 必须在相同 hash slot,否则事务会拆开)。 -| 维度 | MULTI/EXEC | Pipeline | +| 维度 | MULTI/EXEC(TxPipeline) | Pipeline | |------|-----------|----------| -| 原子性 | ✅ 整体执行,中途不插队 | ❌ 命令独立发送执行 | -| 错误处理 | 编译期错误全部拒绝;运行时错误逐个记录 | 每条返回各自结果 | -| Cluster 支持 | ❌ | ✅(但需同 key tag) | -| 性能 | 同 Pipeline | 批量发送减少 RTT | +| 原子性 | ✅ 连续执行,中途不插队 | ❌ 命令独立执行,可能被插队 | +| 错误处理 | 语法错 → 整批拒绝;运行时错 → 只影响当前命令 | 每条命令各自报错 | +| Cluster | ❌ 不支持跨 slot 事务 | ✅ 自动按 slot 分发 | +| 适用场景 | 需要保证执行顺序 | 只求快,不关心顺序 | ### 事务的错误处理 @@ -53,10 +82,10 @@ EXEC ```bash MULTI -SET k 123 # OK, queued -INCR k # OK, queued (先 SET 再 INCR 是合法的) +SET k "hello" # OK, queued +INCR k # OK, queued(语法合法,EXEC 前不会检查类型) EXEC -# → [OK, 124] — 两个都成功了 +# → [OK, ERR value is not an integer...] — SET 成功,INCR 因类型不匹配报错 ``` ```mermaid @@ -84,9 +113,25 @@ sequenceDiagram ## 二、Lua 脚本 -### 为什么需要 Lua? +### 为什么需要 Lua?——MULTI 的局限 -核心目的:**原子性操作多个命令**。MULTI 只能保证"不插队",Lua 能保证"逻辑原子性"(读-判断-写一气呵成)。 +MULTI 能保证命令「不插队」,但它有一个致命缺陷:**你无法根据前一条命令的结果决定下一步操作**。 + +比如你想实现「如果余额 ≥ 100 就扣款」: +1. `GET balance` → 拿到 50 +2. (根据结果判断:50 < 100,不应该扣款) +3. 但在 MULTI 里,`DECRBY balance 100` 已经在排队了——**你无法在 EXEC 之前取消它** + +> [!QUESTION] 怎么办? +> 答案是:把「读 + 判断 + 写」放到**一个脚本**里交给 Redis 执行。这就是 Lua 脚本的核心价值——**逻辑原子性**:整个脚本在 Redis 主线程上一口气跑完,中间不会被任何其他命令插入。 + +简单对比一下两者的原子性: + +| | MULTI/EXEC | Lua 脚本 | +|---|-----------|---------| +| 保证什么 | 命令**连续执行**,不被插队 | 整个脚本**一气呵成**,包含条件判断 | +| 不能做什么 | 根据中间结果做判断 | — | +| 一句话 | 「按顺序执行,但不看中间结果」 | 「读-判断-写,一步到位」 | ```bash EVAL "return redis.call('GET', KEYS[1])" 1 mykey @@ -129,6 +174,14 @@ script := redis.NewScript(` return 0 -- 被其他客户端持有 `) +// 解锁脚本:只有持有者才能释放,防止误删别人的锁 +unlockScript := redis.NewScript(` + if redis.call("get", KEYS[1]) == ARGV[1] then + return redis.call("del", KEYS[1]) + end + return 0 -- 锁不属于自己,拒绝释放 +`) + // KEYS[1] = lock key, ARGV[1] = TTL 秒数, ARGV[2] = 唯一令牌(如 UUID) ok, err := script.Run(ctx, rdb, []string{"lock:order:" + orderId}, lockTTL, clientToken).Int() @@ -144,12 +197,16 @@ if ok == 1 { > [!tip] 分布式锁深入阅读 > 本节仅展示 Lua 实现分布式锁的核心代码。关于安全解锁、锁续期(Watchdog)、Redlock 算法、可重入锁等进阶内容,详见 [[hhs/Redis/09-高级特性/分布式锁]]。 -### 经典场景 2:限流器(令牌桶简化版) +### 经典场景 2:限流器(固定窗口计数器) ```lua --- rate_limit.lua -local limit = tonumber(ARGV[1]) -- 每分钟上限 -local window = tonumber(ARGV[2]) -- 窗口秒数 +-- rate_limit.lua(固定窗口计数器 —— Lua 脚本本身就是原子的,无需 pipeline) +-- KEYS[1] = 计数器 key +-- ARGV[1] = 每窗口上限 +-- ARGV[2] = 窗口秒数 + +local limit = tonumber(ARGV[1]) +local window = tonumber(ARGV[2]) local current = tonumber(redis.call('get', KEYS[1]) or '0') @@ -157,10 +214,11 @@ if current >= limit then return 0 -- 超限,拒绝 end -local pipe = redis.pipeline() -pipe.incr(KEYS[1]) -pipe.expire(KEYS[1], window) -pipe.exec() +-- INCR 返回自增后的新值,首次写入时才设置 TTL(避免每次请求重置窗口) +local new_count = redis.call('incr', KEYS[1]) +if new_count == 1 then + redis.call('expire', KEYS[1], window) +end return 1 -- 通过 ``` @@ -199,18 +257,36 @@ EVALSHA numkeys args... # 只传 SHA1,省带宽 ### Lua 沙箱限制 -```lua --- ❌ 不可用:非确定性函数(传统 Redis < 7.2) -math.random() -- 无法预测结果 -os.date() -- 依赖系统时间 +Redis 为了让 Lua 脚本在**主从同步**和 **AOF 重放**时结果一致,在沙箱中屏蔽了部分标准库函数。核心原则:**「给同样的输入,永远得到同样的输出」的函数才能用**。 --- ✅ 推荐做法:在客户端生成随机值,通过 ARGV 传入 --- local token = client.generateUUID() --- redis.call('set', KEYS[1], ARGV[1]) -- 带 token 写入 +```lua +-- ❌ 不可用:非确定性函数 +math.random() -- 依赖内部 PRNG 状态,跨平台/版本结果不同 +os.time() -- 依赖系统时间 +os.getenv() -- 依赖环境变量 + +-- ✅ 推荐做法:在客户端生成随机值/时间戳,通过 ARGV 传入 +-- local token = ARGV[2] -- 由应用层生成的 UUID +-- local now = ARGV[3] -- 由应用层传入的 unix 时间戳 +-- redis.call('set', KEYS[1], ARGV[1], 'EX', ARGV[4]) ``` > [!NOTE] 关于 Lua 沙箱中的数学函数 -> `math.max`、`math.min`、`math.log`、`math.abs` 等纯函数在沙箱中可用——它们的输出完全由输入决定,属于确定性函数。但 `math.random` **始终不可用**:即使设置相同 seed,伪随机序列也属于非确定性行为,会破坏主从一致性和 AOF 重放。随机值必须由客户端生成后通过 ARGV 传入。 +> - ✅ `math.max`、`math.min`、`math.abs`、`math.floor` 等——纯函数,输入确定输出就确定 +> - ❌ `math.random`——虽然同一平台 + 同一 seed 理论上可复现,但 **Lua 5.1 PRNG 的实现在不同平台(Linux vs macOS、不同 libc 版本)之间不一致**,会导致主节点返回 42 而从节点重放得到 87,数据分裂 +> +> 所以随机值必须在客户端生成,通过 ARGV 传进去。 + +> [!NOTE] Redis 7.0+ 两种沙箱模式 +> Redis 7.0 引入了更严格的 Lua 沙箱。默认行为(等效于 `newdata-safe` 模式): +> +> | 模式 | 行为 | 说明 | +> |------|------|------| +> | Legacy | 仅屏蔽非确定性函数 | 与 Redis < 7.0 兼容,标准库大部分可用 | +> | Newdata-safe(默认) | **仅允许 `redis.*` 和 `math`/`table`/`string`/`cjson`/`cmsgpack`** | 屏蔽 `loadfile`、`os.*`、`io.*`、`pcall` 等,安全性更高 | +> +> 同时 Redis 7.0 改变了 Lua 脚本的**复制方式**:默认将脚本产生的写命令(而非脚本本身)发送给从节点,从而避免从节点重复执行脚本,提升主从一致性。 +> 生产环境建议保持默认——它能防止恶意脚本访问系统资源,并减少主从不一致的风险。 > [!QUESTION] 为什么 Lua 脚本不能有随机函数? > Lua 脚本在主线程中同步执行,如果引入随机性或时间依赖,会导致两个严重问题: @@ -218,16 +294,62 @@ os.date() -- 依赖系统时间 > 2. **AOF 重放不可靠**:重启后从 appendonly.aof 重放脚本,结果与上次不同,状态机紊乱。 > 因此 Redis 在 Lua 环境中屏蔽了所有非确定性 API,保证「同样的输入一定得到同样的输出」。 +### Lua 脚本最佳实践 + +Lua 脚本在 Redis **主线程**上同步执行——脚本运行期间,所有其他命令都会被阻塞。因此: + +| 实践 | 说明 | +|------|------| +| 脚本尽量短小 | 只封装「必须原子执行」的逻辑,不要在里面做大循环 | +| 控制执行时间 | `lua-time-limit`(默认 5s)控制超时;超时后其他客户端可发送 `SCRIPT KILL`,**但前提是脚本尚未执行写操作**——若已写入,只能 `SHUTDOWN NOSAVE` | +| 脚本大小适度 | 超大脚本(几十 MB)会导致 EVAL 传输慢 + 编译耗时,应拆分到客户端或使用 EVALSHA 复用 | +| 复杂计算放客户端 | Lua 只负责「读-判断-写」,数据聚合/排序等重计算应在应用层完成 | +| 优先 EVALSHA | 避免每次都传输完整脚本,利用 SHA1 缓存减少带宽 | + +> [!WARNING] 生产环境踩坑 +> 曾有人在 Lua 里遍历几千个 key 做 `HGETALL`,导致 Redis 阻塞数秒——其他所有请求全部超时。**记住:Lua 脚本 = 临界区代码,越快越好。** 如果脚本已执行了写操作(`SET`/`DEL` 等),`SCRIPT KILL` 将失效,只能强制 `SHUTDOWN NOSAVE`——这在生产环境意味着数据丢失风险。 + +### Lua 脚本调试 + +开发阶段,可以用 `redis-cli --eval` 直接在命令行测试 Lua 脚本,无需启动应用: + +```bash +# 语法:redis-cli --eval <脚本文件> , +# 注意逗号两边必须有空格,左边是 KEYS,右边是 ARGV +redis-cli --eval rate_limit.lua mykey , 100 60 +``` + +> [!TIP] `SCRIPT DEBUG` 模式 +> 在 redis-cli 中执行 `SCRIPT DEBUG SYNC` 后,接下来执行的 `EVAL` 会进入调试模式——你可以在脚本中用 `redis.log(redis.LOG_WARNING, variable)` 打印变量值,日志会输出到 Redis 日志文件。调试完毕后执行 `SCRIPT DEBUG NO` 关闭。 +> +> 生产环境务必关闭调试模式,它会让脚本**同步执行**(正常是异步),显著影响性能。 + ## 三、Pipeline ### 为什么需要 Pipeline? -Redis 协议本质是 request/response,每个命令一次网络往返(RTT)。假设客户端到服务器 RTT 为 1ms,执行 1000 条命令: +前面讲的事务和 Lua 解决的是「正确性」问题——保证命令不出错。但还有一个完全不同的问题:**性能**。 -| 方式 | RTT 次数 | 总耗时估算 | -|------|---------|-----------| -| 单次发送 | 1000 | ~1s | -| Pipeline (batch=200) | 5 次 | ~5ms | +Redis 的通信协议是「你一句我一句」——客户端发一条命令,等服务器返回结果,再发下一条。每次这样的网络来回叫做 **RTT(Round-Trip Time,网络往返时间)**。假设 RTT = 1ms,执行 1000 条命令: + +```mermaid +sequenceDiagram + participant C as "Client" + participant S as "Redis" + Note over C,S: "逐条发送 1000 次 RTT 约 1s" + C->>S: "SET key1 val1" + S-->>C: "OK" + C->>S: "SET key2 val2" + S-->>C: "OK" + C->>S: "... 还有 998 条" +``` + +**Pipeline 的思路很简单:把多条命令攒一批一起发,服务器也一批返回。** 网络只往返一次(或少数几次),而不是每条命令都来回跑。 + +| 方式 | RTT 次数 | 总耗时估算 | 比喻 | +|------|---------|-----------|------| +| 逐条发送 | 1000 | ~1s | 寄 1000 封信,每封等回信再寄下一封 | +| Pipeline (batch=200) | 5 次 | ~5ms | 攒 200 封信一起寄,回信也一批到 | ### Pipeline vs TxPipeline @@ -270,6 +392,13 @@ if err == redis.TxFailedErr { } ``` +> [!QUESTION] WATCH 是什么?为什么需要它? +> `WATCH` 是 Redis 提供的**乐观锁**机制——在事务开始前「监视」一个或多个 key,如果这些 key 在 EXEC 之前被其他客户端修改了,整个事务会被拒绝(返回 `TxFailedErr`),你需要自行重试。 +> +> 回顾上面的转账例子:`WATCH account:alice` → 读余额 → 判断 → `TxPipeline` 扣款。如果在你读余额之后、扣款之前,别人修改了 `account:alice`,`EXEC` 会失败,从而避免「余额不足却被扣款」的竞态条件。 +> +> **一句话**:WATCH 让 MULTI/EXEC 具备了「检查-执行」的能力,弥补了事务不能读中间结果的短板。 + > [!WARNING] Pipeline 的注意事项 > - 单个 pipeline 内的命令可能跨多个节点(Cluster 模式下会分发执行) > - 避免在 pipeline 中混用 WATCH/MULTI(Cluster 不支持) @@ -280,23 +409,32 @@ if err == redis.TxFailedErr { ```go func BatchMigrate(ctx context.Context, rdb *redis.Client, srcKeys []string) error { const batchSize = 100 - pipe := rdb.Pipeline() - + for i := 0; i < len(srcKeys); i += batchSize { batch := srcKeys[i : min(i+batchSize, len(srcKeys))] - - for _, key := range batch { - val, _ := rdb.Get(ctx, key).Result() - newKey := strings.TrimPrefix(key, "old:") - pipe.Set(ctx, newKey, val, 0) + + // 第一轮:用 Pipeline 批量 GET(而非逐条 Get,否则失去 Pipeline 意义) + getPipe := rdb.Pipeline() + getCmds := make([]*redis.StringCmd, len(batch)) + for j, key := range batch { + getCmds[j] = getPipe.Get(ctx, key) } - - _, err := pipe.Exec(ctx) - if err != nil { + getPipe.Exec(ctx) + + // 第二轮:用 Pipeline 批量 SET + setPipe := rdb.Pipeline() + for j, key := range batch { + val, err := getCmds[j].Result() + if err != nil { + continue // 跳过不存在的 key + } + newKey := strings.TrimPrefix(key, "old:") + setPipe.Set(ctx, newKey, val, 0) + } + if _, err := setPipe.Exec(ctx); err != nil { return err } - pipe.Close() - pipe = rdb.Pipeline() // 重建 pipeline 避免内存泄漏 + // Exec 后 pipeline 已清空,无需 Close() / 重建 } return nil } @@ -304,6 +442,13 @@ func BatchMigrate(ctx context.Context, rdb *redis.Client, srcKeys []string) erro ## 四、Pub/Sub —— 发布订阅 +前面讲的事务、Lua、Pipeline 都是「一个客户端跟 Redis 对话」。但有些场景需要**多个客户端之间通信**——比如 WebSocket 服务要广播消息给所有在线用户,或者配置中心更新后通知所有微服务。 + +Pub/Sub 就是 Redis 内置的「广播电台」:一个客户端发布消息,所有订阅了对应频道的客户端都能收到。 + +> [!WARNING] 先记住最关键的一点 +> Pub/Sub 是**即发即忘**的——消息不会存储,订阅者如果下线了就收不到。如果你需要可靠投递(比如任务队列),请直接跳到后面的 [[#Stream]]。 + ### 基础用法 ```bash @@ -337,9 +482,52 @@ rdb.Publish(ctx, "channel:newposts", `{"id":42,"title":"Redis Tips"}`) > - **消息不持久化**:订阅者下线就丢失,重连不会收到离线期间的消息 > - **适合场景**:WebSocket 广播、实时通知、配置变更推送 +### Sharded Pub/Sub(Redis 7.0+) + +普通 Pub/Sub 在集群模式下有个隐患:**所有消息都由接收 PUBLISH 命令的单个节点转发给所有订阅者**,该节点容易成为瓶颈。 + +Redis 7.0 引入了 Sharded Pub/Sub,核心改进是**消息按 slot 分发到各节点**,避免单点压力: + +```bash +# 分片发布(消息只会发送到该 slot 所在节点的订阅者) +SSUBSCRIBE channel:order +SPUBLISH channel:order '{"id":42}' + +# 普通发布(广播到所有节点的订阅者) +SUBSCRIBE channel:order +PUBLISH channel:order '{"id":42}' +``` + +| 维度 | 普通 Pub/Sub | Sharded Pub/Sub | +|------|-------------|----------------| +| 集群消息路由 | 单节点广播 → 所有节点 | 按 slot → 仅目标节点 | +| 适用场景 | 需要全局广播 | 高吞吐、按 key 分区消费 | +| 最低版本 | 所有版本 | Redis 7.0+ | + +> [!QUESTION] 什么时候该用 Sharded Pub/Sub? +> 如果你的频道消息量大、且订阅者只需要消费属于自己 slot 的消息(比如按用户 ID 分片),Sharded Pub/Sub 能显著降低单节点压力。但如果需要**全局广播**(所有节点都收到),还是用普通 Pub/Sub。 + ### Stream —— 可靠的替代方案 -Stream 是 Redis 5.0+ 引入的消息队列数据结构,弥补了 Pub/Sub 的可靠性缺陷。核心设计围绕 **Consumer Group(消费者组)**: +如果 Pub/Sub 是「聊天群」(消息刷过去就没了),那 Stream 就是「消息队列」——**消息会持久化,消费者必须确认收到,挂了可以重新投递**。 + +Stream 是 Redis 5.0+ 引入的数据结构,核心设计围绕 **Consumer Group(消费者组)**: + +```mermaid +flowchart LR + P["Producer"] -->|"XADD: 写入消息"| S["Stream (持久化存储)"] + S -->|"XREADGROUP: 读取并标记 pending"| C1["Consumer A
处理消息"] + S -->|"XREADGROUP: 读取并标记 pending"| C2["Consumer B
处理消息"] + C1 -->|"XACK: 确认处理完成"| S + C2 -->|"XCLAIM: 接管超时消息"| S +``` + +核心流程三步走: +1. **生产者**用 `XADD` 写入消息 +2. **消费者组**中的消费者用 `XREADGROUP` 读取——消息会被标记为 pending(待确认) +3. **消费者**处理完后用 `XACK` 确认——确认后消息才算真正消费完 + +如果消费者中途挂了,pending 消息不会丢失,其他消费者可以用 `XCLAIM` 接管。 ```bash # 生产者写入 @@ -366,18 +554,6 @@ XACK mystream group1 > [!QUESTION] 如果消费者挂了怎么办? > 好消息会被移到 **pending 列表**。用 `XPENDING mystream group1` 查看,用 `XCLAIM` 将 pending 消息转移给其他 consumer 重新处理。这就是 Stream 比 Pub/Sub 可靠的地方:**有状态、可追责**。 -```mermaid -flowchart LR - P["Producer"] -->|"XADD"| S["Stream (持久化)"] - S -->|"读取并标记 pending"| G1["Consumer Group 1
consumer-a"] - S -->|"读取并标记 pending"| G2["Consumer Group 2
consumer-b"] - G1 -->|"XACK"| S - - style S fill:#dfd,stroke:#090 - style G1 fill:#ddf,stroke:#669 - style G2 fill:#ddf,stroke:#669 -``` - > [!TIP] Pub/Sub vs Stream 选型 > | 场景 | 推荐 | > |------|------| @@ -392,19 +568,35 @@ flowchart LR ### 原理与配置 -Redis 会在每个 key 发生变更时发送通知,订阅者通过 Pub/Sub 协议接收。开启方式: +简单来说,Keyspace Notifications 就是 Redis 的「事件总线」——当某个 key 发生变化(被删除、过期、被修改等),Redis 会通过 Pub/Sub 发出一条通知。 + +> [!TIP] 最常见的用途 +> 监听 key 过期事件。比如:用户 session 过期时自动清理关联数据、分布式锁到期时告警。 + +开启方式很简单,在 `redis.conf` 中设置(或运行时 `CONFIG SET`): ```conf -# redis.conf(或运行时 CONFIG SET) notify-keyspace-events KEA -# K = Keyspace events → 发布到 __keyspace@__ prefix -# E = Keyevent events → 发布到 __keyevent@__ prefix -# A = All events 简写 → 覆盖 g/l/s/h/z/x/e(不包含 K/E/d/t/n) -# g = DEL/EXPIRE/RENAMExxx 等通用命令 -# l = list, s = set, h = hash, z = sorted-set -# x = expired, e = evicted (maxmemory 淘汰) ``` +这一串字母是事件类型的组合,不用全部记住,按需选用即可: + +| 字母 | 含义 | 举例 | +|------|------|------| +| `K` | Keyspace 事件(按 key 名通知) | 有人操作了 `user:123` | +| `E` | Keyevent 事件(按事件类型通知) | 有 key 过期了 | +| `A` | All events 简写,覆盖下方所有类型 | 一次性开启全部 | +| `x` | 过期事件 | `EXPIRE` 到期触发 | +| `e` | 淘汰事件 | 内存满时被驱逐 | +| `g` | 通用命令 | `DEL`、`EXPIRE`、`RENAME` 等 | +| `l/s/h/z` | 数据结构相关 | list/set/hash/sorted-set 操作 | + +> [!NOTE] K 和 E 的区别——通知「谁」vs 通知「什么」 +> - `__keyspace@0__:` → "有人动了 `user:123`"(以 key 为中心) +> - `__keyevent@0__:expired` → "有 key 过期了"(以事件为中心) +> +> 两者是同一件事的两种视角,按需选用。**生产环境最常用 `Ex`(只监听过期事件)。** + ```bash # 订阅过期事件(Keyevent 模式:更精确) PSUBSCRIBE __keyevent@0__:expired @@ -434,7 +626,24 @@ for msg := range ch { ## 六、选型指南——如何选择? -面对一个需求,到底该用事务、Lua、Pipeline 还是 Stream?以下是决策流程: +面对一个需求,到底该用事务、Lua、Pipeline 还是 Stream? + +### 决策心法——三个问题搞定 + +拿到需求后,先问自己三个问题: + +1. **我需要「原子性」吗?** → 如果只是批量写入追求性能,Pipeline 就够了 +2. **我需要「读-判断-写」吗?** → 如果答案是 yes,只有 Lua 能做到真正的逻辑原子性 +3. **我在做「消息传递」吗?** → 实时广播用 Pub/Sub,需要可靠投递用 Stream + +> [!TIP] 一句话记忆 +> - **Pipeline**:只求快,不关心顺序 → 「快递员一次送 200 个包裹」 +> - **MULTI/EXEC**:要求连续执行,不被插队 → 「排队结账,不允许别人插队」 +> - **Lua 脚本**:需要根据中间结果做判断 → 「医生先检查再开药,不能分开」 +> - **Pub/Sub**:实时广播,丢了就丢了 → 「微信群消息」 +> - **Stream**:消息必须可靠送达 → 「挂号信,签收才算数」 + +下面是完整的决策流程: ```mermaid flowchart TD @@ -460,12 +669,6 @@ flowchart TD | 任务队列、事件溯源 | Stream | 持久化 + Consumer Group | | 实时通知、WebSocket 广播 | Pub/Sub | 无状态、低延迟 | -> [!TIP] 一句话总结 -> - **Pipeline**:省网络,不保证原子性 -> - **MULTI**:保证顺序,不保证逻辑原子 -> - **Lua**:真正的原子性,但阻塞主线程(脚本要短小精悍) -> - **Stream**:需要消息可靠投递时的首选 - ## 关联笔记 - [[hhs/Redis/09-高级特性/分布式锁]] — 分布式锁完整解析(Redlock、Watchdog、可重入锁) diff --git a/hhs/Redis/17-Go项目集成Redis.md b/hhs/Redis/17-Go项目集成Redis.md new file mode 100644 index 0000000..9013110 --- /dev/null +++ b/hhs/Redis/17-Go项目集成Redis.md @@ -0,0 +1,363 @@ +--- +tags: [Redis, Go, go-redis, 缓存, 后端] +create time: 2026-05-28 10:55 +--- + +# Go 项目集成 Redis + +## 概述 + +Go 生态中最主流的 Redis 客户端是 **go-redis**(`github.com/redis/go-redis/v9`),它支持 Redis 6.0+ 的全部命令、Cluster、Sentinel、Pub/Sub、Lua 脚本、Pipeline 等能力,社区活跃度和文档质量均优于早期的 redigo。本文以 go-redis 为主线,覆盖从初始化到生产实践的完整链路。 + +> [!QUESTION] go-redis vs redigo,怎么选? +> | 维度 | go-redis | redigo | +> |------|----------|--------| +> | API 风格 | 命令式,返回值类型明确 | `Do()` 返回 `interface{}`,需手动断言 | +> | 连接池 | 内置,参数可调 | 需手动封装 `Pool` | +> | Cluster / Sentinel | 原生支持 | 需自行实现 | +> | 维护状态 | 持续更新(v9) | 基本停更 | +> +> **结论**:新项目直接用 go-redis,除非有历史包袱。 + +## 一、客户端初始化与连接池 + +### 1.1 基础连接 + +```go +import ( + "context" + "errors" + "fmt" + "os" + "time" + + "github.com/redis/go-redis/v9" +) + +// 单机模式 +rdb := redis.NewClient(&redis.Options{ + Addr: "localhost:6379", + Password: os.Getenv("REDIS_PASSWORD"), // 生产环境从环境变量读取 + DB: 0, +}) +``` + +### 1.2 连接池参数调优 + +go-redis 内置连接池,**不需要**再手动包一层。关键参数如下: + +```go +rdb := redis.NewClient(&redis.Options{ + Addr: "localhost:6379", + PoolSize: 100, // 连接池最大连接数,默认 10*runtime.GOMAXPROCS + MinIdleConns: 10, // 最小空闲连接,避免冷启动时频繁创建连接 + DialTimeout: 5 * time.Second, // 建立 TCP 连接的超时 + ReadTimeout: 3 * time.Second, // 读超时(含从池中获取连接的等待时间) + WriteTimeout: 3 * time.Second, // 写超时 +}) +``` + +> [!TIP] PoolSize 和 MinIdleConns 的关系 +> - `PoolSize` 是硬上限——同时最多有这么多 TCP 连接到 Redis +> - `MinIdleConns` 是保底——即使没有请求,也会保持这么多空闲连接 +> - 连接不够时,请求会**阻塞等待**(由 `ReadTimeout` / `WriteTimeout` 控制超时) +> - go-redis v9 **没有** `MaxIdleConns` 和 `PoolTimeout`,空闲连接回收由内部 `connMaxIdleTime`(默认 5 分钟)管理 + +> [!TIP] 如何确定 PoolSize? +> - **公式参考**:`PoolSize ≈ 预估并发数 × 单次 Redis 操作耗时(s)`(Little's Law:任意时刻同时占用连接的期望数) +> - 例如并发 1000,每次操作 10ms → `1000 × 0.01 = 10`,再留 2~3 倍余量 → `PoolSize = 30` +> - 监控 `redis.PoolStats()` 中的 `Hits`、`Misses`、`Timeouts` 来动态调整 + +### 1.3 Sentinel 模式 + +```go +rdb := redis.NewFailoverClient(&redis.FailoverOptions{ + MasterName: "mymaster", + SentinelAddrs: []string{"sentinel-1:26379", "sentinel-2:26379", "sentinel-3:26379"}, + Password: os.Getenv("REDIS_PASSWORD"), + SentinelPassword: os.Getenv("SENTINEL_PASSWORD"), + PoolSize: 100, + MinIdleConns: 10, +}) +``` + +### 1.4 Cluster 模式 + +```go +rdb := redis.NewClusterClient(&redis.ClusterOptions{ + Addrs: []string{ + "node-1:6379", "node-2:6379", "node-3:6379", + "node-4:6379", "node-5:6379", "node-6:6379", + }, + Password: os.Getenv("REDIS_PASSWORD"), + PoolSize: 100, + MinIdleConns: 10, + ReadOnly: true, // 从节点可读,分担读压力 + RouteByLatency: true, // 自动路由到延迟最低的节点 +}) +``` + +```mermaid +flowchart TD + A["Go Application"] --> B{"连接模式?"} + B -->|"单机/开发"| C["redis.NewClient"] + B -->|"主从+自动故障转移"| D["redis.NewFailoverClient"] + B -->|"水平扩展/分片"| E["redis.NewClusterClient"] + C --> F["redis.Options"] + D --> G["redis.FailoverOptions"] + E --> H["redis.ClusterOptions"] + F --> I["调优连接池参数"] + G --> I + H --> I +``` + +## 二、基础操作示例 + +go-redis 的命令命名与 Redis 命令一一对应,遵循 `rdb.Xxx(ctx, args...)` 的统一模式。 + +### 2.1 String + +```go +ctx := context.Background() + +// SET key value EX 60 +err := rdb.Set(ctx, "user:1001:name", "张三", 60*time.Second).Err() + +// GET key +name, err := rdb.Get(ctx, "user:1001:name").Result() +// 不存在时返回 redis.Nil 错误,而非空字符串 +if errors.Is(err, redis.Nil) { + fmt.Println("key 不存在,回源 DB") +} +``` + +> [!WARNING] 区分 "key 不存在" 和 "其他错误" +> go-redis 用 `redis.Nil` 表示 key 不存在,不要用 `err != nil` 一概而论——这会把"不存在"误判为"出错"。 + +### 2.2 Hash + +```go +// HSET user:1001 name "张三" age 25 +rdb.HSet(ctx, "user:1001", map[string]interface{}{ + "name": "张三", + "age": 25, +}) + +// HGETALL user:1001 +user, err := rdb.HGetAll(ctx, "user:1001").Result() +``` + +### 2.3 Sorted Set + +```go +// ZADD leaderboard 95.5 "player:A" +rdb.ZAdd(ctx, "leaderboard", redis.Z{ + Score: 95.5, + Member: "player:A", +}) + +// ZREVRANGE leaderboard 0 9 WITHSCORES(取 Top 10) +top10, _ := rdb.ZRevRangeWithScores(ctx, "leaderboard", 0, 9).Result() +``` + +## 三、Pipeline 批量操作 + +逐条发送命令,每条都要等一次 RTT。**Pipeline** 把多条命令打包成一次请求发送,响应也一次性读回,大幅减少网络往返。 + +```go +// 逐条写入 1000 个 key:1000 次 RTT ❌ +// Pipeline 写入 1000 个 key:1 次 RTT ✅ + +pipe := rdb.Pipeline() +for i := 0; i < 1000; i++ { + pipe.Set(ctx, fmt.Sprintf("key:%d", i), i, 0) +} +cmds, err := pipe.Exec(ctx) // 一次发送,批量执行 +``` + +> [!QUESTION] Pipeline 和 Lua 脚本的区别? +> - **Pipeline**:多条命令打包发送,但**不保证原子性**——其他客户端的命令可能穿插在中间 +> - **Lua 脚本**:在 Redis 服务端**原子执行**,适合需要"要么全成功,要么全失败"的场景(如分布式锁解锁) +> - 如果只是"批量读/写",Pipeline 性能更优;如果需要原子语义,用 Lua + +## 四、分布式锁 + +### 4.1 加锁:SET NX EX + +```go +// SET lock:order:1001 $uuid NX EX 10 +ok, err := rdb.SetNX(ctx, "lock:order:1001", uuid, 10*time.Second).Result() +if !ok { + // 获取锁失败,其他协程/进程已持有 + return fmt.Errorf("锁已被持有") +} +defer unlock(ctx, rdb, "lock:order:1001", uuid) // 确保释放 +``` + +### 4.2 解锁:Lua 原子判断 + 删除 + +解锁必须**先判断锁是否是自己加的**,再删除——这两步必须原子执行,否则会误删别人的锁。 + +```go +const unlockScript = ` +if redis.call("GET", KEYS[1]) == ARGV[1] then + return redis.call("DEL", KEYS[1]) +else + return 0 +end` + +func unlock(ctx context.Context, rdb *redis.Client, key, uuid string) error { + result, err := rdb.Eval(ctx, unlockScript, []string{key}, uuid).Int64() + if err != nil { + return err + } + if result == 0 { + return fmt.Errorf("锁已过期或不属于当前持有者") + } + return nil +} +``` + +```mermaid +sequenceDiagram + participant G1 as "Goroutine A" + participant R as "Redis" + participant G2 as "Goroutine B" + G1->>R: "SET lock:order NX EX 10 uuid-a" + R-->>G1: "OK" + G2->>R: "SET lock:order NX EX 10 uuid-b" + R-->>G2: "nil, 失败" + Note over G2: "等待或重试" + G1->>R: "Eval Lua: GET==uuid-a ? DEL : 0" + R-->>G1: "1, 解锁成功" + G2->>R: "SET lock:order NX EX 10 uuid-b" + R-->>G2: "OK" +``` + +> [!WARNING] Redlock 争议 +> Martin Kleppmann 在 [How to do distributed locking](https://martin.kleppmann.com/2016/02/08/how-to-do-distributed-locking.html) 中指出单节点 `SET NX EX` 在主从切换时可能不安全。如果你的业务对锁的安全性要求极高,建议: +> 1. 使用 Redlock 多节点方案(go-redis 的 `redsync` 库) +> 2. 或引入 fencing token + 下游服务校验 + +## 五、缓存模式实战:Cache-Aside + GORM + +### 5.1 Cache-Aside 流程 + +```mermaid +flowchart TD + subgraph "读流程" + R1["查询缓存"] --> R2{"命中?"} + R2 -->|"是"| R3["返回缓存数据"] + R2 -->|"否"| R4["查询 DB"] + R4 --> R5["回填缓存(SET EX)"] + R5 --> R3 + end + subgraph "写流程" + W1["更新 DB"] --> W2["删除缓存 Key"] + end +``` + +### 5.2 实现示例 + +```go +// 读:先缓存,miss 回源 DB +func GetUser(ctx context.Context, rdb *redis.Client, db *gorm.DB, id uint) (*User, error) { + key := fmt.Sprintf("user:%d", id) + + // 1. 查缓存 + cached, err := rdb.Get(ctx, key).Bytes() + if err == nil { + var u User + json.Unmarshal(cached, &u) + return &u, nil + } + if !errors.Is(err, redis.Nil) { + return nil, err // Redis 自身出错,降级到 DB + } + + // 2. 回源 DB + var user User + if err := db.First(&user, id).Error; err != nil { + return nil, err + } + + // 3. 回填缓存(TTL 随机抖动防雪崩) + ttl := 30*time.Minute + time.Duration(rand.Intn(300))*time.Second + data, _ := json.Marshal(user) + rdb.Set(ctx, key, data, ttl) + + return &user, nil +} +``` + +### 5.3 与 GORM 钩子联动 + +在写操作后自动删缓存,通过 GORM 的 `AfterUpdate` / `AfterDelete` 钩子实现。先用类型安全的 context key 注入 Redis 客户端: + +```go +// context key 用自定义类型避免碰撞 +type ctxKey struct{} + +func WithRedis(ctx context.Context, rdb *redis.Client) context.Context { + return context.WithValue(ctx, ctxKey{}, rdb) +} + +func RedisFromCtx(ctx context.Context) *redis.Client { + return ctx.Value(ctxKey{}).(*redis.Client) +} + +// 中间件注入:将 rdb 放入 Gin 的 context,GORM 通过 db.WithContext 传递 +func RedisMiddleware(rdb *redis.Client) gin.HandlerFunc { + return func(c *gin.Context) { + c.Set("db", db.WithContext(WithRedis(c.Request.Context(), rdb))) + c.Next() + } +} +``` + +```go +// 钩子中安全地取出 rdb +func (u *User) AfterUpdate(tx *gorm.DB) error { + rdb := RedisFromCtx(tx.Statement.Context) + key := fmt.Sprintf("user:%d", u.ID) + return rdb.Del(tx.Statement.Context, key).Err() +} +``` + +> [!WARNING] 别用裸字符串做 context key +> `context.Value("rdb")` 这种写法在多包协作时极易碰撞,且缺乏编译期类型检查。始终用**自定义 struct 类型**作为 key。 + +> [!TIP] 延迟双删 +> 在高并发写场景下,"先更新 DB,再删缓存"仍可能出现短暂不一致。**延迟双删**的做法是:第一次删缓存 → 更新 DB → sleep 500ms → 第二次删缓存。第二次删除清理了在 sleep 期间被其他请求回填的旧值。 + +## 六、Go 项目中 Redis 的典型架构位置 + +```mermaid +flowchart LR + Client["客户端请求"] --> Gin["Gin Router"] + Gin --> Handler["Handler 层"] + Handler --> Svc["Service 层"] + Svc --> Cache{"缓存层"} + Cache -->|"命中"| R["Redis"] + Cache -->|"未命中"| Repo["Repository 层"] + Repo --> DB["MySQL / PostgreSQL"] + DB -->|"回填"| R + subgraph "Session 共享" + Gin -->|"gin-contrib/sessions"| R + end +``` + +Redis 在 Go 项目中通常承担三个角色: + +1. **数据缓存**:Cache-Aside 模式,减轻数据库压力 +2. **Session 存储**:`gin-contrib/sessions/redis` 实现多实例共享会话 +3. **分布式协调**:分布式锁、限流器、排行榜等 + +## 关联笔记 + +- [[hhs/Redis/README]] — Redis 知识索引入口 +- [[hhs/Redis/10-缓存架构模式]] — Cache-Aside / 穿透 / 击穿 / 雪崩完整方案 +- [[hhs/Redis/09-高级特性]] — Lua 脚本、Pipeline、Pub/Sub 详解 +- [[hhs/Redis/11-运维与性能调优]] — 连接池监控、慢查询、大 Key 治理 +- [[hhs/GORM/09-钩子函数]] — GORM 生命周期与缓存同步模式 +- [[hhs/EXAM/Week05]] — Docker Compose 中 Redis 容器管理与 Session 共享 diff --git a/hhs/Redis/README.md b/hhs/Redis/README.md index 18adba0..2075b73 100644 --- a/hhs/Redis/README.md +++ b/hhs/Redis/README.md @@ -125,7 +125,7 @@ flowchart TD | 技术栈 | 集成方式 | 相关笔记 | |--------|---------|----------| -| Go | go-redis / redigo | — | +| Go | go-redis / redigo | [[hhs/Redis/17-Go项目集成Redis]] | | Java | lettuce / jedis | — | | GORM | 钩子函数清理/更新缓存 | [[hhs/GORM/09-钩子函数]] | | Docker Compose | depends_on 启动顺序 | [[hhs/EXAM/Week05]] |