用 Go 构建分布式任务调度系统

使用 Go、Machinery、Redis 和 cron 设计并交付一套生产级任务调度器。涵盖分布式锁与重试策略。

zhuermu··15 分钟
GolangDistributed SystemsJob SchedulingMachineryRedisMicroservices

每个后端服务最终都会需要定时任务。发送周报、清理过期会话、从外部 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 到 Completed、Failed、Retrying 或 Dead Letter 的各个状态

任务从 Created(已创建) 状态开始,当调度器将其推送到消息中间件时进入 Queued(已入队),被工作节点取走时进入 Running(运行中),最终以 Completed(已完成)Failed(已失败) 结束。仍有剩余重试次数的失败任务进入 Retrying(重试中) 状态,并以指数退避重新入队。耗尽所有重试次数的任务进入 Dead Letter(死信) 队列,以供人工排查。


技术选型

Machinery——基石

Machinery 是一个用于分布式任务处理的 Go 框架,灵感来自 Python 的 Celery。它开箱即用地提供了任务注册、消息中间件抽象、结果存储和重试逻辑。

不过,单靠 machinery 并不能处理 cron 调度或分布式锁——这些层需要我们自己添加。

框架对比

在确定使用 Machinery 之前,我们先对比一下 Go 生态中的主要选项:

特性MachineryAsynqTemporalgo-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。要点如下:

  1. 停止接收新任务——调用 worker.Quit(),它会停止从中间件消费
  2. 等待进行中的任务——Machinery 的 Quit() 会等待当前正在运行的 goroutine
  3. 设置一个截止时间——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 队列深度自动扩缩
Redis1 主 + 若干副本标准 Redis 高可用部署
数据库1 主 + 只读副本存储任务定义与审计日志
Prometheus + Grafana1监控与告警

所有任务中心实例都运行相同的 cron 调度,但分布式锁确保每个任务只由一个实例分发。如果某个实例宕机,另一个实例会在下一个 cron 触发点接管——无需手动故障切换。

工作节点可以按队列独立扩展。如果 sync_workers 队列开始积压,只需增加同步工作节点,而不会影响通用任务的处理。


总结

构建分布式任务调度系统并不依赖某个单一技巧——它关乎将若干成熟的模式组合起来:

  1. 通过消息队列将调度与执行分离
  2. 分布式锁(Redis SETNX 或 etcd 租约)防止重复调度
  3. 带抖动的指数退避实现重试逻辑
  4. 死信队列处理耗尽重试次数的任务
  5. 幂等的任务处理器,因为恰好一次投递(exactly-once)只是一个神话
  6. Prometheus 指标观测队列深度、失败率和任务耗时
  7. 优雅停机,以避免在部署时杀掉进行中的任务

Machinery 为任务队列和结果存储提供了坚实的基础。在其之上叠加带分布式锁的 cron 调度,再加上监控,你就拥有了一套可投入生产的系统。

完整的源代码在结构上是模块化的——你可以把 Redis 换成 RabbitMQ 作为中间件,如果偏好更简单的 API 可以从 Machinery 切换到 Asynq,或者当工作流复杂度要求时升级到 Temporal。无论你选择哪些库,底层的架构模式始终不变。

参考资料

  1. Go documentation — Go
  2. robfig/cron — GitHub
  3. Distributed Locks with Redis — Redis Documentation