schedule

package
v0.0.5 Latest Latest
Warning

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

Go to latest
Published: Aug 2, 2026 License: MIT Imports: 13 Imported by: 0

README

schedule 包 — 定时任务调度

所属层级: Infrastructure Layer
设计理念: 零外部依赖,Spring 风格调度
设计灵感: Spring @Scheduled + Quartz

概述

schedule 包提供零外部依赖的 Spring 风格定时任务调度框架,为应用提供 Cron 表达式解析、任务调度和生命周期管理能力。

核心功能
功能 说明
Cron 表达式解析 6 字段 Cron 表达式,位图编码,Next() 算法
任务调度 基于最小堆的调度器,支持动态注册/注销
固定延迟任务 从上次任务完成后延迟固定时间再执行
固定频率任务 从任务开始时计算下一次执行时间
自动配置 ScheduleAutoConfigurationschedule.enabled=true 时自动启用
并发控制 可配置并发池大小,限制同时执行的任务数
优雅关闭 支持超时控制的优雅关闭
Builder 模式 链式配置,简化调度器创建
Helper 工具 简化常见调度器操作

快速开始

创建调度器
package main

import (
    "context"
    "fmt"
    "github.com/xudefa/enhance/schedule"
    "time"
)

func main() {
    // 创建调度器
    scheduler := schedule.NewScheduler(
        schedule.WithPoolSize(10),
        schedule.WithErrorHandler(func(taskName string, err error) {
            fmt.Printf("task %s failed: %v\n", taskName, err)
        }),
    )

    // 创建任务
    task := schedule.NewTask("my-task", "0 */5 * * * *", func(ctx context.Context) error {
        fmt.Println("task executed")
        return nil
    })

    // 注册任务
    if err := scheduler.Register(task); err != nil {
        panic(err)
    }

    // 启动调度器
    ctx := context.Background()
    if err := scheduler.Start(ctx); err != nil {
        panic(err)
    }

    // 运行一段时间后关闭
    // ...

    // 优雅关闭
    scheduler.Shutdown(ctx)
}
使用 Builder 模式
// 使用 Builder 模式创建调度器
scheduler := schedule.NewSchedulerBuilder().
    PoolSize(10).
    WithCronTask("cron-task", "0 */5 * * * *", func(ctx context.Context) error {
        fmt.Println("cron task executed")
        return nil
    }).
    WithFixedDelayTask("delay-task", 5*time.Second, func(ctx context.Context) error {
        fmt.Println("fixed delay task executed")
        return nil
    }).
    WithFixedRateTask("rate-task", 10*time.Second, func(ctx context.Context) error {
        fmt.Println("fixed rate task executed")
        return nil
    }).
    Build()

// 启动调度器
ctx := context.Background()
scheduler.Start(ctx)
使用 Helper 工具
// 创建调度器和 Helper
scheduler := schedule.NewScheduler()
helper := schedule.NewScheduleHelper(scheduler)

// 注册任务
helper.RegisterCronTask("my-cron", "0 */5 * * * *", func(ctx context.Context) error {
    return nil
})

helper.RegisterFixedDelayTask("my-delay", 5*time.Second, func(ctx context.Context) error {
    return nil
})

// 检查任务
if helper.HasTask("my-cron") {
    fmt.Println("task exists")
}

fmt.Printf("task count: %d\n", helper.GetTaskCount())

// 注销任务
helper.UnregisterTask("my-cron")
Cron 表达式格式

支持 6 字段 Spring 风格 Cron 表达式:秒 分 时 日 月 周

字段 范围 特殊字符
0-59 * , - /
0-59 * , - /
0-23 * , - /
1-31 * , - /
1-12 或 JAN-DEC * , - /
0-6 或 SUN-SAT * , - /
常用示例
表达式 说明
0 */5 * * * * 每 5 分钟执行
0 0 */1 * * * 每小时执行
0 0 0 * * * 每天零点执行
0 0 0 * * MON-FRI 工作日零点执行
0 0 0 1 * * 每月 1 号零点执行
0 30 9 * * MON-FRI 工作日 9:30 执行

API 参考

