Redis Stream 的 XREADGROUP 配置了批量读取,并不代表业务处理会并发执行。一次实际排查中,Worker 的读取数量是 10,但消息处理函数仍在一个循环里同步调用,因此生成任务表现为“只能一个接一个执行”。修复后,单个 Worker 默认可以同时处理 4 个任务,并且并发上限可以通过环境变量调整。
问题现象
项目是 Go 后端、Redis Streams 队列和独立 Worker 的 AI 视频生成服务。任务创建后先进入 generation.tasks Stream,Worker 消费消息,再调用 Provider 完成提交、轮询、下载和结算。
开发配置中曾经有这样的读取参数:
1 | queue: |
直觉上容易把 readCount: 10 理解成“同时处理 10 个任务”。实际观察却是:一个任务完成前,下一个任务不会进入 Provider 执行。问题表现为队列消息可能已经被 Worker 领取,但生成吞吐量仍接近单任务串行。
这里要区分两个指标:
| 指标 | 含义 | 是否等于并发数 |
|---|---|---|
XREADGROUP COUNT |
一次 Redis 读取最多返回多少条消息 | 否 |
| Worker concurrency | 一个进程同时运行多少个业务处理函数 | 是 |
排查过程
1. 从队列读取入口开始看
server/internal/infrastructure/queue/redisstream/consumer.go 的消费流程先构造 Stream 列表,再调用 XReadGroup:
1 | results, err := c.client.XReadGroup(ctx, &redis.XReadGroupArgs{ |
到这里最多只能证明“消息被批量读出来了”,还不能证明业务处理并发。
2. 检查消息处理循环
原来的 handleMessages 对消息逐条调用处理器:
1 | for _, msg := range messages { |
handler(message) 返回前,循环不会处理下一条消息。更重要的是,ACK 也在处理器成功返回后执行,所以不能简单地让处理器“提前返回”来制造并发,否则会在任务仍未完成时确认消息,导致任务丢失。
3. 排除“是不是 Worker 只有一个实例”的误判
server/internal/application/worker/worker_service.go 只注册了一个 generation.tasks 处理器,并把消费循环交给 QueueConsumer.Run。因此问题不是“注册了多个处理器却互相覆盖”,而是 Redis Consumer 内部只有一个同步消费循环。
同时,readCount 只是读取参数,没有任何并发配置字段。也就是说,配置文件写了 10,并不会自动创建 10 个 Go 执行单元。
根因
直接原因是“批量读取”和“并发处理”被当成了同一件事:
1 | XREADGROUP Count=10 |
只要 Provider 的一次生成需要等待网络请求、上游轮询和结果下载,当前 goroutine 就会被完整占用。Redis 可以把消息交给这个 Consumer,但 Go 代码没有第二个执行单元去处理下一条任务。
解决方案
1. 增加显式并发配置
在 QueueRedis 中增加独立的 Concurrency 字段,并保留安全默认值 4:
1 | type QueueRedis struct { |
开发配置显式写出并发数:
1 | queue: |
生产环境可以通过环境变量覆盖:
1 | QUEUE_REDIS_CONCURRENCY=4 |
配置加载时,非正数会回退到默认值 4;这样不会因为遗漏配置而退回到单任务串行。
2. 启动有界的消费循环
Consumer 根据并发配置启动多个消费循环,每个循环都有同一个 context.Context 和明确的退出条件:
1 | for index := range c.concurrency { |
每个循环只用 XREADGROUP COUNT=1 领取一条新消息:
1 | // 让领取数量与并发槽位保持一致,避免先领取一大批再排队等待。 |
这样,默认 4 个循环最多同时进入 4 个任务处理函数。消息只有在处理器成功返回之后才 ACK;处理失败时仍按原逻辑重新入队或进入死信队列。
3. 只让一个循环负责 pending 回收
旧 pending 的回收仍然通过 XAUTOCLAIM 完成,但只由第一个消费循环执行:
1 | reclaimPending := index == 0 |
其他循环直接领取新消息,避免多个 goroutine 同时扫描和争抢 pending。回收任务与新任务最终仍共用同一个并发上限。
4. 让等待逻辑可以响应取消
原来的重试延迟使用 time.Sleep,进程收到停止信号后可能还要等完整的重试间隔。修复后使用带 Context 的定时等待:
1 | select { |
这不会中断正在执行的 Provider 调用,但可以让空闲消费循环和错误重试循环及时退出。
验证与回归
并发回归测试
新增 consumer_test.go 中的 fake Stream Client,准备 3 条消息并把并发设置为 2。测试让前两个处理器阻塞,确认:
- 两条不同消息可以同时进入处理器;
- 达到并发上限时,第三条消息仍留在队列中,没有被提前领取;
- 释放前两个任务后,第三条消息才开始处理;
- 三条消息最终都完成 ACK;
- 取消 Context 后消费循环可以退出。
这个测试直接验证“并发处理”和“并发上限”两个行为,不依赖真实 Redis,也不会产生 Provider 成本。
实际执行过的检查
在 server/ 目录执行:
1 | go test ./... -count=1 |
以上检查均通过。随后执行:
1 | docker compose config --quiet |
Compose 配置解析通过。开发环境 Worker 热重载后的日志包含:
1 | queue initialized ... group=video-worker ... concurrency=4 |
说明运行时已经读取了新的并发配置。
常见误区与边界
readCount 越大,并发就越高
不是。读取批量大小和业务执行并发是两个独立参数。当前实现故意让新消息每次只领取一条,readCount 主要用于 XAUTOCLAIM 回收旧 pending。
把 ACK 放到 goroutine 外面
不应该。ACK 必须发生在业务处理成功之后。若消息刚交给 goroutine 就立即 ACK,Worker 崩溃时 Redis 不会再把这条消息交给其他 Consumer,可能造成付费任务丢失。
并发设置得越大越好
并发上限是“单个 Worker 进程”的上限。多个 Worker 副本会把总并发继续相乘;同时还要服从 Provider 的 RPM、并发额度、网络带宽、数据库连接池和本地磁盘能力。建议先从 4 开始,根据 Provider 限额和成功率逐步调整。
15 分钟 pending 回收窗口可以解决所有重复执行
不能。长时间同步调用如果超过 pending 空闲窗口,其他 Consumer 仍可能尝试回收这条消息。对于 Provider 已经受理但平台尚未拿到结果的任务,仍需要依赖任务幂等、Attempt checkpoint 和上游查询能力,不能只靠提高 Worker 并发解决可靠性问题。
可复用检查清单
- 分别记录
XREADGROUP COUNT和业务处理并发,不要混用两个概念。 - 检查消息处理器是否在同一个
for循环中同步调用。 - 确认 ACK 只发生在处理成功之后。
- 用显式配置控制单进程并发,并给出安全默认值。
- 新消息领取数量不要超过当前可用执行槽位。
- 为并发上限、ACK、取消和错误重试写回归测试。
- 根据 Provider 限额、数据库连接池和副本数计算总并发。
- 检查长任务、pending 回收和上游已受理任务的重复执行边界。