《Go 语言编程实战》8.1 生产者/消费者与可靠投递

用 Redis Streams 给 TaskHub 搭一套生产者/消费者:XADD 生产并裁剪、XREADGROUP 消费组、XACK 确认、XPENDING 查积压、XAUTOCLAIM 接管崩溃消费者的未确认消息,把「至少一次投递」真正落地并对比三种投递语义。

本节把 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 StreamsKafkaRabbitMQ
消费组 / 多消费者✅ XREADGROUP✅✅
消息确认✅ XACKoffset 提交✅ ack
未确认追踪✅ XPENDINGlag✅ unacked
崩溃接管✅ XAUTOCLAIMrebalance✅ 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 幂等、重试与死信 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「golang」更多文章

  1. 《Go 语言编程实战》目录
  2. 《Go 语言编程实战》18.3 上线、观测与迭代
  3. 《Go 语言编程实战》18.2 故障演练