pitayacron

package module
v0.1.0 Latest Latest
Warning

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

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

README

pitaya-cron

CI Go Reference

pitaya-cron 是 robfig/cron/v3 的轻量扩展,为 Pitaya 服务增加可选的 etcd Leader 选举。它同时支持每个实例独立执行的常规任务,以及整个集群仅由 Leader 执行的分布式任务。

非官方 Pitaya 社区扩展,与 Pitaya 官方项目没有隶属关系。

安装

go get github.com/chuhongliang/pitaya-cron@v0.1.0

项目结构

pitaya-cron/
├─ *.go              # Cron、任务、配置和可替换的 Elector 接口
├─ *_test.go         # 单元测试与条件集成测试
└─ internal/etcd/    # 默认 etcd 选主实现

业务项目只需导入 github.com/chuhongliang/pitaya-cron。etcd Client 由调用方创建并管理,组件不会关闭它。

执行模型

每个进程的 robfig/cron
├─ 常规任务       → 每个进程都执行
└─ 分布式任务     → 只有当前 etcd Leader 执行
  • 常规任务不依赖 etcd,适合清理本地缓存、刷新进程内状态等工作。
  • 分布式任务使用 etcd Lease 选主,适合排行榜结算、全局维护等集群唯一任务。
  • 失去 Leadership 时,正在执行的分布式 Handler 会收到 Context 取消信号。
  • 本组件不包含任务队列、持久化补偿、死信或跨进程重试。

Pitaya 接入

import (
    "context"

    pitayacron "github.com/chuhongliang/pitaya-cron"
)

cronModule, err := pitayacron.New(
    pitayacron.Config{
        Namespace:   "pitaya-starter-dev",
        NodeID:      "global-1",
        WithSeconds: true,
        Location:    "Asia/Shanghai",
        Election: pitayacron.ElectionConfig{
            TTL: 15,
        },
    },
    pitayacron.Dependencies{
        Etcd:   etcdClient,
        Logger: logger,
    },
)
if err != nil {
    return err
}

err = cronModule.AddLocalJob(pitayacron.Job{
    Name: "cleanup-local-cache",
    Spec: "*/10 * * * * *",
    Handler: func(ctx context.Context, task pitayacron.Task) error {
        return cleanupLocalCache(ctx)
    },
})
if err != nil {
    return err
}

err = cronModule.AddDistributedJob(pitayacron.Job{
    Name: "ranking-reset",
    Spec: "0 0 0 * * *",
    Handler: func(ctx context.Context, task pitayacron.Task) error {
        // RunID 可作为同一秒触发的幂等键,业务仍应自行保证幂等。
        return resetRanking(ctx, task.RunID)
    },
})
if err != nil {
    return err
}

return app.RegisterModule(cronModule, "pitayaCron")

所有参与同一个 Namespace 选举的实例必须注册相同的分布式任务。任务只能在 Init 前注册。

配置说明

  • Namespace:选举隔离域,同时用于生成默认 etcd 前缀。
  • NodeID:当前进程标识;留空时自动生成。
  • WithSeconds:开启六字段、秒级 Cron 表达式。
  • Location:Cron 时区,默认使用系统本地时区。
  • Standalone:本地开发模式;不连接 etcd,并在当前进程执行分布式任务。
  • Election.Prefix:默认 /pitaya-cron/<namespace>/election。
  • Election.TTL:etcd Lease TTL,默认 15 秒。

只注册常规任务时,可以不传 etcd:

cronModule, err := pitayacron.New(
    pitayacron.Config{Namespace: "local-maintenance"},
    pitayacron.Dependencies{},
)

本地调试分布式任务时:

pitayacron.Config{
    Namespace:  "local-dev",
    Standalone: true,
}

语义边界

etcd 保证同一选举域只有一个活跃 Leader,但不保存 Cron 任务。Leader 在触发期间宕机时,该次任务可能中断或错过;切换边界也不能提供严格一次执行。重要 Handler 必须幂等、响应 Context 取消,并按业务需要自行记录执行结果。

ScheduledAt 和 RunID 基于 Handler 实际触发时间生成。参与同一选举域的服务器应保持时钟同步;它们不是持久化调度序号。

robfig/cron 默认允许同一任务的多次触发重叠执行。执行时间可能超过调度间隔的任务,应在 Handler 内加锁或自行跳过重复运行。

开发与测试

gofmt -w .
go test -count=1 ./...
go vet ./...
go build ./...

