本节把 TaskHub 推进到「异步」:把发通知、生成报表、跨系统同步这些不该卡在请求线程里的活儿,改由 Redis Streams 承载,用消费组实现可水平扩展的消费者,并保证消息不会因为消费者崩溃而丢失。
适用版本:Go 1.27(实测go1.27.0),Redis 7(redis:7-alpine容器)。
8.1 生产者/消费者与可靠投递
第 7 章解决了读路径的性能,但 TaskHub 里还有一批不该同步做的工作:任务创建后要发邮件通知、要写审计日志、要同步到下游系统。这些操作要么慢,要么可能失败,把它们塞进 HTTP 请求里,用户就得等、还可能因为下游抖动而整个请求失败。消息队列的作用是把「接收请求」和「处理副作用」解耦。
8.1.1 TaskHub 的异步边界
先明确哪些工作该异步。判据是:它是不是请求成功所必需的。
| 工作 | 同步还是异步 | 理由 |
|---|---|---|
| 写入任务记录 | 同步 | 不成功用户就拿不到结果 |
| 发通知邮件 | 异步 | 失败不该影响任务创建 |
| 生成周报 | 异步 | 耗时长,用户不等 |
| 同步到下游系统 | 异步 | 下游可能抖动,需重试 |
| 更新搜索索引 | 异步 | 可最终一致 |
异步的代价是复杂度:消息可能丢、可能重复、可能乱序。本节先把「可靠投递」这块最核心的做对。
8.1.2 为什么用 Redis Streams
本机没有实测 Kafka / RabbitMQ,本节优先用已实测可用的 Redis Streams承载消息语义。这不是妥协——Redis Streams 提供了消费组、确认、pending 列表、自动接管,足以覆盖 TaskHub 这个量级:
| 能力 | Redis Streams | Kafka | RabbitMQ |
|---|---|---|---|
| 消费组 / 多消费者 | ✅ XREADGROUP | ✅ | ✅ |
| 消息确认 | ✅ XACK | offset 提交 | ✅ ack |
| 未确认追踪 | ✅ XPENDING | lag | ✅ unacked |
| 崩溃接管 | ✅ XAUTOCLAIM | rebalance | ✅ requeue |
| 持久化 | ✅ AOF/RDB | ✅ 磁盘日志 | ✅ |
| 吞吐量级 | 万级 QPS | 十万级+ | 万级 |
| 运维成本 | 低(复用 Redis) | 高 | 中 |
本节所有 Redis Streams 示例均在本机 redis:7-alpine 容器实测;Kafka / RabbitMQ 未实测,上表仅为特性对比。 TaskHub 选 Redis Streams 是因为量级够用且复用现有 Redis,等吞吐真的顶不住再迁 Kafka。
8.1.3 生产:XADD 与裁剪
生产者往 stream 里追加消息,用 XADD。必须带 MAXLEN 裁剪,否则 stream 无限增长:
id, err := rdb.XAdd(ctx, &redis.XAddArgs{
Stream: "tasks:events",
MaxLen: 1000, Approx: true, // 近似裁剪,性能更好
Values: map[string]any{"task_id": "1", "event": "created"},
}).Result()
Approx: true 让 Redis 用「大约保留 N 条」的方式裁剪(按节点整块删),比精确裁剪快得多。实测三条消息的返回 ID:
$ go run ./ch8/stream
XADD -> 1791598775858-0
XADD -> 1791598775861-0
XADD -> 1791598775866-0
ID 格式是 毫秒时间戳-序号,天然有序,-0 表示同一毫秒内的第一条。生产端拿到 ID 后可以回写业务表,作为「已投递」的凭据。
8.1.4 消费组:XREADGROUP 与 XACK
消费者不直接读 stream,而是加入一个消费组。同一个组内的多个消费者瓜分消息(每条只被一个消费者拿到),不同组各消费一份。建组时用 0 表示从头消费:
rdb.XGroupCreateMkStream(ctx, "tasks:events", "workers", "0")
msgs, _ := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: "workers",
Consumer: "worker-A",
Streams: []string{"tasks:events", ">"}, // ">" 表示只读从未投递的新消息
Count: 1,
Block: -1, // 非阻塞;见下方说明
}).Result()
for _, m := range msgs[0].Messages {
fmt.Printf("A 收到 %s %v\n", m.ID, m.Values)
rdb.XAck(ctx, "tasks:events", "workers", m.ID) // 处理完确认
}
实测消费一条并确认:
$ go run ./ch8/stream
A 收到 1791598775858-0 map[event:created task_id:1]
这里有个必须记住的坑:go-redis 的 XReadGroupArgs.Block 默认值是 0,而 0 会被翻译成 Redis 的 BLOCK 0——永久阻塞。想要「没有新消息就立刻返回」,要显式写 Block: -1。我第一次写这段时忘了设,程序卡在 XReadGroup 上一直不返回。
一个常驻消费者的骨架大致是这样——循环拉取、处理、ACK,用 Block 的阻塞能力实现「没消息就等一会儿」,避免空转:
for {
msgs, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: "workers", Consumer: consumerID,
Streams: []string{"tasks:events", ">"},
Count: 50,
Block: 2 * time.Second, // 阻塞 2s,有消息立刻返回
}).Result()
if err == redis.Nil { // 超时无消息,继续下一轮
continue
}
if err != nil { // 网络等错误,退避后重试
time.Sleep(time.Second)
continue
}
for _, m := range msgs[0].Messages {
if err := handle(ctx, m); err == nil {
rdb.XAck(ctx, "tasks:events", "workers", m.ID)
}
// 处理失败不 ACK,留给 XAUTOCLAIM 重投
}
}
注意处理失败时故意不 ACK:消息留在 pending,稍后被 XAUTOCLAIM 重新投递,这就是「至少一次」的重投来源。consumerID 建议用「主机名 + PID」,这样 pending 列表里能看出消息卡在哪个实例。
8.1.5 至少一次投递与 pending 列表
关键点来了:消息被 XREADGROUP 投递出去的那一刻,并没有从 stream 里删除。它进入该消费组的 pending(待确认)列表,只有收到 XACK 才算处理完成。这正是「至少一次投递」的实现机制。
如果消费者读到了消息却没来得及 ACK 就崩溃,这条消息会永远留在 pending 里。用 XPENDING 能查到这个积压:
pend, _ := rdb.XPending(ctx, "tasks:events", "workers").Result()
fmt.Printf("XPENDING: count=%d\n", pend.Count)
$ go run ./ch8/stream
B 收到 1791598775861-0(不 ACK,模拟崩溃)
XPENDING: count=1
消费者 B 读走了 task_id=2 这条消息但没确认,pending 计数为 1。这条消息没有丢——它还在,等着被接管。
8.1.6 故障接管:XAUTOCLAIM
pending 里的消息属于「B」这个消费者,但 B 已经死了。Redis 的 XAUTOCLAIM 允许其它消费者把空闲超过一定时间的 pending 消息认领过来:
claimed, _, _ := rdb.XAutoClaim(ctx, &redis.XAutoClaimArgs{
Stream: "tasks:events",
Group: "workers",
Consumer: "worker-C",
MinIdle: 50 * time.Millisecond, // 空闲超过 50ms 才认领
Start: "0",
Count: 10,
}).Result()
for _, m := range claimed {
rdb.XAck(ctx, "tasks:events", "workers", m.ID)
}
MinIdle 是「这条消息多久没人管了」,设得太小会误抢正常消费者正在处理的消息,设得太大则故障恢复慢。实测 C 接管后 pending 清零:
$ go run ./ch8/stream
C 接管 1791598775861-0 map[event:created task_id:2]
接管后 XPENDING: count=0
注意 XAUTOCLAIM 第一次调用可能返回空——因为消息刚投递、空闲时间还没到 MinIdle。必须给足 idle 时间,实际部署里 MinIdle 通常设 30 秒到几分钟,比这里的 50ms 大得多。
8.1.7 三种投递语义
消息系统的投递语义要分清,TaskHub 选的是中间那种:
| 语义 | 做法 | 风险 | 适用 |
|---|---|---|---|
| 最多一次(at-most-once) | 读到即 ACK,再处理 | 处理中崩溃会丢消息 | 可丢的埋点、日志 |
| 至少一次(at-least-once) | 处理完才 ACK | 崩溃重投会重复处理 | TaskHub 的事件流 |
| 恰好一次(exactly-once) | 至少一次 + 消费端幂等 | 实现成本高 | 需要精确计费的场景 |
「至少一次」几乎总是和「消费端幂等」配套出现——因为重投导致的重复处理,必须由消费者自己去重。这正是 8.2 节的主题。Redis Streams 给的是至少一次,TaskHub 接受它,把去重做在业务侧。
8.1.8 吞吐与积压观测
TaskHub 量级不大,但必须知道这条链路的实际吞吐,才能判断要不要扩容消费者。用 pipeline 批量生产 5000 条、批量消费并 ACK 5000 条,实测(本机 colima 容器):
$ go run ./ch8/through
生产 5000 条: 62ms (80958 条/秒)
XLEN=5000 groups=1 pending=0
消费 5000 条: 96ms (51906 条/秒)
单机 Redis 下约 8 万条/秒生产、5 万条/秒消费,远超 TaskHub 当前的事件量。注意这是批量的结果:单条同步 XADD 会受 RTT 拖累(本机 colima 转发下单条约 1ms),生产端能批量就批量。
运维上要盯两个指标:
n, _ := rdb.XLen(ctx, "tasks:events").Result() // stream 总长度 ≈ 积压
groups, _ := rdb.XInfoGroups(ctx, "tasks:events").Result()
// groups[0].Pending 是「已投递未确认」,groups[0].Lag 是「还没投递给任何消费者」
XLEN 持续增长说明消费跟不上生产;Pending 持续增长说明消费者处理慢或卡住;Lag 大说明消费者数量不够。这三个数配合第 10 章的指标体系做成告警,才能第一时间发现积压。
8.1.9 顺序性:消费组会打乱单实体的顺序
一个容易被忽略的点:同一 stream 内的消息天然有序,但消费组会把它们分给不同消费者并行处理,从而丢失全局顺序。对 TaskHub 来说,task:42 的 created 和 updated 两条事件如果被两个消费者并行处理,可能先处理 updated 再处理 created,状态就错了。
Streams 不像 Kafka 那样有 partition key,要实现「同一实体的事件保序」,常用做法是按实体 ID 分片到多个 stream,每个分片由一个消费者独占消费:
// 按 task_id 哈希分到 8 个分片 stream
shard := fnv32(taskID) % 8
stream := fmt.Sprintf("tasks:events:%d", shard)
rdb.XAdd(ctx, &redis.XAddArgs{Stream: stream, Values: payload})
每个分片 stream 配一个独立消费组、只跑一个消费者,这样分片内严格有序,分片间并行。代价是分片数固定(不好动态扩)、热点实体可能压垮单个分片。如果业务允许乱序(比如纯通知),就用单 stream 多消费者追求吞吐,别为不存在的顺序需求付代价。
还有一个更简单的兜底:在消息体里带上版本号或时间戳,消费端遇到乱序时丢弃旧版本。这与 7.3 节的版本守卫是同一个思路——用逻辑顺序代替到达顺序。
8.1.10 常见坑
XReadGroup的Block: 0永久阻塞:默认值0会阻塞,要非阻塞得写Block: -1。- 忘了
XACK:消息永远留在 pending,pending 越积越多,XAUTOCLAIM也救不了。 XADD不带MAXLEN:stream 无限增长,内存被吃光。MinIdle设太小:正常处理中的消息被误抢,导致重复处理。- 只建组不消费:stream 只增不减,忘了消费组这条链路就废了。
- 把消息体做大:stream 存的是消息内容,大对象应该只传引用(如对象存储 key)。
- 多消费组误当多消费者:组内是瓜分、组间是广播,语义完全不同,用错会丢消息。
- 用 pending 数当业务积压:pending 只反映「投递未确认」,真正的业务积压还要看 stream 总长度与消费速率。
小结
- 异步的判据是「是否是请求成功所必需」,非必需的副作用(通知、报表、同步)走消息。
- Redis Streams 用
XADD(带MAXLEN)+ 消费组 +XACK提供至少一次投递;本机实测可用,Kafka/RabbitMQ 未实测。 XREADGROUP投递不等于消费完成,未 ACK 的消息进 pending,XPENDING可查。- 消费者崩溃时用
XAUTOCLAIM认领空闲消息,实现故障接管;MinIdle要按处理耗时合理设置。 - 至少一次必然带来重复,去重是下一节的事。
TaskHub 现在有了可靠投递的骨架,但「至少一次」意味着消息可能被处理两遍。下一节解决重复:消费幂等、失败重试与死信队列。
阅读导航:上一节:7.3 缓存一致性 · 下一节:8.2 幂等、重试与死信 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。