共计 2266 个字符,预计需要花费 6 分钟才能阅读完成。
背景痛点
在分布式系统中,任务调度常常面临几个棘手的问题。首先是任务雪崩,比如某个服务突然不可用,导致大量任务积压,当服务恢复时,这些任务又同时涌向系统,造成二次崩溃。其次是状态不一致,由于网络分区或节点故障,不同节点对任务状态的判断可能产生分歧。还有重复执行的问题,特别是在重试机制下,同一个任务可能被多次处理。

这些问题在电商秒杀、定时报表生成、大数据处理等场景中尤为突出。我曾经遇到过一个案例:一个定时任务因为网络抖动被重复提交了三次,导致下游系统处理了三次相同的订单,最终不得不人工介入修复数据。
架构对比
解决这些问题有多种架构范式可选,我们需要根据具体场景权衡:
- Actor 模型:每个 Agent 作为一个独立 Actor,通过消息传递通信。优势是状态隔离性好,适合有复杂状态转换的场景。但调试相对困难,且消息传递可能成为瓶颈。
- CQRS(命令查询职责分离):将任务提交(命令)和状态查询分离。优点是读写可以独立扩展,适合查询压力大的系统。缺点是实现复杂度高,需要维护最终一致性。
- 事件溯源:将所有状态变化记录为事件流。便于回溯和调试,但存储开销大,不适合高频任务。
对于大多数任务调度场景,我推荐轻量级的 Actor 模型变种:每个任务对应一个 Agent 实例,通过共享队列接收任务,但各自维护执行状态。这样既保证了隔离性,又避免了纯 Actor 模型的消息开销。
核心实现
基础框架(Go 语言示例)
// Agent 核心结构
type TaskAgent struct {
ID string
Queue chan Task // 任务队列
StopChan chan struct{} // 停止信号
IsWorking atomic.Bool // 当前是否正在处理任务
}
// 启动 Agent
go func(a *TaskAgent) Run() {
for {
select {
case task := <-a.Queue:
a.IsWorking.Store(true)
// 幂等性检查
if processed := checkProcessed(task.ID); processed {log.Printf("Task %s already processed", task.ID)
continue
}
// 执行任务(带重试)err := retry(3, 100*time.Millisecond, func() error {return executeTask(task)
})
if err != nil {log.Printf("Task %s failed: %v", task.ID, err)
} else {markProcessed(task.ID) // 标记已完成
}
a.IsWorking.Store(false)
case <-a.StopChan:
return
}
}
}
关键模块说明
- 任务队列消费
- 使用带缓冲的 channel 控制并发度(如
make(chan Task, 100)) -
通过
IsWorking标志防止同一个 Agent 同时处理多个任务 -
心跳检测
go
// 心跳协程
go func() {
ticker := time.NewTicker(5 * time.Second)
for {
select {
case <-ticker.C:
reportHeartbeat(a.ID)
case <-a.StopChan:
ticker.Stop()
return
}
}
}() -
幂等性处理
- 使用 Redis 记录已处理任务 ID,设置合理过期时间
- 对于重复任务,对比任务内容哈希值而不仅是 ID
性能优化
通过基准测试发现几个关键数据点:
- 单个 Agent 的 QPS 与任务复杂度强相关:
- 简单任务(<10ms):约 1200 QPS
- 中等任务(100ms):约 90 QPS
-
复杂任务(1s):约 8 QPS
-
内存占用主要来自:
- 任务队列缓冲(每个任务约 500B)
- 长连接保持(每个 Agent 约 20KB)
优化建议:
- 对于短时任务,适当增大队列缓冲(100-1000)
- 对于长时任务,采用工作池模式限制并发
- 心跳间隔根据网络质量调整(通常 5 -30 秒)
避坑指南
分布式锁的常见误区
-
过度加锁:只在真正需要互斥的地方加锁,比如:
// 错误做法:锁住整个处理流程 lock(task.ID) defer unlock(task.ID) processTask(task) // 正确做法:只锁状态检查 if lock(task.ID) {if !isProcessed(task.ID) {processTask(task) } unlock(task.ID) } -
忽略锁过期:设置合理的锁 TTL,避免节点宕机导致死锁
任务状态机反模式
-
状态爆炸:避免定义过多中间状态(如 ”pre-processing”、”post-processing”)。建议保持:pending -> processing -> succeeded/failed
-
双向转换:不允许状态回退(如 failed -> processing),这会引入复杂性。应该创建新任务实例。
实践建议
上线后需要监控的关键指标:
- 任务积压率:队列长度 / 处理能力
- 平均恢复时间:从故障到重新处理的时间
- 重复执行率:幂等检查拦截的任务比例
推荐监控工具组合:
- Prometheus + Grafana 用于指标可视化
- ELK 用于日志分析
- Jaeger 用于分布式追踪
开放性问题
在实际应用中,我们经常面临这些权衡:
- 如何确定最优的队列容量?太小会导致任务丢弃,太大可能引起内存问题
- 心跳间隔设置多少合适?太频繁会增加负载,间隔太长会影响故障发现速度
- 对于金融级场景,如何在不牺牲性能的情况下实现强一致性?
这些问题的答案往往取决于具体业务场景,需要结合压测数据和业务容忍度来决定。欢迎在评论区分享你的实践经验。