Task 接口
type Task interface {
    Name() string                        // 任务名称
    Cron() string                        // Cron 表达式
    Execute(ctx context.Context) error   // 执行任务
}
任务创建函数
// Cron 表达式任务
func NewTask(name, cron string, fn func(ctx context.Context) error) Task

// 固定延迟任务(从任务完成后开始计时)
func NewFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) Task

// 固定频率任务(从任务开始时开始计时)
func NewFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) Task
Scheduler 接口
type Scheduler interface {
    Start(ctx context.Context) error           // 启动调度器
    Shutdown(ctx context.Context) error        // 优雅关闭
    Register(task Task) error                  // 注册任务
    Unregister(name string) bool               // 注销任务
    IsRunning() bool                           // 是否运行中
    RegisteredTasks() []Task                   // 已注册任务列表
}
配置选项
选项 说明 默认值
WithPoolSize(size int) 任务执行池大小 10
WithErrorHandler(fn func(taskName string, err error)) 错误处理函数 nil
WithLogger(logger log.Logger) 日志记录器 log.NewLoggerBuilder().Build()
Builder 模式
type SchedulerBuilder struct{}

func NewSchedulerBuilder() *SchedulerBuilder
func (b *SchedulerBuilder) PoolSize(size int) *SchedulerBuilder
func (b *SchedulerBuilder) ErrorHandler(fn func(taskName string, err error)) *SchedulerBuilder
func (b *SchedulerBuilder) Logger(logger log.Logger) *SchedulerBuilder
func (b *SchedulerBuilder) WithTask(task Task) *SchedulerBuilder
func (b *SchedulerBuilder) WithCronTask(name, cron string, fn func(ctx context.Context) error) *SchedulerBuilder
func (b *SchedulerBuilder) WithFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) *SchedulerBuilder
func (b *SchedulerBuilder) WithFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) *SchedulerBuilder
func (b *SchedulerBuilder) Build() *DefaultScheduler
Helper 工具
type ScheduleHelper struct{}

func NewScheduleHelper(scheduler *DefaultScheduler) *ScheduleHelper
func (h *ScheduleHelper) RegisterCronTask(name, cron string, fn func(ctx context.Context) error) error
func (h *ScheduleHelper) RegisterFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) error
func (h *ScheduleHelper) RegisterFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) error
func (h *ScheduleHelper) UnregisterTask(name string) bool
func (h *ScheduleHelper) GetTaskCount() int
func (h *ScheduleHelper) HasTask(name string) bool
func (h *ScheduleHelper) StartAndBlock(ctx context.Context) error

架构设计

┌──────────────────────────────────────────────────────┐
│              DefaultScheduler                        │
│                                                      │
│  ┌──────────────┐  ┌──────────────────────────────┐  │
│  │  taskHeap    │  │      taskPool (semaphore)    │  │
│  │  (min-heap)  │  │      (concurrency control)   │  │
│  │              │  │                              │  │
│  │  Register()  │  │  executeTask()               │  │
│  │  Unregister()│  │  calculateNextRun()          │  │
│  └──────────────┘  └──────────────────────────────┘  │
│                                                      │
│  ┌──────────────┐  ┌──────────────────────────────┐  │
│  │ CronExpr     │  │      Task                    │  │
│  │ (bitmap)     │  │  (function wrapper)          │  │
│  │              │  │                              │  │
│  │  Next()      │  │  Execute()                   │  │
│  └──────────────┘  └──────────────────────────────┘  │
└──────────────────────────────────────────────────────┘

最佳实践

1. 合理设置并发池大小
// ✅ 推荐:根据任务特性设置池大小
scheduler := schedule.NewScheduler(
    schedule.WithPoolSize(20), // CPU 密集型:CPU 核心数
)

// ✅ 推荐:IO 密集型可以设置更大
scheduler := schedule.NewScheduler(
    schedule.WithPoolSize(50), // IO 密集型:CPU 核心数 * 2
)
2. 使用错误处理函数
// ✅ 推荐:注册错误处理函数
scheduler := schedule.NewScheduler(
    schedule.WithErrorHandler(func(taskName string, err error) {
        log.Printf("task %s failed: %v", taskName, err)
        // 可以发送告警、记录指标等
    }),
)
3. 优雅关闭
// ✅ 推荐:使用带超时的关闭
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()

