Redis Stream Worker 并发失效:为什么任务只能一个一个执行

  1. 1. 问题现象
  2. 2. 排查过程
    1. 2.1. 1. 从队列读取入口开始看
    2. 2.2. 2. 检查消息处理循环
    3. 2.3. 3. 排除“是不是 Worker 只有一个实例”的误判
  3. 3. 根因
  4. 4. 解决方案
    1. 4.1. 1. 增加显式并发配置
    2. 4.2. 2. 启动有界的消费循环
    3. 4.3. 3. 只让一个循环负责 pending 回收
    4. 4.4. 4. 让等待逻辑可以响应取消
  5. 5. 验证与回归
    1. 5.1. 并发回归测试
    2. 5.2. 实际执行过的检查
  6. 6. 常见误区与边界
    1. 6.1. readCount 越大,并发就越高
    2. 6.2. 把 ACK 放到 goroutine 外面
    3. 6.3. 并发设置得越大越好
    4. 6.4. 15 分钟 pending 回收窗口可以解决所有重复执行
  7. 7. 可复用检查清单

Redis Stream 的 XREADGROUP 配置了批量读取,并不代表业务处理会并发执行。一次实际排查中,Worker 的读取数量是 10,但消息处理函数仍在一个循环里同步调用,因此生成任务表现为“只能一个接一个执行”。修复后,单个 Worker 默认可以同时处理 4 个任务,并且并发上限可以通过环境变量调整。

问题现象

项目是 Go 后端、Redis Streams 队列和独立 Worker 的 AI 视频生成服务。任务创建后先进入 generation.tasks Stream,Worker 消费消息,再调用 Provider 完成提交、轮询、下载和结算。

开发配置中曾经有这样的读取参数:

1
2
3
queue:
redis:
readCount: 10

直觉上容易把 readCount: 10 理解成“同时处理 10 个任务”。实际观察却是:一个任务完成前,下一个任务不会进入 Provider 执行。问题表现为队列消息可能已经被 Worker 领取,但生成吞吐量仍接近单任务串行。

这里要区分两个指标:

指标 含义 是否等于并发数
XREADGROUP COUNT 一次 Redis 读取最多返回多少条消息
Worker concurrency 一个进程同时运行多少个业务处理函数

排查过程

1. 从队列读取入口开始看

server/internal/infrastructure/queue/redisstream/consumer.go 的消费流程先构造 Stream 列表,再调用 XReadGroup

1
2
3
4
5
6
7
results, err := c.client.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: c.group,
Consumer: c.consumerName,
Streams: streams,
Count: c.readCount,
Block: c.blockTimeout,
}).Result()

到这里最多只能证明“消息被批量读出来了”,还不能证明业务处理并发。

2. 检查消息处理循环

原来的 handleMessages 对消息逐条调用处理器:

1
2
3
4
5
6
7
8
for _, msg := range messages {
// 组装 QueueMessage
if err := handler(message); err == nil {
if err := c.ack(ctx, stream, msg.ID); err != nil {
return err
}
}
}

handler(message) 返回前,循环不会处理下一条消息。更重要的是,ACK 也在处理器成功返回后执行,所以不能简单地让处理器“提前返回”来制造并发,否则会在任务仍未完成时确认消息,导致任务丢失。

3. 排除“是不是 Worker 只有一个实例”的误判

server/internal/application/worker/worker_service.go 只注册了一个 generation.tasks 处理器,并把消费循环交给 QueueConsumer.Run。因此问题不是“注册了多个处理器却互相覆盖”,而是 Redis Consumer 内部只有一个同步消费循环。

同时,readCount 只是读取参数,没有任何并发配置字段。也就是说,配置文件写了 10,并不会自动创建 10 个 Go 执行单元。

根因

直接原因是“批量读取”和“并发处理”被当成了同一件事:

1
2
3
4
5
6
7
8
9
XREADGROUP Count=10

handleMessages

for message

handler(message) ← 同步阻塞

XACK

只要 Provider 的一次生成需要等待网络请求、上游轮询和结果下载,当前 goroutine 就会被完整占用。Redis 可以把消息交给这个 Consumer,但 Go 代码没有第二个执行单元去处理下一条任务。

解决方案

1. 增加显式并发配置

QueueRedis 中增加独立的 Concurrency 字段,并保留安全默认值 4:

1
2
3
4
type QueueRedis struct {
ReadCount int64 `yaml:"readCount"`
Concurrency int `yaml:"concurrency"`
}

开发配置显式写出并发数:

1
2
3
4
queue:
redis:
readCount: 10 # 回收旧 pending 时的批量大小
concurrency: 4 # 单个 Worker 同时处理的任务数

生产环境可以通过环境变量覆盖:

1
QUEUE_REDIS_CONCURRENCY=4

配置加载时,非正数会回退到默认值 4;这样不会因为遗漏配置而退回到单任务串行。

2. 启动有界的消费循环

Consumer 根据并发配置启动多个消费循环,每个循环都有同一个 context.Context 和明确的退出条件:

1
2
3
4
5
6
7
for index := range c.concurrency {
reclaimPending := index == 0
c.wg.Go(func() {
c.consumeLoop(ctx, reclaimPending)
})
}
c.wg.Wait()

每个循环只用 XREADGROUP COUNT=1 领取一条新消息:

1
2
// 让领取数量与并发槽位保持一致,避免先领取一大批再排队等待。
Count: 1,

这样,默认 4 个循环最多同时进入 4 个任务处理函数。消息只有在处理器成功返回之后才 ACK;处理失败时仍按原逻辑重新入队或进入死信队列。

3. 只让一个循环负责 pending 回收

旧 pending 的回收仍然通过 XAUTOCLAIM 完成,但只由第一个消费循环执行:

1
reclaimPending := index == 0

其他循环直接领取新消息,避免多个 goroutine 同时扫描和争抢 pending。回收任务与新任务最终仍共用同一个并发上限。

4. 让等待逻辑可以响应取消

原来的重试延迟使用 time.Sleep,进程收到停止信号后可能还要等完整的重试间隔。修复后使用带 Context 的定时等待:

1
2
3
4
5
6
select {
case <-ctx.Done():
return false
case <-timer.C:
return true
}

这不会中断正在执行的 Provider 调用,但可以让空闲消费循环和错误重试循环及时退出。

验证与回归

并发回归测试

新增 consumer_test.go 中的 fake Stream Client,准备 3 条消息并把并发设置为 2。测试让前两个处理器阻塞,确认:

  • 两条不同消息可以同时进入处理器;
  • 达到并发上限时,第三条消息仍留在队列中,没有被提前领取;
  • 释放前两个任务后,第三条消息才开始处理;
  • 三条消息最终都完成 ACK;
  • 取消 Context 后消费循环可以退出。

这个测试直接验证“并发处理”和“并发上限”两个行为,不依赖真实 Redis,也不会产生 Provider 成本。

实际执行过的检查

server/ 目录执行:

1
2
3
4
go test ./... -count=1
go test -race ./internal/infrastructure/queue/redisstream ./internal/infrastructure/config -count=1
go vet ./internal/infrastructure/queue/... ./internal/infrastructure/config/...
go build ./cmd/worker

以上检查均通过。随后执行:

1
docker compose config --quiet

Compose 配置解析通过。开发环境 Worker 热重载后的日志包含:

1
2
queue initialized ... group=video-worker ... concurrency=4
worker starting

说明运行时已经读取了新的并发配置。

常见误区与边界

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 回收和上游已受理任务的重复执行边界。
投喂小莫
给快要饿死的小莫投喂点零食吧~
投喂小莫
分享
分享提示信息