diff --git a/hhs/Redis/01-安装与部署.md b/hhs/Redis/01-安装与部署.md index 4562d68..d29cd8e 100644 --- a/hhs/Redis/01-安装与部署.md +++ b/hhs/Redis/01-安装与部署.md @@ -250,8 +250,8 @@ flowchart TB ## 关联笔记 -- [[hhs/Redis/02-基础数据结构]] — String / Hash / List / Set / ZSet 的使用场景 +- [[hhs/Redis/02-核心数据类型]] — String / Hash / List / Set / ZSet 的使用场景 - [[hhs/Redis/04-RDB持久化]] — RDB 快照原理与配置 - [[hhs/Redis/05-AOF持久化]] — AOF 重写机制与 fsync 策略 - [[hhs/Redis/07-集群方案]] — Cluster 部署与迁移策略 -- [[hhs/Redis/08-常见问题排查]] — OOM、延迟 spike、连接数爆满的诊断思路 \ No newline at end of file +- [[hhs/Redis/11-运维与性能调优]] — OOM、延迟 spike、连接数爆满的诊断思路 \ No newline at end of file diff --git a/hhs/Redis/05-AOF持久化.md b/hhs/Redis/05-AOF持久化.md index 04f762e..a7af560 100644 --- a/hhs/Redis/05-AOF持久化.md +++ b/hhs/Redis/05-AOF持久化.md @@ -181,49 +181,12 @@ ls /var/lib/redis/appendonlydir/ ## 混合持久化(Hybrid Persistence) -混合持久化是 Redis 4.0 引入、Redis 7 默认开启的特性,它解决了 AOF 重写的一个尴尬问题:**重写时到底该生成 RDB 还是 AOF?** - -> [!question] 纯 AOF 重写的痛点 -> 纯 AOF 重写需要把每个 key 的当前状态都转成 SET/HSET 等命令写入文件。当实例有几十万个 key 时,这些命令序列化+解析的过程会显著拖慢恢复速度。而 RDB 是二进制格式,加载速度比逐条执行 AOF 命令快 10~100 倍。 +> [!tip] 混合持久化是 AOF 重写的"终极形态" +> Redis 4.0 引入、Redis 7 默认开启。核心思路:AOF **重写**时先写 RDB 全量快照作为"底子",再追加增量 AOF 命令作为"补丁"——恢复时先加载 RDB(秒级),再回放少量增量 AOF(补齐数据)。 > -> 混合持久化的答案很简单:**两种都用**。重写时先写 RDB 全量快照作为"底子",再追加一小段 AOF 增量命令作为"补丁"。 - -```conf -# redis.conf -aof-use-rdb-preamble yes # 开启混合持久化(默认值) -``` - -重写后生成的文件结构如下: - -```mermaid -flowchart TD - subgraph NEW_AOF["重写后的新 AOF 文件"] - direction LR - RDB["RDB 二进制块
(全量数据快照)"] --> AOF["AOF 增量块
(仅重写期间的新写操作)"] - end - - Fork["fork 子进程"] --> RDB - Fork --> AOF - RDB --> Merge["合并写入临时文件"] - AOF --> Merge - Merge --> Rename["原子 rename 替换旧文件"] - - style RDB fill:#bbf,stroke:#66c - style AOF fill:#dfd,stroke:#090 - style Merge fill:#ffd,stroke:#cc0 -``` - -**好处有两个:** -1. **恢复速度快**:启动时先加载 RDB 二进制快照(顺序读、无需解析命令语法),再回放少量 AOF 增量命令。即使 AOF 增量段有几十 MB,相比从头解析几十 GB 的纯 AOF 命令也快得多 -2. **文件体积小**:RDB 本身就是压缩的二进制格式,比等量数据的纯 AOF 文本命令小 70%~90% - -> [!warning] Redis 7 的多文件目录结构 -> Redis 7 不再使用单个 AOF 文件,而是采用多文件目录结构: -> - `base.rdb`:AOF 重写时生成的 RDB 全量快照 -> - `incr-*.aof`:增量 AOF 文件(记录重写之后的新写操作) -> - `manifest.aof`:文件清单,记录哪些文件属于同一个 AOF 实例 +> **在 AOF 视角下**:重写后生成的文件不再全是文本命令,而是 `RDB 二进制基线 + AOF 增量日志` 的混合结构。好处是恢复速度快(接近纯 RDB)且数据安全性高(接近纯 AOF)。 > -> 恢复时 Redis 会按 manifest 加载:先加载 `base.rdb`,再按顺序回放所有 `incr-*.aof`。如果删除了 `base.rdb`,恢复退化为纯 AOF 模式,速度会大幅下降。 +> 详细的文件结构、原理图解和配置方式,参见 [[hhs/Redis/04-RDB持久化#Redis 混合持久化(推荐)|混合持久化完整章节]]。 ## AOF vs RDB 对比 @@ -308,4 +271,4 @@ flowchart TD - [[hhs/Redis/04-RDB持久化]] — RDB 快照机制 - [[hhs/Redis/06-主从与哨兵]] — 哨兵监控中 AOF 状态检查 -- [[hhs/Redis/09-运维调优]] — AOF 文件膨胀治理 +- [[hhs/Redis/11-运维与性能调优]] — AOF 文件膨胀治理 diff --git a/hhs/Redis/06-主从与哨兵.md b/hhs/Redis/06-主从与哨兵.md index 501309f..4d32471 100644 --- a/hhs/Redis/06-主从与哨兵.md +++ b/hhs/Redis/06-主从与哨兵.md @@ -601,6 +601,6 @@ Sentinel 方案虽然解决了"高可用"问题,但它有明确的边界。理 ## 关联笔记 - [[hhs/Redis/07-集群方案]] — Cluster 架构(水平扩展) -- [[hhs/Redis/09-运维调优]] — 监控与告警配置 -- [[hhs/Redis/04-AOF持久化]] — AOF 与 RDB 结合的高可用策略 +- [[hhs/Redis/11-运维与性能调优]] — 监控与告警配置 +- [[hhs/Redis/05-AOF持久化]] — AOF 与 RDB 结合的高可用策略 - [[hhs/EXAM/Week05]] — Docker Compose 中的部署示例 diff --git a/hhs/Redis/07-集群方案.md b/hhs/Redis/07-集群方案.md index 8272d68..67a1ed5 100644 --- a/hhs/Redis/07-集群方案.md +++ b/hhs/Redis/07-集群方案.md @@ -262,13 +262,16 @@ flowchart LR ### 故障检测三阶段 -``` -┌─────────────┐ ┌─────────────┐ ┌─────────────┐ -│ PFAIL 潜在失效 │ → │ FAIL 确认失效 │ → │ 选举 Leader │ → Failover -│ (单个节点判定) │ │ (多数节点同意) │ │ (Replica 竞选) │ -└─────────────┘ └─────────────┘ └─────────────┘ - ↓ 超时 ↓ 多数派 ↓ Raft 式投票 - cluster-node-timeout 超过半数标记 最高优先级获胜利 +```mermaid +flowchart LR + A["PFAIL 潜在失效
单个节点判定"] -->|"超过 cluster-node-timeout"| B["FAIL 确认失效
多数 Master 同意"] + B -->|"Raft 式投票"| C["选举 Leader
Replica 竞选"] + C -->|"最高优先级获胜"| D["Failover
接管 slots"] + + style A fill:#fff3e0,stroke:#ff9800 + style B fill:#fce4ec,stroke:#e91e63 + style C fill:#e8eaf6,stroke:#3f51b5 + style D fill:#e8f5e9,stroke:#4caf50 ``` **详细流程:** @@ -376,18 +379,19 @@ let val = await cluster.get("key") ### 架构最佳实践 -``` -┌─────────────────────────────────────────────┐ -│ 你的应用服务 │ -└──┬──────────┬──────────┬──────────┬──────────┘ - │ │ │ │ -┌──▼────┐ ┌──▼────┐ ┌──▼────┐ ┌──▼────┐ -│Service A│ │Service B│ │Service C│ │Service D│ -│ :7001-7 │ │ :7004-7 │ │ :7010-7 │ │ :7016-7 │ -└────┬────┘ └────┬────┘ └────┬────┘ └────┬────┘ - │ │ │ │ - Cluster 1 Cluster 2 Cluster 3 Cluster 4 - (独立集群) (独立集群) (独立集群) (独立集群) +```mermaid +flowchart TD + App["你的应用服务"] --> SA["Service A
:7001-7007"] + App --> SB["Service B
:7004-7007"] + App --> SC["Service C
:7010-7017"] + App --> SD["Service D
:7016-7023"] + SA --> C1["Cluster 1
独立集群"] + SB --> C2["Cluster 2
独立集群"] + SC --> C3["Cluster 3
独立集群"] + SD --> C4["Cluster 4
独立集群"] + + style App fill:#e3f2fd,stroke:#1565c0 + style C1,C2,C3,C4 fill:#e8f5e9,stroke:#4caf50 ``` > [!IMPORTANT] 原则 @@ -400,5 +404,5 @@ let val = await cluster.get("key") ## 关联笔记 - [[hhs/Redis/06-主从与哨兵]] — Sentinel 单主高可用方案 -- [[hhs/Redis/09-运维调优]] — 监控指标与告警 +- [[hhs/Redis/11-运维与性能调优]] — 监控指标与告警 - [[hhs/Redis/README]] — 知识索引总览 diff --git a/hhs/Redis/09-高级特性.md b/hhs/Redis/09-高级特性.md index 5f7b89a..ec2f082 100644 --- a/hhs/Redis/09-高级特性.md +++ b/hhs/Redis/09-高级特性.md @@ -97,20 +97,28 @@ return 0 -- 被其他客户端持有 ``` ```go -// Go 调用 +// Go 调用——与上方 Lua 脚本等价的分布式锁实现 script := redis.NewScript(` - if redis.call("get", KEYS[1]) == false then - -- setex(key, seconds, value): ARGV[1]=过期秒数, ARGV[2]=唯一令牌 - redis.call("setex", KEYS[1], ARGV[1], ARGV[2]) + -- 尝试获取锁:SET NX EX(只有 key 不存在时才设置,原子操作) + if redis.call("set", KEYS[1], ARGV[2], "NX", "EX", ARGV[1]) == "OK" then + return 1 -- 成功获取 + end + -- 已持有锁且未过期,续期(防止业务未完成锁就释放) + if redis.call("get", KEYS[1]) == ARGV[2] then + redis.call("expire", KEYS[1], ARGV[1]) return 1 end - return 0 + return 0 -- 被其他客户端持有 `) -ok, err := script.Run(ctx, c, []string{"lock:order:" + orderId}, +// KEYS[1] = lock key, ARGV[1] = TTL 秒数, ARGV[2] = 唯一令牌(如 UUID) +ok, err := script.Run(ctx, rdb, []string{"lock:order:" + orderId}, lockTTL, clientToken).Int() if ok == 1 { - defer unlock(clientToken) + defer func() { + // 解锁:只有持有者才能释放(Lua 保证原子性) + unlockScript.Run(ctx, rdb, []string{"lock:order:" + orderId}, clientToken) + }() doBiz() // 临界区代码 } ``` @@ -355,6 +363,9 @@ flowchart LR > | 任务队列、订单事件 | Stream | > | 海量消息、高吞吐 | Kafka/RabbitMQ | +> [!tip] Stream 深入阅读 +> 本节仅覆盖 Stream 基础用法。Consumer Group 内部机制、消息确认与重试(XCLAIM/XAUTOCLAIM)、Dead Letter Queue 模式、容量管理等进阶内容,详见 [[hhs/Redis/15-Stream]]。 + ## 五、Keyspace Notifications —— 键事件监听 ### 原理与配置 diff --git a/hhs/Redis/12-Bitmap.md b/hhs/Redis/12-Bitmap.md new file mode 100644 index 0000000..626190a --- /dev/null +++ b/hhs/Redis/12-Bitmap.md @@ -0,0 +1,186 @@ +--- +tags: [Redis, 缓存, 数据结构, Bitmap] +create time: 2026-05-24 10:00 +--- + +# Redis Bitmap 位图 + +## 概述 + +Bitmap 并非 Redis 的独立数据结构,而是基于 String 类型的一套位操作指令。它将一个 String 值视为一个巨大的位数组,每个 bit 只占 1 位,因此在处理大规模布尔型数据(如签到、在线状态、DAU 统计)时,内存消耗极低。一个包含 10 亿用户的 Bitmap 仅需约 125MB,这使其成为高并发场景下的利器。 + +## 正文 + +### 1. 核心原理 + +> [!question] 为什么 Bitmap 不是一种新数据结构? +> 因为 Bitmap 本质上就是 String。Redis 的 String 底层是 SDS(Simple Dynamic String),存储的是字节数组。Bitmap 操作只是在这个字节数组上进行位级别的读写,不涉及新的底层结构。 + +Bitmap 的核心指令只有四个: + +| 指令 | 作用 | 时间复杂度 | +|------|------|-----------| +| `SETBIT key offset value` | 设置指定位的值(0 或 1) | O(1) | +| `GETBIT key offset` | 获取指定位的值 | O(1) | +| `BITCOUNT key [start end]` | 统计值为 1 的位数 | O(N),N 为字节数 | +| `BITOP operation destkey key [key ...]` | 对多个 Bitmap 做位运算 | O(N) | + +其中 `BITOP` 支持 `AND`、`OR`、`NOT`、`XOR` 四种位运算,用于多个 Bitmap 之间的聚合分析。 + +> [!tip] offset 是从 0 开始的 +> `SETBIT sign:1001 5 1` 表示将 key `sign:1001` 的第 6 个 bit(offset=5)置为 1。如果 key 不存在,Redis 会自动扩展字符串长度。 + +### 2. 用户签到系统 + +用 Bitmap 实现签到非常直观:每个用户一个 key,每天对应一个 bit 位,签到则置 1。 + +```go +// 用户签到(offset = 一年中的第几天) +func SignIn(ctx context.Context, rdb *redis.Client, userID int64, dayOfYear int) error { + key := fmt.Sprintf("sign:%d:%d", userID, time.Now().Year()) + return rdb.SetBit(ctx, key, int64(dayOfYear), 1).Err() +} + +// 查询某天是否签到 +func IsSigned(ctx context.Context, rdb *redis.Client, userID int64, dayOfYear int) (bool, error) { + key := fmt.Sprintf("sign:%d:%d", userID, time.Now().Year()) + val, err := rdb.GetBit(ctx, key, int64(dayOfYear)).Result() + return val == 1, err +} + +// 本月累计签到天数 +func MonthSignCount(ctx context.Context, rdb *redis.Client, userID int64) (int64, error) { + key := fmt.Sprintf("sign:%d:%d", userID, time.Now().Year()) + // BITCOUNT 支持按字节范围统计,start/end 是字节索引 + // 需要根据月份计算对应的字节范围 + return rdb.BitCount(ctx, key, &redis.BitCount{ + Start: int64((time.Now().Month() - 1) * 4), // 近似值,实际需精确计算 + End: int64(time.Now().Month()*4 - 1), + }).Result() +} +``` + +> [!question] 连续签到 7 天如何判断? +> 有两种思路:(1) 用 `BITCOUNT` 统计最近 7 个 bit 是否全为 1;(2) 更高效的做法是用位移掩码——取出 7 位的值,判断是否等于 `0b1111111`(即 127)。 + +下面是连续签到检测的示例: + +```go +// 检测最近 N 天是否连续签到 +func IsContinuousSign(ctx context.Context, rdb *redis.Client, userID int64, days int) (bool, error) { + key := fmt.Sprintf("sign:%d:%d", userID, time.Now().Year()) + today := time.Now().YearDay() + + // 逐 bit 检查最近 N 天 + for i := 0; i < days; i++ { + val, err := rdb.GetBit(ctx, key, int64(today-i)).Result() + if err != nil { + return false, err + } + if val == 0 { + return false, nil + } + } + return true, nil +} +``` + +### 3. DAU 统计 + +DAU(Daily Active Users)是 Bitmap 最经典的应用场景。思路很简单:每天一个 key,每个用户 ID 对应一个 bit 位。 + +```go +// 记录用户今日活跃 +func MarkActive(ctx context.Context, rdb *redis.Client, userID int64) error { + key := fmt.Sprintf("dau:%s", time.Now().Format("2006-01-02")) + return rdb.SetBit(ctx, key, userID, 1).Err() +} + +// 查询今日 DAU +func GetDAU(ctx context.Context, rdb *redis.Client, date string) (int64, error) { + key := fmt.Sprintf("dau:%s", date) + return rdb.BitCount(ctx, key, nil).Result() +} + +// 查询本周 UV(去重) +func GetWeeklyUV(ctx context.Context, rdb *redis.Client) (int64, error) { + destKey := "dau:weekly:tmp" + var keys []string + // 收集最近 7 天的 key + for i := 0; i < 7; i++ { + day := time.Now().AddDate(0, 0, -i).Format("2006-01-02") + keys = append(keys, fmt.Sprintf("dau:%s", day)) + } + // OR 运算:任意一天活跃即算周活 + err := rdb.BitOpOr(ctx, destKey, keys...).Err() + if err != nil { + return 0, err + } + return rdb.BitCount(ctx, destKey, nil).Result() +} +``` + +> [!tip] BITOP OR 的去重原理 +> 假设用户 A 在周一和周三都活跃,对应的 bit 位已经是 1。OR 运算后该位仍是 1,最终 BITCOUNT 只统计一次,天然实现了去重。这比在应用层用 Set 去重高效得多。 + +### 4. 内存计算 + +> [!tip] Bitmap 到底有多省内存?我们来算一笔账。 + +假设平台有 **10 亿注册用户**,用户 ID 从 0 到 999,999,999。 + +``` +总 bit 数 = 1,000,000,000 bits +总字节数 = 1,000,000,000 / 8 = 125,000,000 bytes ≈ 119.2 MB +``` + +对比一下其他方案存储同样 10 亿用户的信息: + +| 方案 | 数据结构 | 内存消耗 | +|------|---------|---------| +| Bitmap | 1 bit / 用户 | ~125 MB | +| Hash(存 boolean) | ~50 bytes / 用户(含 key 开销) | ~47 GB | +| Set(存用户 ID) | ~16 bytes / 用户 | ~15 GB | + +Bitmap 的内存效率比 Hash 方案低约 **380 倍**,这就是为什么在大规模布尔型统计场景下,Bitmap 是首选。 + +> [!question] 为什么 Bitmap 这么省? +> 因为它只用 1 个 bit 来表示一个布尔值,而 Hash/Set 需要存储完整的 key 和 value。Redis String 底层 SDS 本身也有元数据开销,但分摊到数十亿个 bit 上几乎可以忽略。 + +### 5. Bitmap vs Set vs HyperLogLog 对比 + +| 特性 | Bitmap | Set | HyperLogLog | +|------|--------|-----|-------------| +| **底层结构** | String(位数组) | Hashtable / Ziplist | 概率算法 | +| **单条数据内存** | 1 bit | 16~50 bytes | 固定 12 KB | +| **精确度** | 精确 | 精确 | 误差 ~0.81% | +| **支持操作** | 位运算、计数 | 集合运算、随机取 | 仅计数 | +| **适用场景** | 签到、DAU、布隆过滤器 | 精确集合运算 | 超大规模基数估算 | +| **数据量上限** | 受 String 最大 512 MB 限制(约 40 亿 bit) | 受内存限制 | 固定 12 KB | + +> [!question] 该选哪个? +> - 需要精确统计 + 数据是布尔型(在线/签到/活跃) -> **Bitmap** +> - 需要精确集合运算(交集、差集) -> **Set** +> - 只需要统计基数(UV/PV),允许 0.81% 误差 -> **HyperLogLog** + +### 6. 签到系统流程 + +```mermaid +flowchart TD + A["用户发起签到请求"] --> B["计算当日 offset"] + B --> C["SETBIT sign:uid:year offset 1"] + C --> D["签到成功"] + D --> E{"查询本月签到?"} + E -->|是| F["BITCOUNT 计算本月签到天数"] + E -->|否| G{"检查连续签到?"} + G -->|是| H["逐 bit 检查最近 N 天"] + G -->|否| I["返回结果"] + F --> I + H --> I +``` + +## 关联笔记 + +- [[hhs/Redis/02-核心数据类型]] +- [[hhs/Redis/13-HyperLogLog]] +- [[hhs/Redis/README]] diff --git a/hhs/Redis/13-HyperLogLog.md b/hhs/Redis/13-HyperLogLog.md new file mode 100644 index 0000000..666a172 --- /dev/null +++ b/hhs/Redis/13-HyperLogLog.md @@ -0,0 +1,212 @@ +--- +tags: [Redis, 缓存, 数据结构, HyperLogLog] +create time: 2026-05-24 10:00 +--- + +# Redis HyperLogLog 概率计数 + +## 概述 + +HyperLogLog(HLL)是 Redis 提供的**概率型基数估算数据结构**——用固定的 **12KB 内存**,就能估算高达 2^64 个不同元素的数量,标准误差仅 **0.81%**。它不存储元素本身,只回答"有多少个不同的值",是 UV 统计、去重计数等场景的利器。 + +> [!question] 为什么不直接用 Set 做去重计数? +> 1 亿个用户 ID 存进 Set,每个 ID 至少占几十字节,轻松吃掉几 GB 内存。HyperLogLog 只需要 12KB,代价是允许 0.81% 的误差——对于"今天有多少独立访客"这类问题,这个精度完全够用。 + +## 一、核心原理 + +### 概率思想:用"前导零"推算基数 + +HyperLogLog 的直觉来源于**伯努利试验**:抛一枚均匀硬币,连续出现 k 次正面的概率是 1/2^k。如果我们在所有试验中观察到最长的连续正面次数为 k,就可以反推大约有 2^k 个不同的试验。 + +实际做法: + +1. **哈希映射**:对每个元素做哈希,得到一个均匀分布的比特串 +2. **分桶**:用前 14 位定位到 2^14 = 16384 个桶中的一个 +3. **记录前导零**:对剩余比特,统计第一个 1 出现前连续 0 的个数(取最大值) +4. **调和均值**:所有桶的值通过调和均值公式合并,得到基数估算 + +> [!tip] 关键洞察 +> HLL 不存储元素,只存储每个桶的"最大前导零位数"。每个桶只需 6 个比特(能表示 0~63),所以总内存 = 16384 × 6 bit ≈ **12KB**,无论存 1 个元素还是 10 亿个元素。 + +### Redis 中只有 3 个命令 + +| 命令 | 作用 | 时间复杂度 | +|------|------|-----------| +| `PFADD key element [element ...]` | 向 HLL 添加元素 | O(N),N 为元素个数 | +| `PFCOUNT key [key ...]` | 返回基数估算值 | O(N),N 为 HLL 个数 | +| `PFMERGE dest source [source ...]` | 合并多个 HLL | O(N),N 为 HLL 个数 | + +> [!question] 为什么命令叫 PF 开头? +> PF 是 Philippe Flajolet 的名字缩写——他是 HyperLogLog 算法论文的主要作者。 + +## 二、UV 统计实战 + +### 用 Go 实现页面 UV 计数 + +以下是一个完整的 Web UV 统计示例,使用 go-redis v9: + +```go +package main + +import ( + "context" + "fmt" + "log" + "time" + + "github.com/redis/go-redis/v9" +) + +var rdb *redis.Client +var ctx = context.Background() + +func init() { + rdb = redis.NewClient(&redis.Options{ + Addr: "localhost:6379", + DB: 0, + }) +} + +// RecordVisit 记录一次页面访问,visitorID 为用户唯一标识 +func RecordVisit(page string, visitorID string) error { + key := fmt.Sprintf("uv:%s:%s", page, time.Now().Format("2006-01-02")) + // PFADD:如果 visitorID 已存在,不会重复计数(概率去重) + return rdb.PFAdd(ctx, key, visitorID).Err() +} + +// GetUV 获取某页面当日独立访客数 +func GetUV(page string) (int64, error) { + key := fmt.Sprintf("uv:%s:%s", page, time.Now().Format("2006-01-02")) + // PFCOUNT 返回估算的基数,误差约 0.81% + return rdb.PFCount(ctx, key).Result() +} + +func main() { + // 模拟 10000 个不同用户访问首页 + for i := 0; i < 10000; i++ { + visitorID := fmt.Sprintf("user_%d", i) + if err := RecordVisit("home", visitorID); err != nil { + log.Fatal(err) + } + } + + uv, _ := GetUV("home") + fmt.Printf("首页今日 UV: %d\n", uv) // 输出约 10000,误差极小 +} +``` + +> [!tip] 对比 Set 方案 +> 如果用 `SADD` + `SCARD` 实现同样功能,10000 个用户 ID 至少消耗 **数百 KB**;而 HyperLogLog 的 `uv:home:2026-05-24` 这个 key 始终只占 **12KB**。用户量越大,HLL 的优势越明显。 + +## 三、合并多个统计源 + +实际业务中经常需要合并数据——比如从每日 UV 汇总出每周 UV。 + +```go +// MergeWeeklyUV 将 7 天的日 UV 合并为周 UV +func MergeWeeklyUV(page string, weekStart time.Time) error { + destKey := fmt.Sprintf("uv:%s:week:%s", page, weekStart.Format("2006-01-02")) + + // 构造 7 天的源 key + var sourceKeys []string + for i := 0; i < 7; i++ { + day := weekStart.AddDate(0, 0, i) + sourceKeys = append(sourceKeys, fmt.Sprintf("uv:%s:%s", page, day.Format("2006-01-02"))) + } + + // PFMERGE:将多个 HLL 合并到 dest,结果仍是 HLL(12KB) + // 合并后的基数 ≈ 所有源集合的并集大小 + return rdb.PFMerge(ctx, destKey, sourceKeys...).Err() +} + +// GetWeeklyUV 获取周 UV +func GetWeeklyUV(page string, weekStart time.Time) (int64, error) { + key := fmt.Sprintf("uv:%s:week:%s", page, weekStart.Format("2006-01-02")) + return rdb.PFCount(ctx, key).Result() +} +``` + +> [!question] PFMERGE 是精确合并吗? +> 是的。PFMERGE 会取每个桶在所有源 HLL 中的最大值,数学上等价于对并集重新计算 HLL。合并后再 PFCOUNT,得到的就是**并集的估算基数**,误差不会因为合并而放大。 + +## 四、误差分析 + +HyperLogLog 的标准误差为 **1.04 / sqrt(m)**,其中 m 是桶数。Redis 使用 16384 个桶: + +``` +标准误差 = 1.04 / sqrt(16384) ≈ 0.81% +``` + +| 实际基数 | 估算误差范围(±0.81%) | +|---------|----------------------| +| 1,000 | ±8 | +| 100,000 | ±810 | +| 10,000,000 | ±81,000 | + +> [!question] 0.81% 的误差在实际中能接受吗? +> 绝大多数场景完全可以。UV 统计本身就有缓存穿透、机器人流量等噪声,0.81% 的误差远小于这些业务噪声。但如果你需要精确计数(比如库存扣减),请用 Set 或数据库。 + +## 五、内存计算 + +HyperLogLog 的内存占用是**固定**的,与元素数量无关: + +``` +桶数 = 2^14 = 16384 +每个桶 = 6 bit +总内存 = 16384 × 6 bit = 98304 bit ≈ 12KB +``` + +> [!tip] Redis 的内存优化 +> 当 HLL 计数的基数较小时(稀疏表示),Redis 会用更紧凑的编码存储,可能只占几百字节。只有当基数增长到一定程度,才会扩展到完整的 12KB。这个过程是自动的,对用户透明。 + +| 场景 | Set 内存 | HyperLogLog 内存 | +|------|---------|-----------------| +| 1 万个 ID | ~500KB | 12KB | +| 100 万个 ID | ~50MB | 12KB | +| 1 亿个 ID | ~5GB | 12KB | + +## 六、适用场景 vs 不适用场景 + +### 适合使用 HyperLogLog + +- 页面/接口 UV 统计 +- 搜索关键词去重数 +- 广告曝光独立用户数 +- 任何"只需要知道有多少个不同值"的场景 + +### 不适合使用 HyperLogLog + +- **需要知道具体有哪些元素**:HLL 不存储元素本身,无法枚举 +- **需要精确计数**:0.81% 误差不可接受的场景(如库存、余额) +- **数据量极小(< 1000)**:此时 Set 占用内存也很小,且精确无误 + +> [!question] 三种方案怎么选? + +| 维度 | HyperLogLog | Set | Bitmap | +|------|------------|-----|--------| +| 功能 | 基数估算 | 精确去重 + 列举 | 位标记 + 计数 | +| 内存 | 固定 12KB | 随元素线性增长 | 约 max_id / 8 字节 | +| 精度 | ~99.19% | 100% | 100% | +| 典型场景 | UV 统计 | 标签、好友列表 | 签到、在线状态 | + +> [!tip] 一句话决策 +> 只需要"有多少个"→ HyperLogLog;需要"有哪些"→ Set;需要"某个 ID 在不在"→ Bitmap。 + +## 七、UV 统计流程 + +```mermaid +graph LR + A["用户请求页面"] --> B["PFADD uv:page:date visitor_id"] + B --> C["Redis HyperLogLog"] + C --> D["PFCOUNT uv:page:date"] + D --> E["返回估算 UV 数"] + E --> F["展示在监控面板"] + G["每日凌晨"] --> H["PFMERGE uv:page:week daily_keys"] + H --> I["周报 UV 汇总"] +``` + +## 关联笔记 + +- [[hhs/Redis/02-核心数据类型]] +- [[hhs/Redis/12-Bitmap]] +- [[hhs/Redis/README]] diff --git a/hhs/Redis/14-BloomFilter.md b/hhs/Redis/14-BloomFilter.md new file mode 100644 index 0000000..b233c56 --- /dev/null +++ b/hhs/Redis/14-BloomFilter.md @@ -0,0 +1,240 @@ +--- +tags: [Redis, 缓存, 数据结构, BloomFilter] +create time: 2026-05-24 10:00 +--- + +# Redis Bloom Filter 布隆过滤器 + +## 概述 + +布隆过滤器(Bloom Filter)是一种空间效率极高的概率型数据结构,用于判断一个元素**是否存在于集合中**。核心特性: + +> **"Definitely not or probably yes"** —— 说不存在则一定不存在;说存在则**大概率**存在(有极小误判概率)。 + +这个特性让它成为解决**缓存穿透**问题的利器。 + +> [!question] 为什么需要布隆过滤器? +> 攻击者用大量不存在的 key 轰击系统,每个请求穿透缓存直打 DB,最终压垮数据库。布隆过滤器能在请求到达 DB 之前,用极低的内存成本拦截掉"肯定不存在"的请求。 + +## 一、核心原理 + +布隆过滤器的底层:**一个 bit 数组 + 多个独立的哈希函数**。 + +- **bit 数组**:长度为 `m`,初始全部为 0 +- **哈希函数**:`k` 个相互独立的哈希函数,每个函数将输入映射到 `[0, m-1]` 的某个位置 + +**添加元素**:对元素分别用 `k` 个哈希函数计算,将 bit 数组中对应位置全部置为 1。**查询元素**:用同样的 `k` 个哈希函数计算,所有对应位置都是 1 则判定"可能存在";任意一位是 0 则判定"一定不存在"。 + +> [!tip] 为什么会有误判? +> 不同元素经过哈希计算后可能映射到相同位置(哈希冲突)。随着插入元素增多,查询一个不存在的元素时,恰好所有位置都被其他元素"碰巧"置为 1 的概率就会上升。 + +```mermaid +graph LR + A["element x"] --> H1["Hash1(x) = 3"] + A --> H2["Hash2(x) = 7"] + A --> H3["Hash3(x) = 11"] + H1 --> B1["bit[3] = 1"] + H2 --> B2["bit[7] = 1"] + H3 --> B3["bit[11] = 1"] +``` + +## 二、Redis 中的两种方案 + +### 方案 A:RedisBloom 模块(服务端) + +RedisBloom 是 Redis 官方模块,通过 `BF.ADD` / `BF.EXISTS` 命令直接操作。 + +```go +// go-redis 调用 RedisBloom +func checkWithRedisBloom(ctx context.Context, rdb *redis.Client, key, value string) (bool, error) { + result, err := rdb.Do(ctx, "BF.EXISTS", key, value).Int() // 1=可能存在, 0=一定不存在 + if err != nil { + return false, err + } + return result == 1, nil +} +``` + +### 方案 B:Go 侧布隆过滤器库(客户端) + +使用 `github.com/bits-and-blooms/bloom/v3` 在应用内存中维护过滤器。 + +```go +import "github.com/bits-and-blooms/bloom/v3" + +filter := bloom.NewWithEstimates(1_000_000, 0.0001) // 100万元素, 误判率0.01% +filter.AddString("user:10086") // 添加 +exists := filter.TestString("user:10086") // 查询: true = 可能存在 +``` + +## 三、缓存穿透防护实战 + +这是布隆过滤器最经典的应用场景。下面用一个完整的 Go 示例展示三层防护架构。 + +### 3.1 架构流程 + +```mermaid +graph TD + Req["客户端请求"] --> BF{"布隆过滤器判断"} + BF -- "不存在" --> Ret1["直接返回空, 拒绝请求"] + BF -- "可能存在" --> Cache{"查询 Redis 缓存"} + Cache -- "命中" --> Ret2["返回缓存数据"] + Cache -- "未命中" --> DB["查询 MySQL"] + DB -- "找到数据" --> SetCache["写入 Redis 缓存"] + SetCache --> Ret3["返回数据"] + DB -- "未找到" --> Ret4["返回空并缓存空值"] +``` + +### 3.2 完整代码示例 + +```go +package main + +import ( + "context" + "encoding/json" + "time" + + "github.com/bits-and-blooms/bloom/v3" + "github.com/redis/go-redis/v9" +) + +type User struct { + ID int64 `json:"id"` + Name string `json:"name"` +} + +type BloomCacheService struct { // 三层防护服务 + rdb *redis.Client + filter *bloom.BloomFilter +} + +// NewBloomCacheService 初始化:加载全量 ID 到布隆过滤器 +func NewBloomCacheService(rdb *redis.Client) *BloomCacheService { + filter := bloom.NewWithEstimates(1_000_000, 0.0001) // 100万用户, 误判率0.01% + + // 从 Redis Set 批量加载已有 ID(生产环境可从 DB 加载) + ctx := context.Background() + ids, _ := rdb.SMembers(ctx, "user:all_ids").Result() + for _, id := range ids { + filter.AddString(id) + } + return &BloomCacheService{rdb: rdb, filter: filter} +} + +// GetUser 三层查询:Bloom -> Redis -> MySQL +func (s *BloomCacheService) GetUser(ctx context.Context, userID string) (*User, error) { + cacheKey := "user:" + userID + + // 第一层:布隆过滤器拦截 + if !s.filter.TestString(userID) { + return nil, nil // 一定不存在 + } + + // 第二层:Redis 缓存 + data, err := s.rdb.Get(ctx, cacheKey).Bytes() + if err == nil { + if string(data) == "" { + return nil, nil // 缓存的空值 + } + var user User + json.Unmarshal(data, &user) + return &user, nil + } + + // 第三层:查询数据库 + user, err := queryMySQL(userID) + if err != nil { + return nil, err + } + + if user == nil { + s.rdb.Set(ctx, cacheKey, "", 5*time.Minute) // 缓存空值, 短 TTL + return nil, nil + } + + bytes, _ := json.Marshal(user) + s.rdb.Set(ctx, cacheKey, bytes, 30*time.Minute) + return user, nil +} + +// queryMySQL 模拟数据库查询 +func queryMySQL(userID string) (*User, error) { + return nil, nil // 实际项目中查询 MySQL +} +``` + +> [!question] 为什么还需要缓存空值? +> 布隆过滤器拦截了"肯定不存在"的请求,但如果某个 ID 曾经存在过(布隆过滤器会说"可能存在"),但实际已从 DB 删除,此时仍会穿透到 DB。缓存空值是第二道防线。 + +## 四、参数计算 + +选择布隆过滤器的两个关键参数: + +| 参数 | 含义 | 影响 | +|------|------|------| +| `m` | bit 数组长度 | 越大误判率越低,但内存占用越多 | +| `k` | 哈希函数个数 | 太少冲突多,太多计算慢 | + +**公式**(基于期望元素数 `n` 和目标误判率 `p`): + +$$m = -\frac{n \ln p}{(\ln 2)^2}, \quad k = \frac{m}{n} \ln 2$$ + +> [!tip] 快速估算 +> 对于 n = 100 万、p = 0.01% 的常见场景: +> - m ≈ 19,170,117 bits ≈ **2.3 MB** +> - k ≈ 13 个哈希函数 +> +> 相比之下,用 Redis Set 存 100 万个字符串 key 至少需要 **几十 MB**。布隆过滤器的内存优势是数量级的。 + +```go +import "math" + +func calcBloomParams(n int, p float64) (m int, k int) { + // m = -n * ln(p) / (ln2)^2 + ln2 := math.Ln2 + m = int(-float64(n) * math.Log(p) / (ln2 * ln2)) + // k = (m/n) * ln2 + k = int(math.Round(float64(m) / float64(n) * ln2)) + return m, k +} +``` + +## 五、RedisBloom vs Go-side 对比 + +| 维度 | RedisBloom 模块 | Go-side 布隆库 | +|------|----------------|---------------| +| **部署** | 需安装 RedisBloom 模块 | 只需引入 Go 依赖 | +| **网络开销** | 每次查询一次网络往返 | 本地内存, 零网络开销 | +| **分布式一致性** | 天然共享 | 需自行同步 | +| **性能** | 受网络 RT 限制, < 10 万 QPS | 纯内存, 百万 QPS | +| **持久化** | 随 Redis RDB/AOF 自动持久化 | 需额外持久化机制 | +| **适用场景** | 多实例共享, 数据量大 | 单服务高频查询 | + +> [!tip] 如何选择? +> - **微服务多实例**:优先 RedisBloom,天然共享无需同步 +> - **单体服务、极致性能**:Go-side 布隆库,避免网络开销 +> - **折中方案**:Go-side 为主,定期从 Redis 同步 bit 数组快照 + +## 六、局限性 + +### 6.1 不支持删除 + +标准布隆过滤器不支持删除——将某位从 1 置为 0 可能影响其他元素判断。解决方案是**计数布隆过滤器(Counting Bloom Filter)**,用计数器替代 1 bit,删除时计数器减 1。RedisBloom 的 Cuckoo Filter(`CF.ADD` / `CF.EXISTS`)也原生支持删除。 + +### 6.2 需要预加载 + +布隆过滤器必须在使用前加载全部合法元素。新增数据时需同步更新过滤器,初始化时需遍历全量数据。 + +> [!question] 如果新增了数据但忘记更新布隆过滤器会怎样? +> 新增的 key 在布隆过滤器中查询会返回"不存在",导致合法请求被误拦截。这比缓存穿透更严重——用户明明有数据却查不到。所以务必确保新增数据时同步 `BF.ADD` 或 `filter.AddString()`。 + +### 6.3 误判率随使用上升 + +实际元素数量远超预估值 `n` 时,误判率会显著上升。建议预留 20%-30% 余量,监控实际元素数量,必要时重建过滤器。 + +## 关联笔记 + +- [[hhs/Redis/10-缓存架构模式]] +- [[hhs/Redis/02-核心数据类型]] +- [[hhs/Redis/README]] diff --git a/hhs/Redis/15-Stream.md b/hhs/Redis/15-Stream.md new file mode 100644 index 0000000..8a431f3 --- /dev/null +++ b/hhs/Redis/15-Stream.md @@ -0,0 +1,455 @@ +--- +tags: [Redis, 缓存, Stream, 消息队列] +create time: 2026-05-24 10:00 +--- + +# Redis Stream 消息队列 + +## 概述 + +Redis Stream 是 Redis 5.0 引入的日志型数据结构,专门为**可靠消息传递**设计。它弥补了 Pub/Sub 「发完即忘、不持久化」的先天缺陷,同时又比 Kafka/RabbitMQ 这类外部中间件轻量——不需要额外部署组件,复用已有的 Redis 实例即可。 + +> [!question] Pub/Sub 够用吗?为什么要用 Stream? +> Pub/Sub 的核心问题:消息**不持久化**,订阅者下线就丢失。Stream 将每条消息追加写入内存(可配合 AOF 持久化),并通过 Consumer Group 实现**消费确认**机制——消息只有被 ACK 后才算真正消费完成。这使得 Stream 能胜任任务队列、事件溯源等对可靠性有要求的场景。 + +## 一、核心概念 + +### 基本操作命令 + +```bash +# === 生产者 === +XADD mystream * field1 value1 field2 value2 +# * 表示自动生成 ID(时间戳-序号),如 1685000000000-0 + +# === 消费者(无消费者组模式)=== +XREAD COUNT 10 BLOCK 5000 STREAMS mystream 0 +# BLOCK 5000 = 最多等 5 秒;0 = 从头读取 +# 返回一批 XEntry,每条包含 ID + fields + +# === 消费者组模式(推荐)=== +XGROUP CREATE mystream mygroup $ # $ = 只消费创建后的新消息 +XREADGROUP GROUP mygroup consumer1 COUNT 1 BLOCK 5000 STREAMS mystream > +# > 表示「只取未分配的新消息」 +# 返回后消息进入该消费者的 pending 列表 + +XACK mystream mygroup # 确认消费完成 +``` + +> [!tip] `>` vs `$` vs `0` —— 三个特殊 ID +> | ID | 含义 | 典型场景 | +> |---|------|---------| +> | `0` | Stream 的第一条消息 | 数据重放、全量回溯 | +> | `$` | 当前最新消息的 ID | `XGROUP CREATE` 时指定起始点 | +> | `>` | 只返回**未被任何消费者分配**的新消息 | `XREADGROUP` 正常消费循环 | + +### Stream 底层结构 + +每条消息以 Radix Tree + Listpack 存储,时间复杂度 O(1) 追加。ID 格式为 `<毫秒时间戳>-<序号>`,天然有序。 + +> [!question] Stream 会像 List 一样消费后删除吗? +> 不会。Stream 是**只追加日志**,消息消费后仍在。需要显式 `XDEL` 或通过 `MAXLEN`/`MINID` 策略裁剪。这带来了「可回溯」的优势,但也意味着必须主动管理容量。 + +## 二、Consumer Group 机制详解 + +Consumer Group 是 Stream 的核心抽象,它让多个消费者**协作消费同一条 Stream**,每条消息只会被组内的一个消费者处理。 + +### 消息生命周期 + +```mermaid +flowchart TD + P["Producer"] -->|"XADD"| S["Stream"] + S -->|"XREADGROUP"| C1["Consumer A"] + S -->|"XREADGROUP"| C2["Consumer B"] + C1 -->|"XACK"| S + C2 -->|"XACK"| S + C1 -.->|"crashed"| PDL["Pending List"] + PDL -->|"XCLAIM"| C2 + C2 -->|"XACK"| S + + style S fill:#dfd,stroke:#090 + style C1 fill:#ddf,stroke:#669 + style C2 fill:#ddf,stroke:#669 + style PDL fill:#fdd,stroke:#f00 +``` + +### Pending List(待确认列表) + +当消费者通过 `XREADGROUP GROUP ... >` 读取消息后,消息进入该消费者的 **PEL(Pending Entry List)**。只有显式调用 `XACK` 后才会移除。 + +```bash +# 查看组内 pending 状态 +XPENDING mystream mygroup + +# 更详细的 pending 信息(包含消费者、空闲时间) +XPENDING mystream mygroup - + 10 # 返回最多 10 条 pending 条目 +``` + +XPENDING 返回字段: + +| 字段 | 含义 | +|------|------| +| `message-id` | 消息 ID | +| `consumer` | 被分配的消费者 | +| `idle-ms` | 自读取以来的空闲时间(毫秒) | +| `delivery-count` | 被投递的次数(重试计数器) | + +> [!tip] Pending 是 Stream 的「可靠性根基」 +> Pub/Sub 没有 pending 概念,消息发出去就消失了。Stream 的 PEL 记录了「谁领了哪条消息、领了多久、重试了几次」,这就是「至少一次投递」语义的实现基础。 + +### 消费者组内部机制 + +```mermaid +sequenceDiagram + participant P as Producer + participant S as Stream + participant G as Consumer Group + participant C1 as Consumer A + participant C2 as Consumer B + + P->>S: XADD msg-1 + P->>S: XADD msg-2 + P->>S: XADD msg-3 + C1->>G: XREADGROUP COUNT 2 + G->>S: 分配 msg-1, msg-2 + S-->>C1: msg-1, msg-2 (pending) + C2->>G: XREADGROUP COUNT 2 + G->>S: 分配 msg-3 + S-->>C2: msg-3 (pending) + C1->>S: XACK msg-1 + C1->>S: XACK msg-2 + C2->>S: XACK msg-3 +``` + +> [!question] 消息是如何分配给消费者的? +> Redis 采用**轮询分配**(round-robin):新消息到达时,依次分配给组内活跃的消费者。注意不是负载均衡——如果某个消费者处理慢,它积压的 pending 消息不会自动转移给其他消费者,需要通过 XCLAIM 手动认领。 + +## 三、Go 实战 + +### 生产者 + +```go +// StreamProducer 封装了 Stream 的写入逻辑 +type StreamProducer struct { + rdb *redis.Client + stream string + maxSize int64 // MAXLEN 裁剪阈值 +} + +func NewStreamProducer(rdb *redis.Client, stream string) *StreamProducer { + return &StreamProducer{rdb: rdb, stream: stream, maxSize: 100000} +} + +func (p *StreamProducer) Publish(ctx context.Context, fields map[string]interface{}) (string, error) { + // XADD 的 ApproximateMaxLen 约等于 MAXLEN ~,性能更好 + id, err := p.rdb.XAdd(ctx, &redis.XAddArgs{ + Stream: p.stream, + ApproximateMaxLen: true, // ~ 模式,不精确计数但更快 + MaxLen: p.maxSize, + Values: fields, + }).Result() + if err != nil { + return "", fmt.Errorf("XADD failed: %w", err) + } + return id, nil +} +``` + +### 消费者(含重连与重试) + +```go +type StreamConsumer struct { + rdb *redis.Client + stream string + group string + consumer string + handler func(ctx context.Context, msg redis.XMessage) error + batchSize int64 + blockTime time.Duration +} + +func (c *StreamConsumer) Start(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + default: + // XREADGROUP 阻塞读取 + streams, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{ + Group: c.group, + Consumer: c.consumer, + Streams: []string{c.stream, ">"}, // > = 只取新消息 + Count: c.batchSize, + Block: c.blockTime, + }).Result() + + if err != nil { + if errors.Is(err, redis.Nil) { + continue // 超时无消息,正常 + } + log.Printf("XReadGroup error: %v, retrying...", err) + time.Sleep(time.Second) // 退避重试 + continue + } + + for _, stream := range streams { + for _, msg := range stream.Messages { + if err := c.handler(ctx, msg); err != nil { + log.Printf("handler error for %s: %v", msg.ID, err) + continue // 不 ACK,消息留在 pending 列表 + } + // 处理成功,确认消费 + if err := c.rdb.XAck(ctx, c.stream, c.group, msg.ID).Err(); err != nil { + log.Printf("XAck error for %s: %v", msg.ID, err) + } + } + } + } + } +} +``` + +> [!tip] 不 ACK = 进入 pending +> `handler` 返回 error 时跳过 `XACK`,这条消息会留在消费者的 PEL 中。后续通过 XPENDING + XCLAIM 捞回重试——这就是 Stream 实现「至少一次投递」的关键。 + +### 初始化消费者组 + +```go +func EnsureGroup(ctx context.Context, rdb *redis.Client, stream, group string) error { + // MKSTREAM: 如果 stream 不存在则自动创建 + err := rdb.XGroupCreateMkStream(ctx, stream, group, "$").Err() + if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") { + return err + } + // BUSYGROUP = 组已存在,忽略即可 + return nil +} +``` + +## 四、消息确认与重试 + +### XPENDING —— 检查未确认消息 + +```bash +# 概览:返回 [最小ID, 最大ID, 总条数, 各消费者pending数] +XPENDING mystream mygroup + +# 详细列表:按 ID 范围查询 +XPENDING mystream mygroup - + 10 +# - = 最小ID,+ = 最大ID,10 = 最多返回10条 +``` + +### XCLAIM —— 转移消息所有权 + +当消费者宕机后,需要把它 pending 的消息转给存活的消费者处理: + +```bash +# 将空闲超过 30000ms 的消息转给 consumer-b +XCLAIM mystream mygroup consumer-b 30000 +``` + +```go +// Go 实现:定期扫描超时 pending,认领到自己 +func (c *StreamConsumer) ClaimAbandoned(ctx context.Context, minIdle time.Duration) { + pendings, err := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{ + Stream: c.stream, + Group: c.group, + Idle: minIdle, + Start: "-", + End: "+", + Count: 100, + Consumer: "", // 不限定消费者,扫描所有 + }).Result() + if err != nil || len(pendings) == 0 { + return + } + + ids := make([]string, len(pendings)) + for i, p := range pendings { + ids[i] = p.ID + } + + // XCLAIM 认领 + claimed, err := c.rdb.XClaim(ctx, &redis.XClaimArgs{ + Stream: c.stream, + Group: c.group, + Consumer: c.consumer, + MinIdle: minIdle, + Messages: ids, + }).Result() + if err != nil { + return + } + + // 对认领到的消息重新执行 handler + for _, msg := range claimed { + if err := c.handler(ctx, msg); err == nil { + c.rdb.XAck(ctx, c.stream, c.group, msg.ID) + } + } +} +``` + +### XAUTOCLAIM —— Redis 6.2+ 自动认领 + +`XAUTOCLAIM` 把 XPENDING + XCLAIM 合并为一条原子命令: + +```bash +XAUTOCLAIM mystream mygroup consumer-b 30000 0-0 COUNT 10 +# 30000 = 最小空闲时间(ms) +# 0-0 = 起始游标(首次从头开始) +# COUNT = 每次最多返回 10 条 + +# 返回:next-start-id(下次传入作为游标)+ 消息列表 + 已删除的 ID +``` + +> [!question] XCLAIM 会不会被两个消费者同时认领同一条消息? +> 不会。XCLAIM 是单线程原子操作,Redis 保证同一消息在任意时刻只属于一个消费者。如果 consumer-A 已经 ACK 了某条消息,后续 XCLAIM 不会再返回它。 + +## 五、容量管理 + +Stream 是只追加的,消息不会自动删除。生产环境必须配置裁剪策略。 + +### MAXLEN —— 按数量裁剪 + +```bash +# 写入时裁剪(推荐,边写边裁) +XADD mystream MAXLEN 10000 * field value + +# 精确裁剪 vs 近似裁剪 +XADD mystream MAXLEN 10000 * field value # 精确,O(N) 开销 +XADD mystream MAXLEN ~ 10000 * field value # 近似,~ 模式,O(1) 开销 +``` + +### MINID —— 按 ID 裁剪(Redis 6.2+) + +```bash +# 删除 ID 小于指定值的消息 +XADD mystream MINID 1685000000000-0 * field value +XADD mystream MINID ~ 1685000000000-0 * field value # 近似模式 +``` + +### XTRIM —— 手动裁剪 + +```bash +XTRIM mystream MAXLEN ~ 5000 +XTRIM mystream MINID 1685000000000-0 +``` + +> [!warning] 近似裁剪(`~`)的取舍 +> `~` 模式不会精确保留 N 条,实际保留数量可能略多(因为底层按 Radix Tree 节点边界裁剪)。对于绝大多数场景,近似裁剪足够,而且**性能显著优于精确裁剪**。建议生产环境一律使用 `~`。 + +> [!question] 裁剪会影响 pending 中的消息吗? +> 会。`MAXLEN`/`MINID` 会删除 stream 中的原始消息,但 **pending 列表中的引用不会被清除**。这意味着 `XPENDING` 仍会显示这些 ID,但 `XCLAIM` 尝试读取时会返回 nil。设计裁剪策略时需确保消费者的处理速度跟得上写入速度。 + +## 六、Stream vs Pub/Sub vs Kafka 对比 + +| 维度 | Stream | Pub/Sub | Kafka | +|------|--------|---------|-------| +| **持久化** | 内存 + AOF | 不持久化 | 磁盘日志 | +| **消费者组** | 原生支持 | 不支持 | 原生支持 | +| **投递语义** | 至少一次(配合 ACK) | 最多一次(fire-and-forget) | 至少一次 / 精确一次 | +| **消息回溯** | 支持(按 ID 范围) | 不支持 | 支持(按 offset) | +| **吞吐量** | 中等(10w+ QPS) | 高(无持久化开销) | 极高(百万级 QPS) | +| **运维复杂度** | 低(复用 Redis) | 低 | 高(独立集群) | +| **适用规模** | 中小规模 | 实时通知 | 大规模事件流 | + +> [!tip] 如何选型? +> - **已有 Redis、消息量不大**(< 10w/s):Stream 是最佳选择,零额外运维成本 +> - **实时推送、允许丢消息**:Pub/Sub 更简单直接 +> - **海量事件流、需要精确一次语义**:Kafka 不可替代 + +## 七、死信队列(Dead Letter Queue) + +当某条消息被重试 N 次仍然失败时,继续重试没有意义。此时应将其移入**死信队列**,等待人工介入或补偿处理。 + +### 实现思路 + +Stream 本身没有内置 DLQ,但可以通过 XCLAIM 的 `delivery-count` 自行实现: + +```go +const maxRetries = 5 + +func (c *StreamConsumer) ProcessPendingWithDLQ(ctx context.Context) { + pendings, _ := c.rdb.XPendingExt(ctx, &redis.XPendingExtArgs{ + Stream: c.stream, + Group: c.group, + Start: "-", + End: "+", + Count: 100, + MinIdle: 30 * time.Second, + }).Result() + + for _, p := range pendings { + if p.RetryCount >= maxRetries { + // 超过重试次数 → 写入死信 Stream + c.moveToDLQ(ctx, p.ID) + continue + } + // 认领并重试 + claimed, _ := c.rdb.XClaim(ctx, &redis.XClaimArgs{ + Stream: c.stream, + Group: c.group, + Consumer: c.consumer, + MinIdle: 30 * time.Second, + Messages: []string{p.ID}, + }).Result() + for _, msg := range claimed { + if err := c.handler(ctx, msg); err == nil { + c.rdb.XAck(ctx, c.stream, c.group, msg.ID) + } + } + } +} + +func (c *StreamConsumer) moveToDLQ(ctx context.Context, msgID string) { + // 读取原始消息内容(XCLAIM 的返回可能为空,用 XRANGE 兜底) + msgs, _ := c.rdb.XRange(ctx, c.stream, msgID, msgID).Result() + if len(msgs) == 0 { + return + } + // 写入死信 Stream(附带原始 ID 以便溯源) + dlqFields := msgs[0].Values + dlqFields["original_id"] = msgID + dlqFields["original_stream"] = c.stream + c.rdb.XAdd(ctx, &redis.XAddArgs{ + Stream: c.stream + ":dlq", + Values: dlqFields, + }) + // 从原 Stream 确认移除 + c.rdb.XAck(ctx, c.stream, c.group, msgID) +} +``` + +> [!tip] DLQ 最佳实践 +> - 死信 Stream 命名建议 `:dlq`,保持关联性 +> - 在 DLQ 消息中记录 `original_id`、`original_stream`、`error_reason`,方便排查 +> - 监控 DLQ 长度,超过阈值立即告警——DLQ 堆积意味着业务逻辑有系统性问题 + +## 八、消息完整生命周期 + +```mermaid +flowchart TD + P["Producer"] -->|"XADD"| S["Stream"] + S -->|"XREADGROUP"| C1["Consumer A"] + S -->|"XREADGROUP"| C2["Consumer B"] + C1 -->|"XACK success"| ACK["Confirmed"] + C1 -->|"handler fail"| PEL["Pending Entry List"] + C2 -->|"XACK success"| ACK + C2 -->|"handler fail"| PEL + PEL -->|"retry count < max"| C1 + PEL -->|"retry count >= max"| DLQ["Dead Letter Queue"] + ACK -->|"MAXLEN trim"| TRIM["Trimmed"] + DLQ -->|"人工处理"| FIX["Fixed"] + + style S fill:#dfd,stroke:#090 + style PEL fill:#fdd,stroke:#f00 + style DLQ fill:#ffd,stroke:#f90 + style ACK fill:#dff,stroke:#099 + style TRIM fill:#eee,stroke:#999 +``` + +## 关联笔记 + +- [[hhs/Redis/09-高级特性]] — Pub/Sub 基础、Lua 脚本原子操作 +- [[hhs/Redis/03-基本命令速查]] — Redis 命令速查手册 +- [[hhs/Redis/README]] — 知识索引总览 diff --git a/hhs/Redis/16-GEO.md b/hhs/Redis/16-GEO.md new file mode 100644 index 0000000..97751c6 --- /dev/null +++ b/hhs/Redis/16-GEO.md @@ -0,0 +1,254 @@ +--- +tags: [Redis, 缓存, 数据结构, GEO] +create time: 2026-05-24 10:00 +--- + +# Redis GEO 地理位置 + +## 概述 + +GEO 是 Redis 3.2 引入的地理位置数据类型,底层基于 Sorted Set 实现,通过 GeoHash 编码将经纬度转换为可排序的分数值。它让我们能用简洁的命令完成"附近的人/门店"等 LBS(Location-Based Service)场景,无需依赖外部地理数据库。 + +## 正文 + +### 1. 底层原理 + +> [!question] GEO 既然是独立类型,为什么说它"没有自己的数据结构"? +> 因为 GEO 底层就是 Sorted Set。Redis 只是在 ZSet 之上封装了一层经纬度编解码逻辑,并没有新增底层数据结构。 + +核心编码流程: + +``` +经纬度 (latitude, longitude) + ↓ GeoHash 编码 + Base32 字符串 + ↓ 转换 + 52-bit 整数 → 存为 ZSet 的 score + ↓ + member 名称 → 存为 ZSet 的 value +``` + +GeoHash 的本质是将二维坐标递归二分,交替切分经度和纬度,得到一个二进制串,再编码为 Base32 字符串。切分越细,精度越高(字符串越长)。 + +> [!tip] 为什么 GeoHash 能支持范围查询? +> GeoHash 具有 **前缀匹配特性**:地理上越近的点,其 GeoHash 前缀越相似。将 GeoHash 转为 52-bit 整数后,相邻地理位置的分数也相近,这使得 ZSet 的 `ZRANGEBYSCORE` 天然支持范围查询。 + +**GeoHash 精度表**: + +| 字符串长度 | 精度 (km) | 适用场景 | +|-----------|----------|---------| +| 4 | ~40 | 城市级 | +| 5 | ~5 | 区县级 | +| 6 | ~0.6 | 街道级 | +| 7 | ~0.076 | 建筑级 | +| 8 | ~0.019 | 高精度 | + +Redis GEO 默认使用 **11 位精度**(约 0.019m),足以满足绝大多数 LBS 场景。 + +```mermaid +graph TD + A["输入经纬度"] --> B["GeoHash 编码"] + B --> C["52-bit 整数"] + C --> D["存入 ZSet score"] + D --> E["范围查询时按 score 排序"] + E --> F["解码还原经纬度并计算实际距离"] +``` + +### 2. 核心命令 + +> [!question] GEORADIUS 命令好用,为什么 Redis 6.2 要废弃它? +> `GEORADIUS` / `GEORADIUSBYMEMBER` 参数组合繁多,语义复杂。Redis 6.2 引入 `GEOSEARCH` 和 `GEOSEARCHSTORE`,用更清晰的接口统一了圆形和矩形搜索。 + +#### 常用命令速查 + +| 命令 | 说明 | +|------|------| +| `GEOADD key lng lat member` | 添加一个或多个地理位置 | +| `GEOPOS key member` | 获取成员的经纬度 | +| `GEODIST key m1 m2 [unit]` | 计算两个成员之间的距离 | +| `GEOSEARCH key FROMMEMBER/FROMLONLAT ... BYRADIUS/BYBOX ...` | 搜索附近成员 (6.2+) | +| `GEOSEARCHSTORE dst src ...` | 搜索结果存入目标 key (6.2+) | +| `GEOHASH key member` | 返回成员的 GeoHash 字符串 | + +**已废弃命令**(6.2+): + +- `GEORADIUS` — 用 `GEOSEARCH` 替代 +- `GEORADIUSBYMEMBER` — 用 `GEOSEARCH ... FROMMEMBER` 替代 + +#### 基本操作示例 + +```bash +# 添加门店坐标 +GEOADD stores 116.405285 39.904989 "store:1001" +GEOADD stores 121.473701 31.230416 "store:1002" + +# 获取坐标 +GEOPOS stores "store:1001" + +# 计算距离(默认单位:米) +GEODIST stores "store:1001" "store:1002" km + +# 搜索北京某点 5km 内的门店,返回距离,最多 10 个 +GEOSEARCH stores FROMLONLAT 116.405285 39.904989 \ + BYRADIUS 5 km ASC COUNT 10 WITHDIST +``` + +### 3. 附近门店实战 + +以下是一个完整的 Go 示例,使用 `go-redis/v9` 实现"查找附近门店"功能。 + +```go +package main + +import ( + "context" + "fmt" + "log" + + "github.com/redis/go-redis/v9" +) + +// Store 门店信息 +type Store struct { + Name string + Dist float64 // 距离,单位:米 +} + +func main() { + ctx := context.Background() + rdb := redis.NewClient(&redis.Options{ + Addr: "localhost:6379", + }) + defer rdb.Close() + + // ---- 1. 批量添加门店坐标 ---- + locations := []*redis.GeoLocation{ + {Name: "store:1001", Longitude: 116.405285, Latitude: 39.904989}, + {Name: "store:1002", Longitude: 116.407496, Latitude: 39.908712}, + {Name: "store:1003", Longitude: 116.397827, Latitude: 39.906419}, + } + _, err := rdb.GeoAdd(ctx, "stores", locations...).Result() + if err != nil { + log.Fatal(err) + } + + // ---- 2. 搜索附近 3km 内的门店 ---- + // 按距离升序排列,最多返回 5 个 + searchOpts := &redis.GeoSearchLocationQuery{ + GeoSearchQuery: redis.GeoSearchQuery{ + Longitude: 116.405285, + Latitude: 39.904989, + Radius: 3, + RadiusUnit: "km", + Sort: "ASC", + Count: 5, + }, + WithCoord: true, + WithDist: true, + } + results, err := rdb.GeoSearchLocation(ctx, "stores", searchOpts).Result() + if err != nil { + log.Fatal(err) + } + + // ---- 3. 处理结果 ---- + for _, r := range results { + fmt.Printf("门店: %s, 距离: %.2f km, 坐标: (%.6f, %.6f)\n", + r.Name, r.Dist, r.Longitude, r.Latitude) + } +} +``` + +> [!tip] `COUNT` 参数是软限制 +> `COUNT 5` 表示最多返回 5 个结果,但 Redis 内部仍会遍历所有候选点再做距离校验。如果业务数据量很大,建议结合分页或缩小搜索半径。 + +### 4. 距离计算 + +`GEODIST` 和 `GEOSEARCH` 的距离计算基于 **Haversine 公式**,该公式通过球面三角学计算地球表面两点之间的大圆距离。 + +Haversine 公式核心: + +$$a = \sin^2\left(\frac{\Delta\phi}{2}\right) + \cos\phi_1 \cdot \cos\phi_2 \cdot \sin^2\left(\frac{\Delta\lambda}{2}\right)$$ + +$$d = 2R \cdot \arcsin(\sqrt{a})$$ + +其中 $R \approx 6371$ km(地球半径)。 + +> [!question] Redis 的距离计算精度如何? +> Redis 使用 WGS84 参考椭球体近似为正球体,最大误差约 **0.5%**。对于 LBS 场景(几公里范围),误差通常在几十米以内,完全可接受。 + +### 5. 围栏检测 + +电子围栏(Geofencing)判断某个点是否在指定范围内。对于简单的圆形围栏,可以直接用 `GEODIST`;对于多边形围栏,需要自定义计算。 + +#### 圆形围栏 + +```bash +# 判断 member 是否在圆心 1km 范围内 +GEOSEARCH fence FROMLONLAT 116.405285 39.904989 \ + BYRADIUS 1 km ASC COUNT 1 +``` + +#### 多边形围栏(射线法) + +```go +// PointInPolygon 射线法判断点是否在多边形内 +func PointInPolygon(lat, lng float64, polygon [][2]float64) bool { + inside := false + j := len(polygon) - 1 + for i := 0; i < len(polygon); i++ { + xi, yi := polygon[i][1], polygon[i][0] + xj, yj := polygon[j][1], polygon[j][0] + if (yi > lat) != (yj > lat) && + lng < (xj-xi)*(lat-yi)/(yj-yi)+xi { + inside = !inside + } + j = i + } + return inside +} +``` + +> [!tip] 围栏方案选择 +> - **圆形围栏**:直接使用 GEO 命令,零额外开发成本。 +> - **矩形围栏**:使用 `GEOSEARCH ... BYBOX` 指定宽高。 +> - **不规则多边形**:先用 GEO 圈定一个粗略范围缩小候选集,再用射线法精确判断,兼顾性能与精度。 + +### 6. 附近门店搜索流程 + +```mermaid +graph TD + A["客户端发送请求"] --> B["获取用户经纬度"] + B --> C["GEOSEARCH 搜索半径 R"] + C --> D{"结果数量 > 0?"} + D -->|No| E["扩大半径 R += step"] + E --> C + D -->|Yes| F["返回门店列表及距离"] + F --> G["客户端渲染附近门店"] +``` + +### 7. 性能与限制 + +> [!question] GEO 能存多少数据? +> GEO 底层是 ZSet,受 ZSet 最大容量限制:单个 key 最大约 **150GB**(约 $2^{32}$ 个元素),足以存储数十亿个地理位置。 + +| 限制项 | 说明 | +|-------|------| +| 单 key 最大元素数 | ~$2^{32}$(约 43 亿) | +| 单 key 最大内存 | ~150GB(ZSet 限制) | +| GeoHash 精度 | 52-bit,约 0.019m | +| 最小搜索半径 | 0.000001(单位取决于指定) | +| 经度范围 | -180 ~ 180 | +| 纬度范围 | -85.05112878 ~ 85.05112878 | + +**性能建议**: + +- 搜索操作时间复杂度为 $O(N + \log M)$,其中 $N$ 为结果集大小,$M$ 为 key 中元素总数。大 key 下应控制搜索半径和 COUNT。 +- 如果只需坐标而不需要距离计算,可用 `GEOPOS` 代替 `GEOSEARCH` 的 `WITHDIST`,减少计算开销。 +- 对于超高并发场景,考虑将 GEO 数据加载到应用进程内存中做本地计算,减轻 Redis 压力。 + +## 关联笔记 + +- [[hhs/Redis/02-核心数据类型]] +- [[hhs/Redis/08-SortedSet精解]] +- [[hhs/Redis/README]] diff --git a/hhs/Redis/README.md b/hhs/Redis/README.md index 0c20963..76ceb36 100644 --- a/hhs/Redis/README.md +++ b/hhs/Redis/README.md @@ -34,7 +34,7 @@ create time: 2026-05-15 18:10 |---|------|------| | 2.1 | RDB 快照 | [[hhs/Redis/04-RDB持久化]] — fork 原理、COW、触发时机 | | 2.2 | AOF 追加日志 | [[hhs/Redis/05-AOF持久化]] — everysec 策略、重写流程 | -| 2.3 | 混合持久化 | Redis 7 特性:RDB 快照 + AOF 增量 = 快 + 准 | +| 2.3 | 混合持久化 | [[hhs/Redis/04-RDB持久化#Redis 混合持久化(推荐)]] — RDB 快照 + AOF 增量 = 快 + 准 | | 2.4 ~ 2.5 | 主从复制与 Sentinel | [[hhs/Redis/06-主从与哨兵]] — 全量/增量同步、故障转移流程 | | 2.6 | Cluster 集群 | [[hhs/Redis/07-集群方案]] — 哈希槽、Gossip、扩容迁移 | @@ -57,11 +57,11 @@ flowchart LR | # | 主题 | 链接 | |---|------|------| | 3.1 | Sorted Set 高级玩法 | [[hhs/Redis/08-SortedSet精解]] — 排行榜、延迟队列、Geo、优先级队列 | -| 3.2 | Bitmap | 签到、DAU 统计——每位一天,极致省内存 | -| 3.3 | HyperLogLog | 亿级 UV 去重统计,误差 < 0.1%,仅占 12KB | -| 3.4 | Bloom Filter | 防缓存穿透、URL 去重——false positive 可接受 | -| 3.5 | Stream | 可靠消息队列,Consumer Group 消费组模型 | -| 3.6 | GEO | 附近的人、距离计算、围栏检测——底层就是 ZSet | +| 3.2 | Bitmap | [[hhs/Redis/12-Bitmap]] — 签到、DAU 统计——每位一天,极致省内存 | +| 3.3 | HyperLogLog | [[hhs/Redis/13-HyperLogLog]] — 亿级 UV 去重统计,误差 < 0.1%,仅占 12KB | +| 3.4 | Bloom Filter | [[hhs/Redis/14-BloomFilter]] — 防缓存穿透、URL 去重——false positive 可接受 | +| 3.5 | Stream | [[hhs/Redis/15-Stream]] — 可靠消息队列,Consumer Group 消费组模型 | +| 3.6 | GEO | [[hhs/Redis/16-GEO]] — 附近的人、距离计算、围栏检测——底层就是 ZSet | ```mermaid flowchart TD diff --git a/hhs/Redis/一文吃透Redis.md b/hhs/某公司文档/docs/一文吃透Redis.md similarity index 100% rename from hhs/Redis/一文吃透Redis.md rename to hhs/某公司文档/docs/一文吃透Redis.md