if err := scheduler.Shutdown(ctx); err != nil {
    log.Printf("shutdown timeout: %v", err)
}
4. 任务幂等性
// ✅ 推荐:确保任务可重复执行
task := schedule.NewTask("cleanup-task", "0 0 0 * * *", func(ctx context.Context) error {
    // 确保多次执行结果一致
    return cleanupOldFiles(ctx)
})

Documentation

Overview

Package schedule 提供定时任务调度功能,用于 enhance 框架。

Package schedule 提供定时任务调度功能,用于 enhance 框架。

Package schedule 提供定时任务调度功能,用于 enhance 框架。

该模块提供灵活的定时任务调度机制,支持基于Cron 表达式、固定延迟/固定频率执行等。 参考 Spring 的 Scheduler 设计。

架构设计

  • Task: 任务接口,定义任务执行逻辑
  • Scheduler: 调度器接口,负责任务调度
  • SchedulerOption: 调度器选项函数
  • CronExpression: Cron 表达式解析器
  • DefaultScheduler: 默认调度器实现

核心功能

  • Cron 表达式: 支持标准 Cron 表达式定义执行时间
  • 固定延迟: 支持固定延迟执行任务
  • 固定频率: 支持固定频率执行任务
  • 并发控制: 支持任务并发执行控制
  • 错误处理: 支持任务执行失败后的重试

Cron 表达式

支持 6 位 Cron 表达式:秒 分 时 日 月 周

  • "0 */5 * * * *": 每 5 分钟执行
  • "0 0 */1 * * *": 每小时执行
  • "0 0 0 * * *": 每天零点执行
  • "0 0 0 * * MON-FRI": 工作日零点执行

配置选项

  • WithPoolSize: 设置任务执行池大小(最大并发数)
  • WithErrorHandler: 设置执行错误处理函数

Package schedule 提供定时任务调度功能,用于 enhance 框架。

Package schedule 提供定时任务调度功能,用于 enhance 框架。

Index

Constants

View Source
const (
	// ScheduleEnabled 调度器启用开关。
	ScheduleEnabled = "schedule.enabled"
	// SchedulePoolSize 调度器线程池大小。
	SchedulePoolSize = "schedule.pool-size"
	// ScheduleScanAnnotations 是否扫描注解注册任务。
	ScheduleScanAnnotations = "schedule.scan-annotations"
)

调度器配置键常量。

View Source
const (
	// DefaultSchedulePoolSize 默认线程池大小。
	DefaultSchedulePoolSize = 10
	// DefaultScheduleScanAnnotations 默认是否扫描注解。
	DefaultScheduleScanAnnotations = true
	// ConditionTrue 条件判断的真值常量。
	ConditionTrue = "true"
)

默认值常量。

Variables

This section is empty.

Functions

This section is empty.

Types

type CronExpression

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

CronExpression Cron 表达式解析器。

支持 6 字段 Spring 风格 Cron 表达式:秒 分 时 日 月 周 使用位图编码实现高效的下次执行时间计算。

func ParseCronExpression

func ParseCronExpression(expr string) (*CronExpression, error)

ParseCronExpression 解析 Cron 表达式字符串。

支持格式:秒 分 时 日 月 周 示例:

  • "0 */5 * * * *" : 每 5 分钟
  • "0 0 */1 * * *" : 每小时
  • "0 0 0 * * *" : 每天零点
  • "0 0 0 * * MON-FRI" : 工作日零点

func (*CronExpression) Next

func (ce *CronExpression) Next(from time.Time) time.Time

Next 计算给定时间之后的下次执行时间。

使用逐字段匹配算法,从当前时间开始向后搜索最近的匹配时间。

type DefaultScheduler

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

DefaultScheduler 默认调度器实现。

基于最小堆实现任务调度,支持动态注册/注销任务、并发控制和优雅关闭。

func NewScheduler

func NewScheduler(ctx context.Context, opts ...SchedulerOption) *DefaultScheduler

NewScheduler 创建调度器实例。

参数:

  • ctx: 父级 context,用于控制调度器生命周期
  • opts: 配置选项

func (*DefaultScheduler) Close added in v0.0.4

func (s *DefaultScheduler) Close()

Close 关闭调度器,释放资源。

