用 Go 构建分布式任务调度系统
使用 Go、Machinery、Redis 和 cron 设计并交付一套生产级任务调度器。涵盖分布式锁与重试策略。
每个后端服务最终都会需要定时任务。发送周报、清理过期会话、从外部 API 同步数据、重试失败的支付。当你只运行单个实例时,一个简单的 cron 任务就够用了。但一旦扩展到多个实例,问题就来了:任务被重复执行、无法观测正在运行的内容、失败信息悄无声息地消失。
本文将完整讲解如何用 Go 设计并构建一套生产级的分布式任务调度系统。我们会从架构讲起,用可运行的代码实现核心组件,并覆盖那些教程通常会跳过的生产级问题——分布式锁、重试策略、优雅停机和监控。
为什么不直接用 Cron?
在单台服务器上,一个基础的 cron 任务(或 Go 的 time.Ticker)工作得很好。问题出现在你部署多个实例时:
- 重复执行:所有实例在同一时刻触发同一个任务
- 无可观测性:你无法知道哪些任务正在运行、哪些失败了、上次成功是什么时候
- 无重试逻辑:如果任务失败,它就消失了,直到下一个调度周期
- 强耦合:任务逻辑内嵌在负责调度的服务里,难以独立扩展工作节点
分布式任务调度系统通过将调度(何时运行什么)与执行(运行实际逻辑)分离,并在两者之间加入一个消息队列,解决了上述所有问题。
架构概览
系统由两个主要组件构成,通过一个消息队列相连:
任务中心(Task Center)——系统的大脑:
- 调度器(Scheduler):运行 cron 触发器,配合分布式锁确保每个任务恰好触发一次
- 任务管理器(Task Manager):用于管理任务定义的 CRUD API
- 监控器(Monitor):跟踪任务状态、采集指标、在失败时发送告警
- 结果处理器(Result Handler):从队列消费执行结果并更新任务状态
工作节点(Workers)——无状态、可水平扩展的执行单元:
- 任务接收器(Task Receiver):从队列拉取任务并分发
- 任务执行器(Task Executor):运行实际的任务逻辑
- 结果上报器(Result Reporter):将执行结果推回队列
- 重试处理器(Retry Handler):为失败任务实现指数退避
任务生命周期
每个任务都遵循一个定义清晰的状态机:
任务从 Created(已创建) 状态开始,当调度器将其推送到消息中间件时进入 Queued(已入队),被工作节点取走时进入 Running(运行中),最终以 Completed(已完成) 或 Failed(已失败) 结束。仍有剩余重试次数的失败任务进入 Retrying(重试中) 状态,并以指数退避重新入队。耗尽所有重试次数的任务进入 Dead Letter(死信) 队列,以供人工排查。
技术选型
Machinery——基石
Machinery 是一个用于分布式任务处理的 Go 框架,灵感来自 Python 的 Celery。它开箱即用地提供了任务注册、消息中间件抽象、结果存储和重试逻辑。
不过,单靠 machinery 并不能处理 cron 调度或分布式锁——这些层需要我们自己添加。
框架对比
在确定使用 Machinery 之前,我们先对比一下 Go 生态中的主要选项:
| 特性 | Machinery | Asynq | Temporal | go-cron |
|---|---|---|---|---|
| 任务队列 | 支持(Redis、AMQP、SQS) | 支持(仅 Redis) | 支持(内置) | 不支持 |
| Cron 调度 | 不支持(需手动添加) | 支持(内置) | 支持(内置) | 支持(核心功能) |
| 分布式锁 | 不支持 | 通过 Redis 实现唯一任务 | 内置 | 不支持 |
| 结果存储 | 支持(Redis、Memcache、MongoDB) | 支持(Redis) | 支持(内置) | 不支持 |
| 带退避的重试 | 支持 | 支持 | 支持(策略复杂) | 不支持 |
| 工作流/编排 | 支持(groups、chains、chords) | 不支持 | 支持(完整工作流引擎) | 不支持 |
| 复杂度 | 中等 | 低 | 高 | 极低 |
| 最适合场景 | 通用任务队列 | 简单异步任务 | 复杂工作流 | 简单的 cron 替代方案 |
当你需要一个具备中间件灵活性和工作流支持(chains、groups)的通用任务队列时,Machinery 是正确的选择。如果你已经确定使用 Redis,Asynq 更简单。Temporal 是应对复杂、长时间运行、带 saga 模式工作流的重量级选项。go-cron 适合单实例调度,但它本身无法解决分布式问题。
在本文中,我们将使用 Machinery 作为任务队列,并在其之上添加带分布式锁的 cron 调度。
实现
项目结构
scheduler/
├── cmd/
│ ├── server/main.go # Task center entry point
│ └── worker/main.go # Worker entry point
├── internal/
│ ├── config/config.go # Configuration
│ ├── lock/redis.go # Distributed lock
│ ├── scheduler/cron.go # Cron scheduler
│ ├── tasks/registry.go # Task definitions
│ ├── tasks/handlers.go # Task handler implementations
│ └── monitor/metrics.go # Monitoring and metrics
├── go.mod
└── go.sum
配置
// internal/config/config.go
package config
import "time"
type Config struct {
RedisURL string `env:"REDIS_URL" default:"redis://localhost:6379"`
BrokerURL string `env:"BROKER_URL" default:"redis://localhost:6379"`
ResultBackend string `env:"RESULT_BACKEND" default:"redis://localhost:6379"`
LockTTL time.Duration `env:"LOCK_TTL" default:"30s"`
MaxRetries int `env:"MAX_RETRIES" default:"3"`
DefaultQueue string `env:"DEFAULT_QUEUE" default:"machinery_tasks"`
MetricsPort int `env:"METRICS_PORT" default:"9090"`
}
任务定义与注册
首先,定义工作节点可以执行的任务。每个任务都是一个普通的 Go 函数——Machinery 负责处理序列化和分发。
// internal/tasks/handlers.go
package tasks
import (
"context"
"fmt"
"net/http"
"time"
"github.com/RichardKnop/machinery/v2/tasks"
)
// SendReport generates and emails a weekly report.
// Parameters are passed as primitive types (Machinery serialization requirement).
func SendReport(reportType string, recipientEmail string) (string, error) {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancel()
report, err := generateReport(ctx, reportType)
if err != nil {
return "", fmt.Errorf("generate report %s: %w", reportType, err)
}
if err := emailReport(ctx, recipientEmail, report); err != nil {
return "", fmt.Errorf("email report to %s: %w", recipientEmail, err)
}
return fmt.Sprintf("report_%s_sent_to_%s", reportType, recipientEmail), nil
}
// CleanExpiredSessions removes sessions older than the given threshold.
func CleanExpiredSessions(maxAgeDays int64) (int64, error) {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
count, err := deleteExpiredSessions(ctx, time.Duration(maxAgeDays)*24*time.Hour)
if err != nil {
return 0, fmt.Errorf("clean sessions older than %d days: %w", maxAgeDays, err)
}
return count, nil
}
// SyncExternalData pulls data from an external API and upserts it locally.
func SyncExternalData(apiEndpoint string) (string, error) {
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
req, err := http.NewRequestWithContext(ctx, http.MethodGet, apiEndpoint, nil)
if err != nil {
return "", fmt.Errorf("create request for %s: %w", apiEndpoint, err)
}
resp, err := http.DefaultClient.Do(req)
if err != nil {
return "", fmt.Errorf("fetch %s: %w", apiEndpoint, err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("unexpected status %d from %s", resp.StatusCode, apiEndpoint)
}
count, err := upsertData(ctx, resp.Body)
if err != nil {
return "", fmt.Errorf("upsert data from %s: %w", apiEndpoint, err)
}
return fmt.Sprintf("synced_%d_records", count), nil
}
现在,把这些任务注册到 Machinery:
// internal/tasks/registry.go
package tasks
import (
"github.com/RichardKnop/machinery/v2"
)
// RegisterAll registers all available task handlers with the Machinery server.
func RegisterAll(server *machinery.Server) error {
return server.RegisterTasks(map[string]interface{}{
"send_report": SendReport,
"clean_expired_sessions": CleanExpiredSessions,
"sync_external_data": SyncExternalData,
})
}
使用 Redis 实现分布式锁
cron 调度器会在任务中心的每个实例上运行。如果没有分布式锁,每个实例都会同时触发同一个任务。我们使用 Redis 的 SET NX(不存在则设置)搭配过期时间,确保恰好一个实例抢到锁。
// internal/lock/redis.go
package lock
import (
"context"
"crypto/rand"
"encoding/hex"
"errors"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
var ErrLockNotAcquired = errors.New("lock not acquired")
// RedisLock implements a distributed lock using Redis SET NX with automatic expiry.
type RedisLock struct {
client *redis.Client
key string
value string // unique value to prevent releasing someone else's lock
ttl time.Duration
}
// NewRedisLock creates a new distributed lock.
func NewRedisLock(client *redis.Client, key string, ttl time.Duration) *RedisLock {
// Generate a random value so only the holder can release the lock
b := make([]byte, 16)
rand.Read(b)
return &RedisLock{
client: client,
key: fmt.Sprintf("dlock:%s", key),
value: hex.EncodeToString(b),
ttl: ttl,
}
}
// Acquire attempts to acquire the lock. Returns ErrLockNotAcquired if
// another instance already holds it.
func (l *RedisLock) Acquire(ctx context.Context) error {
ok, err := l.client.SetNX(ctx, l.key, l.value, l.ttl).Result()
if err != nil {
return fmt.Errorf("redis SETNX: %w", err)
}
if !ok {
return ErrLockNotAcquired
}
return nil
}
// Release releases the lock, but only if we still own it.
// Uses a Lua script for atomic check-and-delete.
func (l *RedisLock) Release(ctx context.Context) error {
script := redis.NewScript(`
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("DEL", KEYS[1])
end
return 0
`)
_, err := script.Run(ctx, l.client, []string{l.key}, l.value).Result()
if err != nil {
return fmt.Errorf("release lock %s: %w", l.key, err)
}
return nil
}
// Extend resets the lock's TTL. Useful for long-running tasks that need
// to hold the lock beyond the initial TTL.
func (l *RedisLock) Extend(ctx context.Context, ttl time.Duration) error {
script := redis.NewScript(`
if redis.call("GET", KEYS[1]) == ARGV[1] then
return redis.call("PEXPIRE", KEYS[1], ARGV[2])
end
return 0
`)
_, err := script.Run(ctx, l.client, []string{l.key},
l.value, ttl.Milliseconds()).Result()
if err != nil {
return fmt.Errorf("extend lock %s: %w", l.key, err)
}
return nil
}
为什么 Release 要用 Lua 脚本? 如果不用,就存在竞态条件:实例 A 检查值发现匹配,但在它删除 key 之前,锁过期了,实例 B 抢到了锁。于是实例 A 删掉了实例 B 的锁。Lua 脚本让检查与删除成为原子操作。
替代方案:基于 etcd 租约的锁
对于已经在运行 etcd 的系统(例如 Kubernetes 环境),etcd 的租约机制比 Redis 提供了更强的保证:
// Alternative: etcd-based distributed lock
import (
clientv3 "go.etcd.io/etcd/client/v3"
"go.etcd.io/etcd/client/v3/concurrency"
)
func acquireEtcdLock(client *clientv3.Client, lockName string) (*concurrency.Mutex, error) {
session, err := concurrency.NewSession(client, concurrency.WithTTL(30))
if err != nil {
return nil, err
}
mutex := concurrency.NewMutex(session, "/locks/"+lockName)
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := mutex.TryLock(ctx); err != nil {
session.Close()
return nil, err
}
return mutex, nil
}
etcd 使用 Raft 共识,因此提供真正的线性一致性。相比之下,Redis 在 Redis Sentinel 或 Cluster 部署中发生故障切换时可能丢失锁。对于大多数任务调度场景,Redis 已经足够——在故障切换期间偶尔出现一次重复执行是可以接受的。对于金融或安全攸关的任务,应优先选择 etcd,或使用基于多个独立 Redis 实例的 Redlock 算法。
带分布式锁的 Cron 调度器
现在,我们将 robfig/cron 与分布式锁结合,创建一个能在所有实例间恰好触发一次任务的调度器:
// internal/scheduler/cron.go
package scheduler
import (
"context"
"errors"
"log/slog"
"time"
"github.com/RichardKnop/machinery/v2"
"github.com/RichardKnop/machinery/v2/tasks"
"github.com/redis/go-redis/v9"
"github.com/robfig/cron/v3"
"scheduler/internal/lock"
)
// JobDefinition describes a scheduled job.
type JobDefinition struct {
Name string // unique job name, used as the lock key
CronExpr string // cron expression, e.g. "0 */5 * * *"
TaskName string // registered Machinery task name
Args []tasks.Arg // task arguments
Queue string // target queue (for routing to specific workers)
MaxRetries int // max retry count
}
// CronScheduler wraps robfig/cron with distributed locking.
type CronScheduler struct {
cron *cron.Cron
server *machinery.Server
redisClient *redis.Client
lockTTL time.Duration
logger *slog.Logger
}
func NewCronScheduler(
server *machinery.Server,
redisClient *redis.Client,
lockTTL time.Duration,
logger *slog.Logger,
) *CronScheduler {
return &CronScheduler{
cron: cron.New(cron.WithSeconds()),
server: server,
redisClient: redisClient,
lockTTL: lockTTL,
logger: logger,
}
}
// Schedule registers a job definition with the cron scheduler.
func (s *CronScheduler) Schedule(job JobDefinition) error {
_, err := s.cron.AddFunc(job.CronExpr, func() {
s.fireJob(job)
})
if err != nil {
return err
}
s.logger.Info("scheduled job",
"name", job.Name,
"cron", job.CronExpr,
"task", job.TaskName,
)
return nil
}
// fireJob attempts to acquire the distributed lock and dispatch the task.
func (s *CronScheduler) fireJob(job JobDefinition) {
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
// Try to acquire the distributed lock
dlock := lock.NewRedisLock(s.redisClient, job.Name, s.lockTTL)
if err := dlock.Acquire(ctx); err != nil {
if errors.Is(err, lock.ErrLockNotAcquired) {
s.logger.Debug("another instance holds the lock, skipping",
"job", job.Name)
return
}
s.logger.Error("failed to acquire lock", "job", job.Name, "error", err)
return
}
defer dlock.Release(ctx)
// Build and send the Machinery task
signature := &tasks.Signature{
Name: job.TaskName,
Args: job.Args,
RetryCount: job.MaxRetries,
}
if job.Queue != "" {
signature.RoutingKey = job.Queue
}
result, err := s.server.SendTask(signature)
if err != nil {
s.logger.Error("failed to send task",
"job", job.Name,
"task", job.TaskName,
"error", err,
)
return
}
s.logger.Info("dispatched task",
"job", job.Name,
"task", job.TaskName,
"taskID", result.Signature.UUID,
)
}
// Start begins the cron scheduler.
func (s *CronScheduler) Start() {
s.cron.Start()
s.logger.Info("cron scheduler started")
}
// Stop gracefully stops the cron scheduler, waiting for running jobs to finish.
func (s *CronScheduler) Stop() context.Context {
return s.cron.Stop()
}
工作节点实现
工作节点是一个独立的二进制程序,连接到相同的消息中间件并处理任务:
// cmd/worker/main.go
package main
import (
"context"
"log/slog"
"os"
"os/signal"
"syscall"
"github.com/RichardKnop/machinery/v2"
backendsiface "github.com/RichardKnop/machinery/v2/backends/iface"
brokersiface "github.com/RichardKnop/machinery/v2/brokers/iface"
configMachinery "github.com/RichardKnop/machinery/v2/config"
lockiface "github.com/RichardKnop/machinery/v2/locks/iface"
eagerlock "github.com/RichardKnop/machinery/v2/locks/eager"
redisbackend "github.com/RichardKnop/machinery/v2/backends/redis"
redisbroker "github.com/RichardKnop/machinery/v2/brokers/redis"
"scheduler/internal/config"
"scheduler/internal/tasks"
)
func main() {
logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
cfg := config.Load()
// Configure Machinery
mCfg := &configMachinery.Config{
DefaultQueue: cfg.DefaultQueue,
ResultsExpireIn: 3600,
Redis: &configMachinery.RedisConfig{
MaxIdle: 3,
IdleTimeout: 240,
ReadTimeout: 15,
WriteTimeout: 15,
ConnectTimeout: 15,
NormalTasksPollPeriod: 1000,
DelayedTasksPollPeriod: 500,
},
}
broker := redisbroker.NewGR(mCfg, []string{cfg.BrokerURL}, 0)
backend := redisbackend.NewGR(mCfg, []string{cfg.ResultBackend}, 0)
lock := eagerlock.New()
server := machinery.NewServer(mCfg, broker, backend, lock)
// Register task handlers
if err := tasks.RegisterAll(server); err != nil {
logger.Error("failed to register tasks", "error", err)
os.Exit(1)
}
// Create a worker with a unique tag (hostname + PID works well)
hostname, _ := os.Hostname()
workerTag := fmt.Sprintf("worker-%s-%d", hostname, os.Getpid())
worker := server.NewWorker(workerTag, 10) // concurrency = 10
// Set up pre/post task hooks for logging
worker.SetPreTaskHandler(func(signature *tasks.Signature) {
logger.Info("starting task",
"taskID", signature.UUID,
"name", signature.Name,
"retry", signature.RetryCount,
)
})
worker.SetPostTaskHandler(func(signature *tasks.Signature) {
logger.Info("completed task",
"taskID", signature.UUID,
"name", signature.Name,
)
})
worker.SetErrorHandler(func(err error) {
logger.Error("task error", "error", err)
})
// Graceful shutdown
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
go func() {
<-ctx.Done()
logger.Info("shutting down worker...")
worker.Quit()
}()
logger.Info("starting worker", "tag", workerTag, "concurrency", 10)
if err := worker.Launch(); err != nil {
logger.Error("worker stopped with error", "error", err)
os.Exit(1)
}
}
任务中心(Server)入口
// cmd/server/main.go
package main
import (
"context"
"log/slog"
"os"
"os/signal"
"syscall"
"github.com/RichardKnop/machinery/v2"
"github.com/RichardKnop/machinery/v2/tasks"
"github.com/redis/go-redis/v9"
"scheduler/internal/config"
internalTasks "scheduler/internal/tasks"
"scheduler/internal/scheduler"
)
func main() {
logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
cfg := config.Load()
// Initialize Machinery server (same setup as worker)
server := initMachineryServer(cfg)
if err := internalTasks.RegisterAll(server); err != nil {
logger.Error("failed to register tasks", "error", err)
os.Exit(1)
}
// Initialize Redis client for distributed locking
redisOpts, _ := redis.ParseURL(cfg.RedisURL)
redisClient := redis.NewClient(redisOpts)
// Create and configure the cron scheduler
cronScheduler := scheduler.NewCronScheduler(server, redisClient, cfg.LockTTL, logger)
// Register scheduled jobs
jobs := []scheduler.JobDefinition{
{
Name: "weekly-report",
CronExpr: "0 0 9 * * 1", // Every Monday at 9:00 AM
TaskName: "send_report",
Args: []tasks.Arg{
{Type: "string", Value: "weekly"},
{Type: "string", Value: "team@example.com"},
},
MaxRetries: 3,
},
{
Name: "session-cleanup",
CronExpr: "0 0 */6 * * *", // Every 6 hours
TaskName: "clean_expired_sessions",
Args: []tasks.Arg{
{Type: "int64", Value: 30}, // 30 days
},
MaxRetries: 2,
},
{
Name: "external-sync",
CronExpr: "0 */5 * * * *", // Every 5 minutes
TaskName: "sync_external_data",
Args: []tasks.Arg{
{Type: "string", Value: "https://api.example.com/data"},
},
Queue: "sync_workers",
MaxRetries: 5,
},
}
for _, job := range jobs {
if err := cronScheduler.Schedule(job); err != nil {
logger.Error("failed to schedule job", "name", job.Name, "error", err)
os.Exit(1)
}
}
// Start the scheduler
cronScheduler.Start()
// Wait for shutdown signal
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
<-ctx.Done()
logger.Info("shutting down task center...")
shutdownCtx := cronScheduler.Stop()
<-shutdownCtx.Done()
logger.Info("task center stopped")
}
带指数退避的重试机制
Machinery 提供了内置的重试能力,但其默认行为是立即重试。对于生产系统,你会希望使用指数退避,以避免不断冲击一个正在失败的依赖。
// Configure retry behavior on task signatures
func newTaskSignatureWithBackoff(taskName string, args []tasks.Arg, maxRetries int) *tasks.Signature {
return &tasks.Signature{
Name: taskName,
Args: args,
RetryCount: maxRetries,
RetryTimeout: 10, // base retry delay in seconds
// Machinery multiplies RetryTimeout by 2^(attempt number)
// Attempt 1: 10s, Attempt 2: 20s, Attempt 3: 40s
}
}
若想对重试行为有更多控制,可以实现一个自定义的错误处理器来决定是否重试:
// internal/tasks/retry.go
package tasks
import (
"errors"
"fmt"
"math"
"time"
)
// PermanentError wraps an error to indicate it should NOT be retried.
type PermanentError struct {
Err error
}
func (e *PermanentError) Error() string { return e.Err.Error() }
func (e *PermanentError) Unwrap() error { return e.Err }
// RetryableError wraps an error with a specific delay before the next retry.
type RetryableError struct {
Err error
RetryIn time.Duration
}
func (e *RetryableError) Error() string {
return fmt.Sprintf("%s (retry in %s)", e.Err.Error(), e.RetryIn)
}
// CalculateBackoff returns the delay for the given attempt using
// exponential backoff with jitter.
func CalculateBackoff(attempt int, baseDelay time.Duration, maxDelay time.Duration) time.Duration {
delay := time.Duration(float64(baseDelay) * math.Pow(2, float64(attempt)))
if delay > maxDelay {
delay = maxDelay
}
// Add up to 25% jitter to prevent thundering herd
jitter := time.Duration(float64(delay) * 0.25 * rand.Float64())
return delay + jitter
}
在你的任务处理器中,使用这些错误类型来控制重试行为:
func SyncExternalData(apiEndpoint string) (string, error) {
// ... (fetch logic)
if resp.StatusCode == http.StatusNotFound {
// 404 won't fix itself — don't retry
return "", &PermanentError{Err: fmt.Errorf("endpoint %s returned 404", apiEndpoint)}
}
if resp.StatusCode == http.StatusTooManyRequests {
// Rate limited — retry with longer backoff
return "", &RetryableError{
Err: fmt.Errorf("rate limited by %s", apiEndpoint),
RetryIn: 60 * time.Second,
}
}
// Other errors: use default retry behavior
if resp.StatusCode >= 500 {
return "", fmt.Errorf("server error %d from %s", resp.StatusCode, apiEndpoint)
}
// ... (success path)
}
任务状态监控
在分布式系统中,可观测性并非可选项。你需要知道哪些任务正在运行、哪些失败了,以及它们各自耗时多久。
查询任务状态
Machinery 将任务状态存储在结果后端中。你可以通过编程方式查询:
// internal/monitor/status.go
package monitor
import (
"fmt"
"time"
"github.com/RichardKnop/machinery/v2"
"github.com/RichardKnop/machinery/v2/backends/result"
)
type TaskStatus struct {
ID string `json:"id"`
State string `json:"state"`
Result string `json:"result,omitempty"`
Error string `json:"error,omitempty"`
CreatedAt time.Time `json:"created_at"`
}
// GetTaskStatus retrieves the current state of a task by its ID.
func GetTaskStatus(server *machinery.Server, taskID string) (*TaskStatus, error) {
asyncResult := result.NewAsyncResult(&tasks.Signature{UUID: taskID}, server.GetBackend())
taskState := asyncResult.GetState()
status := &TaskStatus{
ID: taskID,
State: taskState.State,
CreatedAt: taskState.CreatedAt,
}
if taskState.Error != "" {
status.Error = taskState.Error
}
if taskState.IsSuccess() {
for _, r := range taskState.Results {
status.Result += fmt.Sprintf("%v ", r.Value)
}
}
return status, nil
}
// WaitForResult blocks until the task completes or the timeout expires.
func WaitForResult(server *machinery.Server, taskID string, timeout time.Duration) (*TaskStatus, error) {
asyncResult := result.NewAsyncResult(&tasks.Signature{UUID: taskID}, server.GetBackend())
results, err := asyncResult.GetWithTimeout(timeout, 500*time.Millisecond)
if err != nil {
return nil, fmt.Errorf("task %s: %w", taskID, err)
}
status := &TaskStatus{
ID: taskID,
State: "SUCCESS",
}
for _, r := range results {
status.Result += fmt.Sprintf("%v ", r.Interface())
}
return status, nil
}
Prometheus 指标
导出指标,让你的监控栈能够跟踪任务健康状况:
// internal/monitor/metrics.go
package monitor
import (
"net/http"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
var (
TasksDispatched = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "scheduler_tasks_dispatched_total",
Help: "Total number of tasks dispatched to the queue",
},
[]string{"task_name"},
)
TasksCompleted = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "scheduler_tasks_completed_total",
Help: "Total number of tasks completed successfully",
},
[]string{"task_name"},
)
TasksFailed = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "scheduler_tasks_failed_total",
Help: "Total number of tasks that failed (including retries exhausted)",
},
[]string{"task_name"},
)
TaskDuration = promauto.NewHistogramVec(
prometheus.HistogramOpts{
Name: "scheduler_task_duration_seconds",
Help: "Time taken to execute a task",
Buckets: prometheus.ExponentialBuckets(0.1, 2, 10), // 0.1s to ~51s
},
[]string{"task_name"},
)
TaskRetries = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "scheduler_task_retries_total",
Help: "Total number of task retries",
},
[]string{"task_name"},
)
QueueDepth = promauto.NewGaugeVec(
prometheus.GaugeOpts{
Name: "scheduler_queue_depth",
Help: "Current number of tasks waiting in the queue",
},
[]string{"queue_name"},
)
)
// StartMetricsServer starts a Prometheus metrics endpoint.
func StartMetricsServer(port int) {
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
go http.ListenAndServe(fmt.Sprintf(":%d", port), mux)
}
之后你可以配置 Grafana 仪表盘和 Prometheus 告警规则:
# prometheus-alerts.yml
groups:
- name: scheduler
rules:
- alert: TaskFailureRateHigh
expr: rate(scheduler_tasks_failed_total[5m]) > 0.1
for: 5m
labels:
severity: warning
annotations:
summary: "High task failure rate"
description: "Task {{ $labels.task_name }} failing at {{ $value }} per second"
- alert: QueueBacklog
expr: scheduler_queue_depth > 1000
for: 10m
labels:
severity: critical
annotations:
summary: "Task queue backlog growing"
description: "Queue {{ $labels.queue_name }} has {{ $value }} pending tasks"
生产环境考量
优雅停机
在部署新版本时,你需要让工作节点在停止前完成当前正在执行的任务。上面的工作节点代码已经通过 signal.NotifyContext 处理了 SIGTERM。要点如下:
- 停止接收新任务——调用
worker.Quit(),它会停止从中间件消费 - 等待进行中的任务——Machinery 的
Quit()会等待当前正在运行的 goroutine - 设置一个截止时间——Kubernetes 会先发送
SIGTERM,在terminationGracePeriodSeconds(默认 30 秒)之后再发送SIGKILL。将其设置为高于你运行时间最长的任务
# kubernetes deployment snippet
spec:
terminationGracePeriodSeconds: 300 # 5 minutes for long-running tasks
containers:
- name: worker
lifecycle:
preStop:
exec:
command: ["sleep", "5"] # Allow load balancer to drain
任务幂等性
网络故障、工作节点崩溃以及中间件的重投递,都可能导致一个任务被运行不止一次。每个任务处理器都必须是幂等的——用相同参数运行两次应当产生相同的结果。
实用技巧:
- 使用唯一请求 ID:将 Machinery 的任务 UUID 传入你的业务逻辑,并将其用作去重键
- 用 upsert 而非 insert:使用
INSERT ... ON DUPLICATE KEY UPDATE或等价语句 - 先检查后动作:在做出更改前先校验当前状态(例如,如果
email_sent标志已置位,就不要再发送邮件)
func SendReport(reportType string, recipientEmail string) (string, error) {
// Idempotency: check if this report was already sent today
key := fmt.Sprintf("report:%s:%s:%s", reportType, recipientEmail,
time.Now().Format("2006-01-02"))
exists, err := redisClient.Exists(ctx, key).Result()
if err == nil && exists > 0 {
return "already_sent", nil
}
// ... generate and send report ...
// Mark as sent (expire after 24h)
redisClient.Set(ctx, key, "1", 24*time.Hour)
return "sent", nil
}
死信队列
耗尽所有重试次数的任务不应被悄无声息地丢弃。把它们推入死信队列以供人工排查:
func (w *Worker) SetErrorHandler(fn func(err error)) {
// In addition to the error handler, check if retries are exhausted
worker.SetPostTaskHandler(func(signature *tasks.Signature) {
state := server.GetBackend().GetState(signature.UUID)
if state.IsFailure() && signature.RetryCount <= 0 {
// Push to dead letter queue
deadLetterSignature := *signature
deadLetterSignature.RoutingKey = "dead_letter"
deadLetterSignature.RetryCount = 0
server.SendTask(&deadLetterSignature)
logger.Error("task moved to dead letter queue",
"taskID", signature.UUID,
"name", signature.Name,
"error", state.Error,
)
}
})
}
健康检查
暴露一个健康检查端点,用于验证与 Redis 和消息中间件的连通性:
func healthHandler(redisClient *redis.Client) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
ctx, cancel := context.WithTimeout(r.Context(), 3*time.Second)
defer cancel()
if err := redisClient.Ping(ctx).Err(); err != nil {
w.WriteHeader(http.StatusServiceUnavailable)
fmt.Fprintf(w, `{"status":"unhealthy","error":"%s"}`, err)
return
}
w.WriteHeader(http.StatusOK)
fmt.Fprint(w, `{"status":"healthy"}`)
}
}
把它们组装起来
以下是一套生产环境部署的拓扑结构:
| 组件 | 实例数 | 扩展策略 |
|---|---|---|
| 任务中心 | 2-3(高可用) | 固定——分布式锁防止重复调度 |
| 通用工作节点 | 3-10 | 根据队列深度自动扩缩 |
| 同步工作节点 | 2-5 | 根据 sync_workers 队列深度自动扩缩 |
| Redis | 1 主 + 若干副本 | 标准 Redis 高可用部署 |
| 数据库 | 1 主 + 只读副本 | 存储任务定义与审计日志 |
| Prometheus + Grafana | 1 | 监控与告警 |
所有任务中心实例都运行相同的 cron 调度,但分布式锁确保每个任务只由一个实例分发。如果某个实例宕机,另一个实例会在下一个 cron 触发点接管——无需手动故障切换。
工作节点可以按队列独立扩展。如果 sync_workers 队列开始积压,只需增加同步工作节点,而不会影响通用任务的处理。
总结
构建分布式任务调度系统并不依赖某个单一技巧——它关乎将若干成熟的模式组合起来:
- 通过消息队列将调度与执行分离
- 用分布式锁(Redis SETNX 或 etcd 租约)防止重复调度
- 用带抖动的指数退避实现重试逻辑
- 用死信队列处理耗尽重试次数的任务
- 幂等的任务处理器,因为恰好一次投递(exactly-once)只是一个神话
- 用 Prometheus 指标观测队列深度、失败率和任务耗时
- 优雅停机,以避免在部署时杀掉进行中的任务
Machinery 为任务队列和结果存储提供了坚实的基础。在其之上叠加带分布式锁的 cron 调度,再加上监控,你就拥有了一套可投入生产的系统。
完整的源代码在结构上是模块化的——你可以把 Redis 换成 RabbitMQ 作为中间件,如果偏好更简单的 API 可以从 Machinery 切换到 Asynq,或者当工作流复杂度要求时升级到 Temporal。无论你选择哪些库,底层的架构模式始终不变。
参考资料
- Go documentation — Go
- robfig/cron — GitHub
- Distributed Locks with Redis — Redis Documentation