redqueue

package
v1.74.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 19, 2026 License: MIT Imports: 6 Imported by: 0

Documentation

Overview

Package redqueue 提供基于 Redis 的可靠任务队列: 即时任务走 List,延迟任务走 ZSet(score=执行时间),多实例可并行消费。 消息带重试计数:handler 失败自动重投(延迟 1s),超过上限进入死信队列, 并通过死信回调(OnDeadLetter)通知运维(可接 alert/notify)。 与进程内 taskqueue 互补:进程内队列重启丢失,redqueue 持久化且支持多实例。

Index

Constants

View Source
const DefaultMaxRetries = 5

默认重试上限。

Variables

This section is empty.

Functions

This section is empty.

Types

type DeadLetter added in v1.71.0

type DeadLetter struct {
	Payload  []byte    `json:"payload"`   // 业务负载
	Retries  int       `json:"retries"`   // 进入死信时的重试次数
	FailedAt time.Time `json:"failed_at"` // 进入死信时间
}

DeadLetter 是死信消息(超过重试上限)。

type DeadLetterHook added in v1.73.0

type DeadLetterHook func(ctx context.Context, letter DeadLetter)

DeadLetterHook 死信回调:死信产生时调用(可接 alert/notify/日志)。 注意:多实例部署时每个实例都会收到回调,业务侧需自行去重 (如以 failed_at+payload 指纹做 Redis SETNX,或接受重复告警)。

type Queue

type Queue struct {
	// contains filtered or unexported fields
}

Queue 是基于 Redis 的可靠任务队列。

func NewQueue

func NewQueue(client *redis.Client, keyPrefix string) *Queue

NewQueue 创建队列;client 为 go-redis 客户端(可用 cache.RedisCache.Client())。 keyPrefix 用于多队列隔离(如 "order-tasks")。

func (*Queue) Consume

func (q *Queue) Consume(ctx context.Context, handler func(ctx context.Context, payload []byte) error) error

Consume 阻塞消费任务:延迟任务到期后原子搬入即时队列,再阻塞取任务执行。 handler 返回错误时任务重试计数 +1 并延迟 1s 重投;超过上限进入死信队列 并触发死信回调(如有设置)。 ctx 取消时优雅退出(处理中的任务完成后返回)。

func (*Queue) DeadLetterCount added in v1.71.0

func (q *Queue) DeadLetterCount(ctx context.Context) (int64, error)

DeadLetterCount 返回死信数量。

func (*Queue) DeadLetters added in v1.71.0

func (q *Queue) DeadLetters(ctx context.Context, offset, count int64) ([]DeadLetter, error)

DeadLetters 查询死信列表(倒序,offset/count 分页);损坏条目跳过。

func (*Queue) MaxRetries added in v1.71.0

func (q *Queue) MaxRetries() int

MaxRetries 返回当前重试上限。

func (*Queue) Pending

func (q *Queue) Pending(ctx context.Context) (int64, error)

Pending 返回待处理任务总数(即时队列 + 延迟队列)。

func (*Queue) RequeueDeadLetter added in v1.71.0

func (q *Queue) RequeueDeadLetter(ctx context.Context, index int64) error

RequeueDeadLetter 把指定位置的死信重新投递到即时队列并移除。 index 为死信列表索引(0 = 最新一条);重投时保留原重试计数。

func (*Queue) Submit

func (q *Queue) Submit(ctx context.Context, payload []byte, delay time.Duration) error

Submit 提交任务:delay<=0 走即时队列,否则走延迟队列(执行时间 = now+delay)。

func (*Queue) WithDeadLetterHook added in v1.73.0

func (q *Queue) WithDeadLetterHook(hook DeadLetterHook) *Queue

WithDeadLetterHook 设置死信回调(死信产生时调用,不阻塞消费)。

func (*Queue) WithMaxRetries added in v1.71.0

func (q *Queue) WithMaxRetries(maxRetries int) *Queue

WithMaxRetries 设置重试上限(0 表示无限重试;超过上限进死信队列)。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL