From f7c3c83c2be2dc84be4698f87368dec6541d5bf3 Mon Sep 17 00:00:00 2001 From: hhs <386998068@qq.com> Date: Tue, 26 May 2026 00:03:40 +0800 Subject: [PATCH] vault backup: 2026-05-26 00:03:40 --- hhs/Redis/15-Stream.md | 73 +++++++++++++++++++++++++++++++++++++++--- hhs/Redis/16-GEO.md | 2 +- 2 files changed, 69 insertions(+), 6 deletions(-) diff --git a/hhs/Redis/15-Stream.md b/hhs/Redis/15-Stream.md index 8a431f3..4d90bd0 100644 --- a/hhs/Redis/15-Stream.md +++ b/hhs/Redis/15-Stream.md @@ -40,15 +40,55 @@ XACK mystream mygroup # 确认消费完成 > |---|------|---------| > | `0` | Stream 的第一条消息 | 数据重放、全量回溯 | > | `$` | 当前最新消息的 ID | `XGROUP CREATE` 时指定起始点 | -> | `>` | 只返回**未被任何消费者分配**的新消息 | `XREADGROUP` 正常消费循环 | +> | `>` | 只返回组内**尚未投递给任何消费者**的消息 | `XREADGROUP` 正常消费循环 | +> +> 注意:`>` 只在 `XREADGROUP` 中有意义。如果传 `0` 读取消费者组,则返回该消费者的 **pending 消息**而非新消息——这在"重启后恢复未完成任务"场景中非常有用。 -### Stream 底层结构 +### Stream ID 自动生成机制 -每条消息以 Radix Tree + Listpack 存储,时间复杂度 O(1) 追加。ID 格式为 `<毫秒时间戳>-<序号>`,天然有序。 +ID 格式为 `<毫秒时间戳>-<序号>`,由 Redis 服务器生成(用 `*` 时)。规则如下: + +1. **时间戳部分**:取当前毫秒级 Unix 时间戳 +2. **序号部分**:同一毫秒内递增(从 0 开始);如果时间戳前进到下一毫秒,序号重置为 0 + +```bash +XADD mystream * k v # 返回 1685000000000-0 +XADD mystream * k v # 同一毫秒内 → 1685000000000-1 +XADD mystream * k v # 跨毫秒后 → 1685000000001-0 +``` + +> [!question] 能手动指定 ID 吗? +> 可以,但必须**严格递增**。新 ID 必须大于 Stream 中已有最大 ID,否则报错 `(ERR) The ID specified in XADD is equal or smaller than the target stream top item`。手动指定 ID 的典型场景:数据迁移、事件溯源回放。 + +### 底层存储结构 + +每条消息以 Radix Tree + Listpack 存储,追加写入时间复杂度 O(1)。同一时间戳内的多条消息会紧凑打包在 Listpack 节点中,极大节省内存。 > [!question] Stream 会像 List 一样消费后删除吗? > 不会。Stream 是**只追加日志**,消息消费后仍在。需要显式 `XDEL` 或通过 `MAXLEN`/`MINID` 策略裁剪。这带来了「可回溯」的优势,但也意味着必须主动管理容量。 +### 历史消息读取 + +消费者组模式适用于"实时消费",但有时你需要**按范围回溯历史消息**——比如排障、数据对账、事件溯源。这时用 `XRANGE`/`XREVRANGE`: + +```bash +# 按 ID 范围读取:从最小 ID 到最大 ID,最多 10 条 +XRANGE mystream - + COUNT 10 +# - = 最小 ID,+ = 最大 ID + +# 按时间范围读取(ID 只需写时间戳部分,序号默认为 0) +XRANGE mystream 1685000000000 1685000060000 COUNT 10 + +# 反向读取(从新到旧) +XREVRANGE mystream + - COUNT 5 + +# 获取 stream 长度 +XLEN mystream +``` + +> [!tip] `XRANGE` vs `XREAD` +> `XREAD` 适合**持续阻塞消费**(配合 `BLOCK`),是生产者的标准读法。`XRANGE` 适合**一次性按范围查询**,不阻塞,常用于排障和数据审计。两者互补,不要混淆。 + ## 二、Consumer Group 机制详解 Consumer Group 是 Stream 的核心抽象,它让多个消费者**协作消费同一条 Stream**,每条消息只会被组内的一个消费者处理。 @@ -123,6 +163,24 @@ sequenceDiagram > [!question] 消息是如何分配给消费者的? > Redis 采用**轮询分配**(round-robin):新消息到达时,依次分配给组内活跃的消费者。注意不是负载均衡——如果某个消费者处理慢,它积压的 pending 消息不会自动转移给其他消费者,需要通过 XCLAIM 手动认领。 +### XINFO —— 运维监控利器 + +排查 Stream 问题时,`XINFO` 是你的第一选择: + +```bash +# 查看 stream 元信息(长度、首尾 ID、消费者组数量) +XINFO STREAM mystream + +# 查看消费者组列表 +XINFO GROUPS mystream + +# 查看组内各消费者状态(pending 数、空闲时间) +XINFO CONSUMERS mystream mygroup +``` + +> [!tip] 运维小贴士 +> 监控脚本可以定期执行 `XINFO CONSUMERS`,如果某个消费者的 `idle` 值远超阈值且 pending 不为 0,说明该消费者可能已宕机,需要触发 XCLAIM 或告警。 + ## 三、Go 实战 ### 生产者 @@ -405,10 +463,15 @@ 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 { + // 消息可能已被裁剪删除,仅 ACK 清理 pending + c.rdb.XAck(ctx, c.stream, c.group, msgID) return } - // 写入死信 Stream(附带原始 ID 以便溯源) - dlqFields := msgs[0].Values + // 拷贝一份,避免污染原始消息的 Values map + dlqFields := make(map[string]interface{}, len(msgs[0].Values)+2) + for k, v := range msgs[0].Values { + dlqFields[k] = v + } dlqFields["original_id"] = msgID dlqFields["original_stream"] = c.stream c.rdb.XAdd(ctx, &redis.XAddArgs{ diff --git a/hhs/Redis/16-GEO.md b/hhs/Redis/16-GEO.md index 97751c6..7d43fcf 100644 --- a/hhs/Redis/16-GEO.md +++ b/hhs/Redis/16-GEO.md @@ -43,7 +43,7 @@ GeoHash 的本质是将二维坐标递归二分,交替切分经度和纬度, | 7 | ~0.076 | 建筑级 | | 8 | ~0.019 | 高精度 | -Redis GEO 默认使用 **11 位精度**(约 0.019m),足以满足绝大多数 LBS 场景。 +Redis GEO 将 GeoHash 编码为 **52-bit 整数**(对应约 10 位字符精度),水平方向误差约 **±1m**,垂直方向约 **±0.6m**,足以满足绝大多数 LBS 场景。 ```mermaid graph TD