一、为什么需要自研调度系统
在百度广告业务线,我们有上千个定时任务需要调度执行,涵盖广告创意审核、账户预算检查、投放数据聚合、报表生成等多个业务场景。初期团队使用的是基于 Quartz 集群的方案,但随着业务规模增长,这套方案逐渐暴露出几个关键问题:
- 单点瓶颈:Quartz 的集群模式本质上还是依赖数据库行锁做选主,当任务量超过阈值后,数据库的压力急剧上升。
- 缺乏动态分片能力:广告创意审核任务每天需要处理千万级数据,单节点执行耗时过长,但 Quartz 不支持将一个任务动态拆分到多节点并行执行。
- 容灾恢复慢:当某个执行节点宕机时,Quartz 的 misfire 机制处理不够精细,导致大量任务积压。
经过技术选型评估,我们最终决定基于 Zookeeper 的临时节点和 Watcher 机制自研调度系统。核心思路是:用 ZK 做协调,用时间轮做调度,用消息队列做解耦。
二、Zookeeper 选主与分布式锁
调度系统的第一个核心问题是 选主。我们需要在多个 Scheduler 节点中选出一个 Master,由 Master 负责任务的分发和调度决策,而 Worker 节点只负责任务的实际执行。
public class ZkElection {
private final CuratorFramework zkClient;
private final String electionPath = "/scheduler/election/master";
private LeaderSelectorListenerAdapter listener;
public void startElection() {
listener = new LeaderSelectorListenerAdapter() {
@Override
public void takeLeadership(CuratorFramework client) throws Exception {
// 当前的节点被选举为 Master
while (true) {
// Master 职责:任务扫描、分片、分发
scheduleTasks();
Thread.sleep(5000);
}
}
};
LeaderSelector selector = new LeaderSelector(zkClient, electionPath, listener);
selector.autoRequeue();
selector.start();
}
}
这里有几个关键设计决策:
- 使用 Curator 的 LeaderSelector 而非手动创建临时节点。LeaderSelector 内部封装了创建临时顺序节点的逻辑,并且支持
autoRequeue(),当 Master 失去领导权后可以重新参与选举。 - Session 超时设置:我们将 ZK Session 超时设为 8 秒(默认 30 秒太长),这样当 Master 网络分区时,能在较短时间内完成重新选举。
- 分布式锁兜底:对于某些需要互斥执行的任务(如预算结算),我们在 ZK 上创建了分布式锁,使用
InterProcessMutex实现:
InterProcessMutex lock = new InterProcessMutex(zkClient, "/locks/budget-settle/" + accountId);
if (lock.acquire(10, TimeUnit.SECONDS)) {
try {
executeBudgetSettlement(accountId);
} finally {
lock.release();
}
}
相比 Redis 分布式锁,ZK 锁的优势在于天然保证锁的持有时长与 Session 绑定,即使持有锁的进程崩溃,Session 过期后锁也会自动释放,不会出现死锁。
三、时间轮调度与执行器解耦
任务调度层我们采用了 Hierarchical Timing Wheel(层级时间轮) 的设计,参考了 Netty 的 HashedWheelTimer 但做了分布式适配:
- 第一层(秒级):tick 间隔 1 秒,1000 个 slot,覆盖 0~999 秒的任务。
- 第二层(分钟级):tick 间隔 1 分钟,60 个 slot,覆盖 0~59 分钟的任务。
- 第三层(小时级):tick 间隔 1 小时,24 个 slot,覆盖 0~23 小时的任务。
当一个任务被触发时,Master 节点不会直接在本地执行,而是将任务信息序列化后发送到 Kafka 的指定 topic。Worker 节点消费这些消息并执行:
// Master 端:任务触发后投递到 Kafka
TaskDispatchEvent event = new TaskDispatchEvent();
event.setTaskId(task.getId());
event.setShardIndex(shardIndex);
event.setTotalShards(totalShards);
event.setParams(task.getParams());
kafkaTemplate.send("scheduler-dispatch", task.getId(), JSON.toJSONString(event));
// Worker 端:消费并执行
@KafkaListener(topics = "scheduler-dispatch", groupId = "scheduler-workers")
public void onDispatch(ConsumerRecord<String, String> record) {
TaskDispatchEvent event = JSON.parseObject(record.value(), TaskDispatchEvent.class);
taskExecutor.submit(() -> {
TaskResult result = executeTask(event);
reportResult(result); // 结果回写
});
}
这种 Master-Scheduler + Kafka + Worker 的三段式架构带来了几个重要优势:
- 执行解耦:Master 只负责调度决策,Worker 只负责执行,两者通过 Kafka 解耦,可以独立扩缩容。
- 削峰填谷:凌晨 2 点集中触发的结算任务不会瞬间压垮 Worker,Kafka 天然提供了消息缓冲能力。
- 故障隔离:如果某个 Worker 执行某个任务时 OOM 或死循环,不影响 Master 和其他 Worker 的正常工作。
四、K8s 容器化部署与弹性伸缩
在部署层面,我们将 Scheduler 和 Worker 分别部署为 Kubernetes Deployment,并利用 HPA(Horizontal Pod Autoscaler)实现弹性伸缩:
# Worker HPA 配置
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
name: scheduler-worker
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: scheduler-worker
minReplicas: 3
maxReplicas: 20
metrics:
- type: Pods
pods:
metric:
name: kafka_consumer_lag
target:
type: AverageValue
averageValue: 100
这里我们使用 Kafka Consumer Lag 作为 Worker 扩缩容的核心指标,比 CPU 和内存利用率更能反映实际负载。通过自定义 Metrics Server 采集 consumer lag,当 lag 超过阈值时自动扩容 Worker 副本数。
整体上线后,系统的关键指标有了显著提升:
- 任务调度延迟从 Quartz 时代的 30s 降至 2s 以内。
- 千万级创意审核任务从串行执行 4 小时缩减到 20 分钟(10 分片并行)。
- 单节点故障恢复时间 从分钟级降到秒级(ZK Session 超时 8s + Kafka rebalance)。
这套系统在百度广告线上稳定运行超过两年,日均处理任务量从初期的 50 万增长到 300 万+,期间没有出现过因调度系统自身故障导致的任务丢失或重复执行。