在构建后端系统时,Golang常被用于编写任务调度服务,例如定时生成报表、清理过期数据或同步第三方接口。当服务以多副本方式部署,如何避免同一任务被多个实例同时执行,以及在执行中断后如何恢复状态,是决定系统可靠性的核心问题。单纯使用time.Ticker或robfig/cron只能在单进程内生效,集群环境下必须引入外部协调组件。

为什么本地调度无法保证一致性
多数初学者会直接在代码里启动一个定时协程,逻辑看起来简单清晰。然而当我们将二进制部署到三台机器上,每台机器上的定时器都会独立触发,导致任务被执行三次。对于写文件或发通知这类操作,重复执行可能只是浪费资源;但对于扣款、对账等场景,重复会直接引发资损。
另一种隐蔽问题是时钟漂移。不同宿主机的系统时间可能存在秒级偏差,即使使用分布式cron框架,也可能出现两个节点几乎同时认为“到了执行时间”。因此,可靠的调度方案必须把“谁有权执行”的决策交给一个外部可信源,而不是依赖各节点的本地时钟。
package main
import (
"time"
"log"
)
// 错误示例:仅本地调度,多副本会重复执行
func main() {
ticker := time.NewTicker(10 * time.Second)
for range ticker.C {
log.Println("执行任务") // 每个实例都会打印
}
}
基于分布式锁的抢占方案
最直观的解决办法是利用Redis或etcd提供的分布式锁。某个Golang实例在任务触发时,先尝试获取一个带有过期时间的锁,比如键名为task:report:lock。只有成功设值的实例才继续执行,其他实例立即返回。这样在任意时刻,集群中至多只有一个执行者。
为了避免执行时间超过锁过期时间,主节点需要启动一个续租协程,以锁TTL的三分之一为间隔重新设置过期时间。当主节点崩溃,续租停止,锁在TTL到期后自动释放,备用节点即可接管。下面是使用Redis的SET NX EX命令实现抢锁与续租的示例。
package main
import (
"context"
"time"
"github.com/go-redis/redis/v8"
"log"
)
var rdb = redis.NewClient(&redis.Options{Addr: "127.0.0.1:6379"})
var lockKey = "task:report:lock"
// tryLock 尝试获取锁并启动续租
func tryLock(ctx context.Context) (bool, func()) {
ok, err := rdb.SetNX(ctx, lockKey, "1", 30*time.Second).Result()
if err != nil || !ok {
return false, nil
}
// 续租
stop := make(chan struct{})
go func() {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for {
select {
case <-stop:
return
case <-ticker.C:
rdb.Expire(ctx, lockKey, 30*time.Second)
}
}
}()
return true, func() { close(stop); rdb.Del(ctx, lockKey) }
}
func main() {
ctx := context.Background()
ok, unlock := tryLock(ctx)
if !ok {
log.Println("未获取到锁,退出")
return
}
defer unlock()
// 执行实际任务
log.Println("开始执行任务")
time.Sleep(5 * time.Second)
}
任务幂等与去重表设计
分布式锁解决了“同时只有一人执行”,但无法覆盖这样的边缘情况:主节点执行完任务、在释放锁之前网络闪断,此时锁过期,备节点重新执行。若任务本身不幂等,仍会产生重复副作用。因此,关键任务应设计幂等能力。
常见做法是引入去重表,以任务类型加日期或唯一业务标识作为主键。执行前先插入一行,如果冲突则说明已处理。配合数据库事务,可以保证即使重复调度也不会重复写数据。下例展示用MySQL唯一索引做屏障:
package main
import (
"database/sql"
"log"
_ "github.com/go-sql-driver/mysql"
)
// 假设表结构: CREATE TABLE task_log (id VARCHAR(64) PRIMARY KEY, done INT);
func markDone(db *sql.DB, taskID string) bool {
_, err := db.Exec("INSERT INTO task_log(id, done) VALUES(?, 1)", taskID)
if err != nil {
log.Println("任务已执行或插入失败:", err)
return false
}
return true
}
func main() {
db, _ := sql.Open("mysql", "user:pass@tcp(127.0.0.1:3306)/test")
if markDone(db, "report-20240101") {
log.Println("首次执行,继续处理")
} else {
log.Println("重复调度,直接返回")
}
}
消息队列解耦调度与执行
除了锁方案,还可以将调度器仅负责“投递消息”,真正的执行交给消费者组。例如用Kafka或Redis Stream,调度节点在定时点向队列发送一条任务消息,队列保证至少投递一次,消费者使用上文去重表屏蔽重复。这种架构下,调度层即使全部重启也不会丢失任务,因为消息已持久化。
该方案的优点是执行节点可动态扩缩容,且任务状态由消息偏移量记录。缺点是引入中间件复杂度,且在网络分区时可能暂时无法消费。对于报表类非实时任务,这是性价比很高的选择。
| 方案 | 一致性强度 | 实现复杂度 | 适用场景 |
|---|---|---|---|
| 本地cron | 无 | 低 | 单实例测试 |
| Redis分布式锁 | 弱一致 | 中 | 通用定时任务 |
| 消息队列 | 最终一致 | 高 | 高吞吐异步处理 |
故障恢复与监控建议
无论采用哪种方案,都必须对任务执行结果做监控。推荐在任务开始和结束时向Prometheus推送指标,并在锁异常或续租失败时报警。同时保留任务执行日志,便于在资损发生时回溯是哪个节点在何时执行。
对于跨天的长任务,建议将大任务拆成小批次,每批更新进度到数据库。这样即使进程被杀,重启后也能从断点继续,而不是从头再来。Golang的context包可用来传递取消信号,让清理逻辑及时释放资源。
可靠性不是某一个组件的功劳,而是锁、幂等、监控三者配合的结果。在设计之初就把重复执行当作必然会发生的事,系统才真正健壮。