一、为什么需要自研调度系统

在百度广告业务线,我们有上千个定时任务需要调度执行,涵盖广告创意审核、账户预算检查、投放数据聚合、报表生成等多个业务场景。初期团队使用的是基于 Quartz 集群的方案,但随着业务规模增长,这套方案逐渐暴露出几个关键问题:

经过技术选型评估,我们最终决定基于 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();
    }
}

这里有几个关键设计决策:

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 但做了分布式适配:

当一个任务被触发时,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 的三段式架构带来了几个重要优势:

四、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 副本数。

整体上线后,系统的关键指标有了显著提升:

这套系统在百度广告线上稳定运行超过两年,日均处理任务量从初期的 50 万增长到 300 万+,期间没有出现过因调度系统自身故障导致的任务丢失或重复执行。