You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go操作Redis Stream查询Pending消息为空及定时巡检实现问题

Redis Stream 邮件队列问题排查与实现

场景说明

  • 任务需求:实现队列存储待发送给用户的邮件数据
  • 预期逻辑:创建队列时同步创建绑定的消费组,队列内置worker协程,通过ticker定时检查Pending状态消息,将这类消息重新投递到队列实现重发
  • 异常现象:
    1. 调用CheckPending方法查询Pending状态消息时,始终返回空数组[]
    2. 需要实现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)
    }
}

问题原因

  1. Pending列表查询为空符合Redis设计逻辑:Redis Stream的Pending列表仅存储被消费组消费者读取、但未执行XACK确认的消息。已执行XACK确认的历史消息会直接从Pending列表移除,自然查询不到。
  2. 代码逻辑错误导致巡检提前退出:
    • CheckPending方法中,第一次查询到Pending列表为空就直接return nil退出整个循环,无法实现持续巡检。
    • New函数中创建的context会在30秒(60*500ms)后主动触发cancel,就算循环逻辑正常也会被强制终止。
    • New函数中调用CheckPending的时机早于消息写入、消息读取操作,此时本来就不存在未确认的Pending消息,返回空是正常现象。
  3. 定时配置不符合需求:现有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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 10:18:19