此方法确保调度器的 context 被正确取消,防止 goroutine 泄漏。 如果调度器未启动或已关闭,调用此方法不会产生任何副作用。

func (*DefaultScheduler) IsRunning

func (s *DefaultScheduler) IsRunning() bool

IsRunning 返回调度器是否正在运行。

func (*DefaultScheduler) Register

func (s *DefaultScheduler) Register(task Task) error

Register 注册定时任务,任务名唯一,重复返回 error。

func (*DefaultScheduler) RegisteredTasks

func (s *DefaultScheduler) RegisteredTasks() []Task

RegisteredTasks 返回所有已注册任务。

func (*DefaultScheduler) Shutdown

func (s *DefaultScheduler) Shutdown(ctx context.Context) error

Shutdown 优雅关闭,等待正在执行的任务完成。

func (*DefaultScheduler) Start

func (s *DefaultScheduler) Start(ctx context.Context) error

Start 启动调度器,开始触发定时任务。

func (*DefaultScheduler) Unregister

func (s *DefaultScheduler) Unregister(name string) bool

Unregister 注销定时任务。

type ScheduleAutoConfiguration

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

ScheduleAutoConfiguration 调度器自动配置类。

当配置文件中启用调度器时自动生效(schedule.enabled=true)。 负责创建调度器实例并注册到 IoC 容器中。

func (*ScheduleAutoConfiguration) Configure

Configure 配置调度器。

该方法在自动配置阶段调用,负责:

  1. 从 Environment 中读取调度器配置
  2. 创建调度器实例
  3. 注册 Scheduler 到 IoC 容器

type ScheduleConfig

type ScheduleConfig struct {
	Enabled  bool `json:"enabled" mapstructure:"enabled"`
	PoolSize int  `json:"pool-size" mapstructure:"pool-size"`
}

ScheduleConfig 调度器配置。

type ScheduleHelper

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

ScheduleHelper 调度器辅助工具,简化常见操作。

func NewScheduleHelper

func NewScheduleHelper(scheduler *DefaultScheduler) *ScheduleHelper

NewScheduleHelper 创建调度器辅助工具。

func (*ScheduleHelper) GetTaskCount

func (h *ScheduleHelper) GetTaskCount() int

GetTaskCount 获取已注册任务数量。

func (*ScheduleHelper) HasTask

func (h *ScheduleHelper) HasTask(name string) bool

HasTask 检查任务是否已注册。

func (*ScheduleHelper) RegisterCronTask

func (h *ScheduleHelper) RegisterCronTask(name, cron string, fn func(ctx context.Context) error) error

RegisterCronTask 注册 Cron 表达式任务。

func (*ScheduleHelper) RegisterFixedDelayTask

func (h *ScheduleHelper) RegisterFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) error

RegisterFixedDelayTask 注册固定延迟任务。

func (*ScheduleHelper) RegisterFixedRateTask

func (h *ScheduleHelper) RegisterFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) error

RegisterFixedRateTask 注册固定频率任务。

func (*ScheduleHelper) StartAndBlock

func (h *ScheduleHelper) StartAndBlock(ctx context.Context) error

StartAndBlock 启动调度器并阻塞,直到收到停止信号。

func (*ScheduleHelper) UnregisterTask

func (h *ScheduleHelper) UnregisterTask(name string) bool

UnregisterTask 注销任务。

type ScheduleStarter

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

ScheduleStarter 调度器启动器。

管理调度器的生命周期,在应用启动时启动调度器, 在应用停止时优雅关闭调度器。

func (*ScheduleStarter) Configure

func (s *ScheduleStarter) Configure(ctx boot.ApplicationContext) error

Configure 配置阶段调用。

func (*ScheduleStarter) Dependencies

func (s *ScheduleStarter) Dependencies() []string

Dependencies 返回依赖的其他启动器名称。

func (*ScheduleStarter) GetCondition

func (s *ScheduleStarter) GetCondition() condition.Condition

GetCondition 返回启动条件。

func (*ScheduleStarter) Name

func (s *ScheduleStarter) Name() string

Name 返回启动器名称。

func (*ScheduleStarter) Start

Start 启动阶段调用。

func (*ScheduleStarter) Stop

