不能直接用redis.Client做任务队列,因其缺乏原子性「取-处理-确认」流程,易因网络超时、进程崩溃或未ACK退出导致任务丢失;需组合LPUSH+RPOPLPUSH+EXPIRE+LREM或升级至Redis 6.2+使用Streams(XADD/XREADGROUP)保障可靠性。

为什么不能直接用 redis.Client 做任务队列
直接用 redis.Client 调用 LPUSH/BRPOP 容易丢任务:网络超时、进程崩溃、消费者没 ACK 就退出,都没法保证至少一次投递。Redis 本身不提供原子性「取-处理-确认」流程,必须靠客户端补全语义。
真正可用的方案是组合使用 Redis 的几种能力:LPUSH + RPOPLPUSH + EXPIRE + LREM,或者直接上 Redis Streams(6.2+)。前者兼容老版本但逻辑重,后者原生支持 consumer group 和 pending list,推荐优先考虑。
- 若 Redis 版本 ≥ 6.2,用
XADD/XREADGROUP,天然支持失败重试与多消费者负载均衡 - 若版本 RPOPLPUSH 模拟「正在处理队列」,配合定时任务扫
EXPIRE过期的 stuck job - 避免用
BLPOP+ 单独 ACK —— 中间任何一环失败都会导致任务丢失
用 go-redis/v9 接入 Streams 的最小可行代码
go-redis/v9 对 Streams 支持完整,但要注意 group 创建和消息确认的顺序。很多开发者在没创建 consumer group 的情况下直接 XREADGROUP,会报错 NOGROUP No such key 'xxx' or consumer group 'xxx'。
正确流程是:先 XGROUP CREATE,再 XREADGROUP,处理完必须调 XACK,否则消息永远留在 pending list 里。
// 初始化 client 后
ctx := context.Background()
// 创建 group(仅需执行一次,生产环境建议幂等检查)
_, err := rdb.XGroupCreate(ctx, &redis.XGroupCreateArgs{
Key: "task_stream",
Group: "worker_group",
ID: "$", // 从最新消息开始
}).Result()
if err != nil && !strings.Contains(err.Error(), "BUSYGROUP") {
log.Fatal(err)
}
<p>// 消费
streamMsgs, err := rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: "worker_group",
Consumer: "consumer_1",
Streams: []string{"task<em>stream", ">"},
Count: 1,
Block: 0,
}).Result()
if err != nil {
log.Printf("read failed: %v", err)
return
}
for </em>, msg := range streamMsgs[0].Messages {
// 处理业务逻辑
handleTask(msg.Values)</p><pre class="brush:php;toolbar:false;">// 必须确认,否则下次还会读到
rdb.XAck(ctx, "task_stream", "worker_group", msg.ID)}
如何防止任务重复消费或堆积
Streams 的 pending list 不会自动清理,如果 worker crash 且没发 XACK,消息会一直卡住。Redis 提供 XPENDING + XCLAIM 来接管超时任务,但需要你自己实现「reclaim logic」。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 设置 consumer 的
Idle时间(比如 30s),超过该时间未确认的任务视为失败 - 定期跑一个 goroutine 执行
XPENDING task_stream worker_group - + 10,拿到 pending 消息 ID 列表 - 对每个 pending ID 调
XCLAIM,指定新 consumer 名并重置 idle 计时器 - 注意
XCLAIM返回的是消息内容,不是布尔值 —— 它实际完成了「转移 + 重置」两件事
别依赖 DEL 或 LREM 清理旧消息:Streams 是 append-only,过期靠 XTRIM 控制长度,例如 rdb.XTrim(ctx, "task_stream", &redis.XTrimArgs{MaxLen: 10000})。
并发模型下怎么控制 worker 数量和背压
goroutine 开太多会把 Redis 连接打爆,开太少又吞不下流量。关键不是「起多少 goroutine」,而是「每条连接能扛多少并发命令」。
- 用
redis.NewClusterClient或连接池配置PoolSize(默认10,高并发建议设为 50–100) - 每个 worker goroutine 应复用同一个
*redis.Client,不要每次新建 - 避免在 handler 里同步调
time.Sleep或阻塞 IO —— 用context.WithTimeout包裹外部调用,超时就XACK并记录 error - 如果下游服务响应慢,靠
XREADGROUP的Count参数限流(比如设为 5),比在 Go 层做 semaphore 更准
最常被忽略的一点:Streams 的 group 和 consumer 是逻辑概念,不绑定连接或进程。同一个 consumer name 在多个实例里同时运行,会导致消息被重复分配 —— 所以 Consumer 名必须带唯一标识(如 hostname + pid)。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










