共计 1518 个字符,预计需要花费 4 分钟才能阅读完成。
背景痛点:传统方案的致命缺陷
在电商大促期间,我们曾因传统定时任务框架的缺陷付出惨痛代价。典型的单机 Cron 方案在分布式环境下暴露出三大致命伤:

- 单点故障:当唯一执行节点宕机时,整个促销活动定时任务全部停滞
- 雪崩效应:某次数据库连接超时导致任务线程阻塞,后续任务堆积引发 OOM
- 时间漂移:各节点 NTP 同步差异导致重复执行,优惠券被超额发放
graph TD
A[传统 Cron] --> B{单点故障}
A --> C{雪崩效应}
A --> D{时钟漂移}
架构选型:Agent 模式脱颖而出
对比三种主流方案后,我们选择 Agent 架构实现解耦:
- Cron 方案
- 优点:实现简单
-
缺点:无法水平扩展,缺乏容错
-
Work Queue 方案
- 优点:解耦生产者消费者
-
缺点:中间件成为新单点
-
Agent 方案
- 动态任务分片
- 自动故障转移
- 横向扩展能力
核心实现:Go 语言实战
任务抢占机制
通过 Redis 分布式锁实现安全抢占,关键代码:
// 获取分布式锁(设置 10 秒 TTL 避免死锁)func acquireLock(rdb *redis.Client, key string) bool {return rdb.SetNX(ctx, key, "locked", 10*time.Second).Val()}
// 任务分片示例
func dispatchTasks(agents []string) map[string][]Task {// 一致性哈希算法分配任务}
幂等性保障
每个任务携带唯一指纹,处理前校验状态:
type Task struct {
ID string // UUID
Payload string
Fingerprint string // md5(Payload+Timestamp)
}
func handleTask(t Task) error {if isProcessed(t.Fingerprint) {return nil // 已处理则跳过}
// ... 业务逻辑
}
生产级优化
Goroutine 泄漏防护
使用 context+sync.WaitGroup 双重保障:
func workerPool(ctx context.Context) {
var wg sync.WaitGroup
defer wg.Wait() // 等待所有协程退出
for {
select {case <-ctx.Done():
return
default:
wg.Add(1)
go func() {defer wg.Done()
// ... 任务处理
}()}
}
}
监控埋点示例
Prometheus 指标采集关键指标:
var (
tasksCounter = prometheus.NewCounterVec(
prometheus.CounterOpts{
Name: "scheduler_tasks_total",
Help: "Total processed tasks",
},
[]string{"status"}, // success/failure
)
)
func init() {prometheus.MustRegister(tasksCounter)
}
血泪教训:三大生产事故
- NTP 不同步事件
- 现象:跨机房部署时,5 秒时钟差导致重复调度
-
解决:采用租约机制 (lease) 替代绝对时间
-
Redis 网络抖动
- 现象:锁提前释放引发任务冲突
-
优化:增加锁续期心跳检测
-
内存泄漏事件
- 现象:任务日志未清理,半年占满磁盘
- 方案:增加日志轮转和自动归档
延伸思考
- 如何实现 Agent 的灰度发布?可以考虑版本标记 + 流量分发
- 当任务量突增 10 倍时,如何动态调整 Agent 规模?需要结合 K8s HPA 实现自动扩缩容
写在最后
这套 Agent 系统目前已稳定运行 2 年,日均处理任务量超过 300 万。最大的收获是:分布式系统没有银弹,必须根据业务特点持续迭代。建议读者从简单场景入手,逐步完善监控和容错机制。
正文完
