共计 2865 个字符,预计需要花费 8 分钟才能阅读完成。
背景痛点
在微服务架构中,任务调度一直是一个棘手的问题。传统的调度方案通常依赖于中心化的调度器,比如使用 CronJob 或者一些开源的调度框架。这种方案在分布式环境中存在几个明显的缺陷:

- 单点故障:中心化调度器一旦宕机,整个系统的任务调度就会中断。
- 缺乏弹性:任务执行器的数量和资源无法根据负载动态调整,容易导致资源浪费或任务积压。
- 任务状态管理复杂:分布式环境下,任务的执行状态和结果难以统一管理,尤其是在网络分区的情况下。
技术选型
为了解决这些问题,我们对比了几种常见的任务调度方案:
- 传统调度框架(如 Quartz、Celery):
- 优点:成熟稳定,社区支持丰富。
-
缺点:中心化设计,弹性不足,扩展性有限。
-
基于消息队列的方案(如 Kafka、RabbitMQ):
- 优点:解耦生产者和消费者,支持高吞吐。
-
缺点:需要额外开发任务状态管理和失败重试逻辑。
-
Agent 工具调用助手:
- 优点:去中心化设计,支持动态扩缩容,内置智能路由和失败重试机制。
- 缺点:实现复杂度较高,需要开发者对分布式系统有一定了解。
综合考虑后,我们选择了基于 Agent 的方案,因为它能够更好地满足我们对高可用性和弹性的需求。
核心架构
交互流程
以下是 Agent 与任务执行器的交互流程(用序列图表示):
sequenceDiagram
participant Scheduler
participant Agent
participant Executor
Scheduler->>Agent: 下发任务
Agent->>Executor: 执行任务
Executor->>Agent: 返回结果
Agent->>Scheduler: 更新任务状态
关键组件设计
- 任务队列:
- 使用优先级队列管理待执行任务,支持任务优先级设置。
-
任务队列持久化到数据库,确保系统重启后任务不丢失。
-
状态机:
- 每个任务的生命周期通过状态机管理,状态包括
PENDING、RUNNING、SUCCESS、FAILED、RETRY。 -
状态转换逻辑清晰,便于调试和监控。
-
心跳检测:
- Agent 定期向调度器发送心跳,报告自身状态和负载情况。
- 调度器根据心跳信息动态分配任务,避免任务堆积在某个 Agent 上。
代码实现
Agent 核心代码(Python)
import threading
import time
from queue import PriorityQueue
class TaskAgent:
def __init__(self):
self.task_queue = PriorityQueue()
self.lock = threading.Lock()
self.running = True
def add_task(self, task, priority=0):
"""添加任务到队列"""
with self.lock:
self.task_queue.put((priority, task))
def run(self):
"""启动 Agent,持续从队列中取出任务执行"""
while self.running:
if not self.task_queue.empty():
_, task = self.task_queue.get()
try:
result = task.execute()
self.update_task_status(task.id, 'SUCCESS', result)
except Exception as e:
self.handle_failure(task, e)
time.sleep(0.1)
def handle_failure(self, task, error):
"""处理任务失败"""
if task.retry_count < task.max_retries:
task.retry_count += 1
self.add_task(task, priority=task.retry_count)
self.update_task_status(task.id, 'RETRY', str(error))
else:
self.update_task_status(task.id, 'FAILED', str(error))
任务幂等性处理
class Task:
def __init__(self, id, func, args, kwargs, max_retries=3):
self.id = id
self.func = func
self.args = args
self.kwargs = kwargs
self.max_retries = max_retries
self.retry_count = 0
def execute(self):
"""执行任务,确保幂等性"""
# 检查任务是否已经执行过
if self._is_executed():
return self._get_executed_result()
# 执行任务
result = self.func(*self.args, **self.kwargs)
# 保存执行结果
self._save_result(result)
return result
性能考量
吞吐量测试
我们在不同负载下测试了系统的吞吐量,以下是测试结果(单位:任务 / 秒):
| 并发 Agent 数量 | 平均吞吐量 | 峰值吞吐量 |
|---|---|---|
| 1 | 100 | 120 |
| 5 | 450 | 500 |
| 10 | 850 | 900 |
容错策略
在网络分区的情况下,系统通过以下策略保证任务的可靠性:
- 任务状态持久化:所有任务状态变更都会持久化到数据库,确保系统恢复后能够继续处理。
- 心跳超时检测:如果 Agent 长时间未发送心跳,调度器会将其标记为不可用,并重新分配其任务。
- 最终一致性:任务状态通过异步方式同步,确保在网络分区恢复后系统能够最终一致。
避坑指南
资源泄漏
在开发过程中,我们遇到过以下几种资源泄漏的场景:
- 数据库连接未关闭:任务执行完毕后忘记关闭数据库连接,导致连接池耗尽。
- 线程未正确终止:Agent 停止时未正确终止后台线程,导致线程泄漏。
解决方案是确保所有资源在使用完毕后正确释放,例如使用 with 语句管理资源:
with get_db_connection() as conn:
conn.execute("SELECT * FROM tasks")
分布式锁
分布式锁的正确使用方式:
- 锁的粒度:锁的粒度要尽可能小,避免长时间持有锁导致系统性能下降。
- 超时设置:为锁设置合理的超时时间,防止死锁。
- 锁的可重入性:确保同一个线程可以多次获取同一把锁,避免死锁。
延伸思考
结合 K8s 实现弹性扩缩容
我们可以通过 K8s 的 HPA(Horizontal Pod Autoscaler)功能实现 Agent 的弹性扩缩容。具体步骤如下:
- 定义自定义指标:通过 Prometheus 采集 Agent 的负载指标(如任务队列长度、CPU 使用率)。
- 配置 HPA:根据指标动态调整 Agent 的副本数量。
- 平滑缩容:在缩容时,确保正在执行的任务能够正常完成,避免任务中断。
结尾
通过 Agent 工具调用助手,我们实现了一个高可用、高弹性的任务调度系统。但在实际生产环境中,仍然有一些问题需要进一步优化,例如:
- 如何在大规模集群中降低调度器的负载?
- 如何进一步减少任务执行的延迟?
欢迎大家在评论区分享自己的经验和想法!
正文完
