Go操作Redis Stream查询Pending消息为空及定时巡检实现问题
Redis Stream 邮件队列问题排查与实现
场景说明
- 任务需求:实现队列存储待发送给用户的邮件数据
- 预期逻辑:创建队列时同步创建绑定的消费组,队列内置worker协程,通过ticker定时检查Pending状态消息,将这类消息重新投递到队列实现重发
- 异常现象:
- 调用
CheckPending方法查询Pending状态消息时,始终返回空数组[] - 需要实现
CheckPending方法在后台按每分钟一次的频率运行,巡检Pending状态数据并将其重新加入消息列表
- 调用
- 补充背景:Redis实例中存有3天前写入的历史数据,仅对这些数据执行过
XACK命令完成消息确认,未做其他操作。
附现有代码
Queue.go
package main import ( "context" "errors" "github.com/go-redis/redis/v8" "github.com/rs/zerolog" "strings" "sync" "time" ) // ErrNoGroup 未创建消费组时请求数据返回该错误 var ErrNoGroup = errors.New("no group has been created") // ErrNoStream 未创建Stream时请求数据返回该错误 var ErrNoStream = errors.New("") type Options struct { Name string Redis *redis.Client Logger *zerolog.Logger } type Queue struct { Client *redis.Client Name string Group string WG sync.WaitGroup Logger *zerolog.Logger } const interval = 500 func New(options *Options) *Queue { logger := options.Logger.With().Str("service", "queue").Logger() q := &Queue{ Client: options.Redis, Name: options.Name + "stream", Group: options.Name + "group", Logger: &logger, } // 创建消费组 err := q.Client.XGroupCreateMkStream(context.Background(), q.Name, q.Group, "0").Err() if err != nil { return nil } ctx, cancel := context.WithCancel(context.Background()) q.WG.Add(1) go func() { time.Sleep(60 * interval * time.Millisecond) defer q.WG.Done() cancel() }() errP := q.CheckPending(ctx) if errP != nil { return nil } return q } // CheckPending 存在逻辑问题,无法查询到数据 func (q *Queue) CheckPending(ctx context.Context) error { l := q.methodLogger(ctx, "Queue CheckPending") l.Debug().Msg("check pending task") ticker := time.NewTicker(interval * time.Millisecond) for { select { case <-ticker.C: pending, err := q.Client.XPendingExt(ctx, &redis.XPendingExtArgs{ Stream: q.Name, Group: q.Group, Start: "-", End: "+", Count: 10, }).Result() if err != nil { if strings.HasPrefix(err.Error(), "NOGROUP") { return ErrNoGroup } if strings.HasPrefix(err.Error(), "NOSTREAM") { return ErrNoStream } return err } if len(pending) == 0 { return nil } fmt.Print(pending) // 实际返回[] break case <-ctx.Done(): return nil } } }
Main.go
func main() { rds := redis.NewClient(&redis.Options{ Addr: ":6379", }) q := queue.New(rds, "email") data := map[string]interface{}{"email": "email@gmail.com", "message": "We have received you order and we are working on it."} err := q.Add(context.Background(), data) fmt.Print(err) q.Read(context.TODO()) // 调用示例 err := q.CheckPending(context.TODO()) // 返回[] if err != nil { fmt.Print(err) } }
问题原因
- Pending列表查询为空符合Redis设计逻辑:Redis Stream的Pending列表仅存储被消费组消费者读取、但未执行XACK确认的消息。已执行XACK确认的历史消息会直接从Pending列表移除,自然查询不到。
- 代码逻辑错误导致巡检提前退出:
CheckPending方法中,第一次查询到Pending列表为空就直接return nil退出整个循环,无法实现持续巡检。New函数中创建的context会在30秒(60*500ms)后主动触发cancel,就算循环逻辑正常也会被强制终止。New函数中调用CheckPending的时机早于消息写入、消息读取操作,此时本来就不存在未确认的Pending消息,返回空是正常现象。
- 定时配置不符合需求:现有interval配置为500ms,和要求的每分钟巡检频率不符。
修复方案
1. 修正Pending巡检逻辑
- 巡检间隔调整为1分钟,查询Pending消息时增加Idle时间过滤,避免认领正在正常处理的消息。
- 移除“查询为空直接return”的逻辑,保持循环运行直到context被取消。
- 查到Pending消息后,先通过
XCLAIM将消息认领至巡检专用消费者,再执行重发/处理逻辑,处理完成后执行XACK避免消息重复Pending。 - 消费组创建时增加
BUSYGROUP错误处理,避免消费组已存在时初始化失败。
核心修正代码参考:
const pendingCheckInterval = 1 * time.Minute const claimIdleThreshold = 2 * time.Minute // 认领超过2分钟未确认的消息 // StartPendingInspector 启动后台Pending巡检协程,ctx与服务生命周期绑定 func (q *Queue) StartPendingInspector(ctx context.Context) { q.WG.Add(1) go func() { defer q.WG.Done() ticker := time.NewTicker(pendingCheckInterval) defer ticker.Stop() l := q.methodLogger(ctx, "Queue pending inspector") for { select { case <-ctx.Done(): l.Info().Msg("pending inspector exit") return case <-ticker.C: l.Debug().Msg("start scanning pending messages") pending, err := q.Client.XPendingExt(ctx, &redis.XPendingExtArgs{ Stream: q.Name, Group: q.Group, Start: "-", End: "+", Count: 10, Idle: claimIdleThreshold, }).Result() if err != nil { l.Error().Err(err).Msg("query pending failed") continue } if len(pending) == 0 { l.Debug().Msg("no pending messages found") continue } // 提取待认领消息ID msgIDs := make([]string, 0, len(pending)) for _, p := range pending { msgIDs = append(msgIDs, p.ID) } // 认领超时消息到巡检消费者 msgs, err := q.Client.XClaim(ctx, &redis.XClaimArgs{ Stream: q.Name, Group: q.Group, Consumer: "pending-inspector", MinIdle: claimIdleThreshold, Messages: msgIDs, }).Result() if err != nil { l.Error().Err(err).Msg("claim pending messages failed") continue } // 执行重发逻辑,可根据重试次数判断是否移入死信队列 for _, msg := range msgs { l.Info().Interface("msg", msg).Msg("reprocessing pending message") // 消息处理完成后执行ACK _ = q.Client.XAck(ctx, q.Name, q.Group, msg.ID).Err() } } } }() }
2. 调整初始化逻辑
移除New函数中自动启动30秒巡检的逻辑,初始化仅完成Stream和消费组创建,巡检协程、消费协程通过单独的Start方法启动,生命周期context由业务侧传入,服务退出时调用cancel触发协程退出,再通过WaitGroup等待所有协程退出即可。
内容的提问来源于stack exchange,提问作者alex
相关产品推荐
相关产品推荐