Stop 停止阶段调用。

type Scheduler

type Scheduler interface {
	// Start 启动调度器,开始触发定时任务。
	Start(ctx context.Context) error
	// Shutdown 优雅关闭,等待正在执行的任务完成。
	Shutdown(ctx context.Context) error
	// Register 注册定时任务,任务名唯一,重复返回 error。
	Register(task Task) error
	// Unregister 注销定时任务。
	Unregister(name string) bool
	// IsRunning 返回调度器是否正在运行。
	IsRunning() bool
	// RegisteredTasks 返回所有已注册任务。
	RegisteredTasks() []Task
}

Scheduler 定时任务调度器接口。

type SchedulerBuilder

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

SchedulerBuilder 调度器构建器,支持链式配置。

func NewSchedulerBuilder

func NewSchedulerBuilder() *SchedulerBuilder

NewSchedulerBuilder 创建调度器构建器。

func (*SchedulerBuilder) Build

func (b *SchedulerBuilder) Build() *DefaultScheduler

Build 构建调度器。

func (*SchedulerBuilder) Context added in v0.0.4

Context 设置父级 context。

func (*SchedulerBuilder) ErrorHandler

func (b *SchedulerBuilder) ErrorHandler(fn func(taskName string, err error)) *SchedulerBuilder

ErrorHandler 设置错误处理函数。

func (*SchedulerBuilder) Logger

func (b *SchedulerBuilder) Logger(logger log.Logger) *SchedulerBuilder

Logger 设置日志记录器。

func (*SchedulerBuilder) MustBuild

func (b *SchedulerBuilder) MustBuild() *DefaultScheduler

MustBuild 构建调度器,失败则 panic。

func (*SchedulerBuilder) PoolSize

func (b *SchedulerBuilder) PoolSize(size int) *SchedulerBuilder

PoolSize 设置任务执行池大小。

func (*SchedulerBuilder) WithCronTask

func (b *SchedulerBuilder) WithCronTask(name, cron string, fn func(ctx context.Context) error) *SchedulerBuilder

WithCronTask 添加 Cron 表达式任务。

func (*SchedulerBuilder) WithFixedDelayTask

func (b *SchedulerBuilder) WithFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) *SchedulerBuilder

WithFixedDelayTask 添加固定延迟任务。

func (*SchedulerBuilder) WithFixedRateTask

func (b *SchedulerBuilder) WithFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) *SchedulerBuilder

WithFixedRateTask 添加固定频率任务。

func (*SchedulerBuilder) WithTask

func (b *SchedulerBuilder) WithTask(task Task) *SchedulerBuilder

WithTask 添加任务。

type SchedulerOption

type SchedulerOption func(*schedulerConfig)

SchedulerOption 调度器配置选项函数。

func WithErrorHandler

func WithErrorHandler(fn func(taskName string, err error)) SchedulerOption

WithErrorHandler 设置执行错误处理函数。

func WithLogger

func WithLogger(logger log.Logger) SchedulerOption

WithLogger 设置日志记录器。

func WithPoolSize

func WithPoolSize(size int) SchedulerOption

WithPoolSize 设置任务执行池大小(最大并发数)。

type Task

type Task interface {
	// Name 返回任务名称,用于注册和查找。
	Name() string
	// Cron 返回任务的 cron 表达式(6 字段 Spring 风格)。
	Cron() string
	// Execute 执行任务逻辑。
	Execute(ctx context.Context) error
}

Task 定时任务接口。

每个定时任务包含名称、Cron 表达式和执行逻辑。 通过 NewTask 创建基于函数的任务实例。

func NewFixedDelayTask

func NewFixedDelayTask(name string, delay time.Duration, fn func(ctx context.Context) error) Task

NewFixedDelayTask 创建固定延迟任务。

延迟时间从上次任务完成时开始计算。

func NewFixedRateTask

func NewFixedRateTask(name string, interval time.Duration, fn func(ctx context.Context) error) Task

NewFixedRateTask 创建固定频率任务。

从任务开始执行时计算下一次执行时间。

func NewTask

func NewTask(name, cron string, fn func(ctx context.Context) error) Task

NewTask 创建基于函数的任务实例。

Jump to

Keyboard shortcuts

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