本地 etcd 已启动时,可以验证真实选主和故障接管:

$env:PITAYA_CRON_INTEGRATION = "1"
go test -count=1 -run "TestEtcdFailoverIntegration|TestDistributedCronIntegration" -v .

Documentation

Overview

Package pitayacron extends robfig/cron with optional etcd leader election for Pitaya services.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func DefaultNodeID

func DefaultNodeID() string

DefaultNodeID builds a readable process-unique identity for election and local execution metadata.

func NewRunID

func NewRunID(jobName string, scheduledAt time.Time) string

NewRunID returns the cluster-stable execution id for a distributed job.

Types

type Config

type Config struct {
	Namespace       string
	NodeID          string
	WithSeconds     bool
	Location        string
	Standalone      bool
	Election        ElectionConfig
	RetryInterval   time.Duration
	ShutdownTimeout time.Duration
}

Config controls cron parsing, election retry, and shutdown behavior.

type Cron

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

Cron runs ordinary jobs on every process and distributed jobs only on the current etcd leader. Cron implements Pitaya's interfaces.Module method set.

func New

func New(config Config, dependencies Dependencies) (*Cron, error)

New creates a Cron using the built-in etcd election driver when an etcd dependency is supplied. The client remains owned by the caller.

func NewWithElector

func NewWithElector(config Config, elector Elector, logger *slog.Logger) (*Cron, error)

NewWithElector creates a Cron with a custom election driver. It is mainly useful for tests or non-etcd election implementations.

func (*Cron) AddDistributedJob

func (c *Cron) AddDistributedJob(job Job) error

AddDistributedJob registers a job that executes only on the current etcd leader. Every election participant must register the same distributed jobs.

func (*Cron) AddLocalJob

func (c *Cron) AddLocalJob(job Job) error

AddLocalJob registers a job that executes independently on every process. Jobs must be registered before Init.

func (*Cron) AfterInit

func (c *Cron) AfterInit()

AfterInit implements Pitaya's Module API.

func (*Cron) BeforeShutdown

func (c *Cron) BeforeShutdown()

BeforeShutdown cancels local handlers and the active election term.

func (*Cron) Epoch

func (c *Cron) Epoch() int64

Epoch returns the current etcd election revision, or zero when not leader.

func (*Cron) Init

func (c *Cron) Init() error

Init starts cron scheduling and, when necessary, etcd election.

func (*Cron) IsLeader

func (c *Cron) IsLeader() bool

IsLeader reports whether this process currently owns distributed leadership.

func (*Cron) Shutdown

func (c *Cron) Shutdown() error

Shutdown stops future triggers and waits for handlers and election work.

type Dependencies

type Dependencies struct {
	Etcd   *clientv3.Client
	Logger *slog.Logger
}

Dependencies are caller-owned resources. Cron never closes the etcd client.

type ElectionConfig

type ElectionConfig struct {
	Prefix          string
	TTL             int
	ResignTimeout   time.Duration
	CallbackTimeout time.Duration
}

ElectionConfig controls the built-in etcd election driver.

type Elector

type Elector interface {
	Run(ctx context.Context, nodeID string, leader LeaderFunc) error
}

Elector serializes distributed cron execution across processes.

type Handler

type Handler func(ctx context.Context, task Task) error

Handler executes a cron task. Distributed handlers should be idempotent and stop promptly when ctx is canceled after leadership loss.

type Job

type Job struct {
	Name    string
	Spec    string
	Handler Handler
}

Job describes a statically registered cron job.

type JobType

type JobType string

JobType describes where a cron job executes.

const (
	// JobTypeLocal executes independently on every process.
	JobTypeLocal JobType = "local"
	// JobTypeDistributed executes only on the current etcd leader.
	JobTypeDistributed JobType = "distributed"
)

type LeaderFunc

type LeaderFunc = func(ctx context.Context, epoch int64) error

LeaderFunc runs while the caller owns distributed leadership. Epoch is a monotonically increasing fencing token supplied by the election backend.

type Task

type Task struct {
	RunID       string
	Name        string
	Type        JobType
	NodeID      string
	ScheduledAt time.Time
	Epoch       int64
}

Task contains metadata for one direct handler invocation.

Directories

Path Synopsis
internal
etcd
Package etcd implements pitaya-cron's leader election with etcd leases.
Package etcd implements pitaya-cron's leader election with etcd leases.

Jump to

Keyboard shortcuts

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