Documentation
¶
Overview ¶
Package jobqueue 提供带优先级、进度上报、生命周期事件的任务队列。
借鉴 BullMQ 的架构分层(Queue / Worker / Events):
- Queue: 投递任务,支持优先级(数值越小优先级越高)、延迟投递;
- Worker: N 个 goroutine 并发消费,支持 Pause/Resume;
- EventHook: 生命周期事件回调(Submit/Start/Progress/Complete/Fail), 用于构建可观测性(日志、metrics、dashboard)。
与 pkg/orchestration/scheduler 的区别:
- scheduler 是 FIFO channel + 工作池,不支持优先级/进度/事件;
- jobqueue 用 priority 堆做排序,每个 Job 有状态机 + 进度回调。
与 pkg/orchestration/delayqueue 的区别:
- delayqueue 仅"到点触发回调",不关心执行状态与并发控制;
- jobqueue 关注"完整的 Job 生命周期":排队→执行→报告进度→完成/失败。
进程内实现(不持久化);分布式持久化版本见 contrib/redisqueue。 零值不可用,用 New 构造。并发安全。
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func ReportProgress ¶
ReportProgress 在 Job.Fn 内调用,上报执行进度(0~100)。 若 ctx 中无 reporter(非 jobqueue 执行),则静默忽略。
Types ¶
type Config ¶
type Config struct {
Workers int // worker 数量,默认 4
QueueSize int // 内部就绪信号缓冲,默认 1024
Hook EventHook // 事件钩子(nil=不回调)
OnPanic func(job *Job, r any, stack []byte)
}
Config 配置。
type EventHook ¶
type EventHook interface {
OnEvent(event Event)
}
EventHook 事件回调接口。实现者可选择性处理感兴趣的事件。 回调在 worker goroutine 中同步调用,应尽量轻量(如写 channel / 递增 metric)。
type EventHookFunc ¶
type EventHookFunc func(Event)
EventHookFunc 函数适配器。
func (EventHookFunc) OnEvent ¶
func (f EventHookFunc) OnEvent(e Event)
type Job ¶
type Job struct {
// ID 任务唯一标识。
ID string
// Name 任务名称(用于日志/metrics 分组)。
Name string
// Priority 优先级:数值越小越优先(0 最高)。默认 0。
Priority int
// Payload 任务数据(业务自定义)。
Payload any
// Fn 执行函数。ctx 携带 ProgressReporter,可通过 ReportProgress 上报进度。
Fn func(ctx context.Context, job *Job) error
// MaxRetries 最大重试次数(0=不重试)。
MaxRetries int
// RetryDelay 重试基础延迟(第 n 次 = delay * 2^n)。
RetryDelay time.Duration
// Delay 延迟执行:投递后等 Delay 再进入就绪队列。0=立即就绪。
Delay time.Duration
// Timeout 单次执行超时(0=不限)。
Timeout time.Duration
// State 当前状态。
State JobState
// Attempts 已尝试次数。
Attempts int
// Progress 当前进度(0~100)。
Progress float64
// Err 最近一次执行错误。
Err error
// CreatedAt 创建时间。
CreatedAt time.Time
// StartedAt 开始执行时间。
StartedAt time.Time
// CompletedAt 完成时间。
CompletedAt time.Time
// contains filtered or unexported fields
}
Job 是一个待执行的任务。
type Option ¶
type Option func(*Config)
Option 配置函数。
func WithPanicHandler ¶
WithPanicHandler 设置 panic 处理。
Click to show internal directories.
Click to hide internal directories.