scheduler

package
v0.7.0 Latest Latest
Warning

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

Go to latest
Published: Aug 10, 2026 License: Apache-2.0 Imports: 3 Imported by: 0

Documentation

Overview

Package scheduler 提供工作池式任务调度器:生产者把任务投递到队列, N 个 worker goroutine 并发消费;支持运行时 Pause/Resume 与优雅停止。

与 pkg/service/cron 的区别:cron 按 cron 表达式定时触发,scheduler 按 事件驱动(手动 Submit),且支持运行时暂停/恢复回调处理——适合"回调可能很重、 需要在业务高峰期暂停处理"的场景(如发奖、批量通知、过期清理)。

工作池 + Pause/Resume 模型。

零值不可用,用 New 构造。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type ErrorHandler

type ErrorHandler func(taskName string, err error, panicStack []byte)

ErrorHandler 处理 worker 执行 Task 时的 panic / 返回 error。

type Option

type Option func(*config)

Option 配置 Scheduler。

func WithErrorHandler

func WithErrorHandler(h ErrorHandler) Option

WithErrorHandler 设置错误处理回调。默认忽略 error,panic 会被 recover 不中断 worker。

func WithQueueSize

func WithQueueSize(n int) Option

WithQueueSize 设置任务队列容量,默认 256。队列满时 Submit 阻塞(默认)或返回 false。

func WithWorkers

func WithWorkers(n int) Option

WithWorkers 设置 worker 数量,默认 4。0 表示不启动 worker(纯排队模式, 仅用于测试或需要外部自行消费 queue 的场景)。

type Scheduler

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

Scheduler 管理工作池。

func New

func New(opts ...Option) *Scheduler

New 创建调度器(未启动)。用 Start 启动。

func (*Scheduler) Pause

func (s *Scheduler) Pause()

Pause 暂停所有 worker:正在执行的 Task 继续完成,之后 worker 阻塞不再取新任务。 已入队任务不丢失,Resume 后继续消费。幂等。

func (*Scheduler) Paused

func (s *Scheduler) Paused() bool

Paused 返回是否暂停。

func (*Scheduler) Pending

func (s *Scheduler) Pending() int

Pending 返回队列中待处理任务数(近似)。

func (*Scheduler) Resume

func (s *Scheduler) Resume()

Resume 恢复消费。幂等。

func (*Scheduler) Start

func (s *Scheduler) Start(ctx context.Context)

Start 启动 worker 池。幂等。ctx 取消时优雅停止(排空队列后退出)。

func (*Scheduler) Stop

func (s *Scheduler) Stop()

Stop 停止调度器:关闭队列,等所有 worker 排空并退出。幂等。 排空保证已 Submit 的任务被执行(除非 ctx 已取消)。

func (*Scheduler) Submit

func (s *Scheduler) Submit(t *Task) bool

Submit 投递一个任务。若已停止返回 false;队列满时阻塞(保证不丢任务)。 在 Pause 期间仍可 Submit(任务入队,Resume 后消费)。

func (*Scheduler) TrySubmit

func (s *Scheduler) TrySubmit(t *Task) bool

TrySubmit 非阻塞投递。队列满立即返回 false。

func (*Scheduler) Wait

func (s *Scheduler) Wait()

Wait 阻塞直到所有 worker 退出。

type Task

type Task struct {
	Name string
	Fn   func(ctx context.Context) error
}

Task 是被调度的工作单元。业务在 fn 中执行具体逻辑。

Jump to

Keyboard shortcuts

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