欢迎光临
我们一直在努力

分布式任务调度+多级缓存系统:从零实现一个自研调度框架

一、为什么做这个项目?

在实习和校招面试中,经常会被问到:“你做过分布式系统吗?” 大多数校招简历上的项目,要么是基于某个开源框架(如 XXL-JOB)做二次开发,要么是简单的 CRUD 业务系统。我想做一个真正从零开始的分布式系统——不是为了用某个框架,而是为了理解分布式调度背后的每一个细节。

这个项目是我对“分布式调度”这个经典问题的回答:调度中心如何管理执行器?任务如何分片?节点宕机了怎么办?任务失败了怎么重试?

带着这些问题,我一步步搭建了这个系统。

二、整体架构

系统采用经典的调度中心(Scheduler)+ 执行器(Worker) 架构,两者通过基于 Netty 的长连接通信。

模块职责

scheduler-common

公共协议、数据模型、JSON 工具

scheduler-core

调度中心服务端(8080)+ 管理 HTTP 接口(8081)

scheduler-worker

Worker 客户端,连接调度中心并执行任务

三、核心功能详解

1. 自定义通信协议:从底层解决粘包问题

  • TCP 是流式协议,消息边界是一个必须解决的问题。我没有用 Netty 自带的 LengthFieldBasedFrameDecoder,而是自己实现了一套二进制协议:
  • [魔数 4B][版本 1B][类型 1B][状态 1B][长度 4B][Body]
  • 消息类型:任务请求、任务响应、心跳、Worker 注册、缓存迁移
  • 长度字段:解决粘包/半包的核心

  • 解码器核心逻辑:

    // 1. 检查可读字节是否足够解析头部
    if (in.readableBytes() < HEADER_LENGTH) return;
    // 2. 标记读位置,以便回滚
    in.markReaderIndex();
    // 3. 读取魔数、版本、类型、状态…
    // 4. 读取 Body 长度
    int length = in.readInt();
    // 5. 如果 Body 不完整,重置读位置,等待下一次数据
    if (in.readableBytes() < length) {
    in.resetReaderIndex();
    return;
    }
    // 6. 读取完整 Body,解析为 Message

2. Worker 注册与心跳:Worker 如何被调度中心感知

  • Worker 启动后向调度中心注册(workerId)
  • 每 30 秒发送心跳;服务端 60 秒无心跳则断开并移除 Worker
  • 断连、异常、超时时自动从路由表移除
  • 核心逻辑:

    private void startHeartbeat() {
    Thread heartbeatThread = new Thread(() -> {
    while (channel != null && channel.isActive()) {
    Thread.sleep(30000);
    channel.writeAndFlush(Message.heartbeat());
    }
    });
    heartbeatThread.setDaemon(true);
    heartbeatThread.start();
    }

  • 这里有一个容易被忽视的问题:为什么是 60 秒而不是 30 秒? 因为网络抖动可能导致一个心跳包丢失。60 秒意味着允许连续丢失 1-2 个心跳包,避免了网络抖动导致的误判。

3. 分片任务调度:如何让 3 台机器并行处理 10 个任务

  • 定时调度器每 30 秒触发一次,将 10 个分片 任务下发
  • 使用 一致性哈希(每节点 150 个虚拟节点)把分片路由到对应 Worker
  • 一致性 Hash 路由:

    public class ConsistentHashRouter {
    private final SortedMap<Integer, String> hashRing = new TreeMap<>();
    private final int virtualNodeCount = 150;

    public String route(String key) {
    int hash = hash(key);
    if (!hashRing.containsKey(hash)) {
    // 沿环顺时针寻找最近的虚拟节点
    SortedMap<Integer, String> tailMap = hashRing.tailMap(hash);
    hash = tailMap.isEmpty() ? hashRing.firstKey() : tailMap.firstKey();
    }
    return hashRing.get(hash);
    }
    }

  • 为什么用一致性 Hash? 如果直接用 Hash(key)% nodeCount,当节点数量变化时,几乎所有任务的映射关系都会改变。一致性 Hash 通过引入虚拟节点,将节点变化的影响降到最低——只有少量任务需要迁移。

4. 任务执行与幂等

  • Worker 收到任务后解析 JobContext,执行业务逻辑(当前为模拟:读用户数据 + sleep 500ms)
  • 每个任务携带 JobContext,包含:

     shardingTotal:总分片数(如 10)

     shardingItem:当前分片序号(如 0-9)

     taskId:全局唯一任务 ID(用于幂等)

        Worker 收到任务后,根据 shardingItem 决定处理哪部分数据。这就是“分片”的本质:同一个任务的不同部分被分配到不同机器上并行执行。

  • 执行结果通过 ExecutionResult 回传调度中心

5. 高可用机制:节点挂了怎么办

  • 心跳检测 + 故障转移是系统高可用的核心。

    当调度中心检测到 Worker 心跳超时后,会执行以下流程:

    1. 从注册表中移除该 Worker
    2. 将该 Worker 未完成的任务重新放回重试队列
    3. 重建一致性 Hash 环
    4. 重试调度器将任务分配给其他健康的 Worker

  • 关键设计:removeWorker 方法中,任务转移和缓存迁移是独立且互补的:

    // 1. 任务转移
    List<JobContext> orphanTasks = runningTasks.remove(workerId);
    // … 将任务放入重试队列

    // 2. 缓存迁移
    List<String> hotKeys = CacheMigrationService.getHotKeys(workerId);
    // … 将热点 Key 迁移到其他 Worker

    这样,即使一个节点突然宕机,它正在执行的任务不会丢失,它的缓存也不会“失温”。

6. 超时控制与重试:任务卡住了怎么办?

每个任务下发时都会指定超时时间(默认 30 秒)。调度中心使用 ScheduledExecutorService 为每个任务创建独立的超时定时器:

ScheduledFuture<?> timeoutFuture = timeoutExecutor.schedule(() -> {
JobContext timedOutJob = pendingJobs.remove(jobId);
if (timedOutJob != null) {
// 标记超时,放入重试队列
retryQueue.offer(timedOutJob);
}
}, timeout, TimeUnit.SECONDS);

超时后,任务进入重试队列,重试调度器每 10 秒消费一次。

  • 重试策略:最多重试 3 次;采用指数退避(重试间隔逐渐增加)。
  • 当任务成功执行并返回结果后,onJobCompleted 会做三件事:取消超时定时器(timeoutFuture.cancel(false));从 pendingJobs 和 runningTasks 中移除任务;更新数据库状态

7. 任务幂等:怎么保证任务不重复执行?

在分布式环境中,网络抖动可能导致调度中心重发同一个任务。我的做法是:

数据库唯一索引:

CREATE TABLE schedule_job (
job_id BIGINT PRIMARY KEY,
task_id VARCHAR(128) UNIQUE, — 唯一索引

);

Worker 侧幂等检查:

private final Set<String> executedTasks = ConcurrentHashMap.newKeySet();

// 执行前检查
if (executedTasks.contains(job.getTaskId())) {
return ExecutionResult.success(job.getJobId(), job.getTaskId(), "Already executed");
}
// 执行成功后记录
executedTasks.add(job.getTaskId());

这样,即使调度中心因网络超时重发了同一个任务,Worker 也能识别并跳过。


8. 多级缓存:怎么加速数据访问?

Worker 在执行任务时可能需要读取数据。我实现了两级缓存:

层级实现过期时间特点
L1 本地缓存 Caffeine 60 秒 微秒级访问,无网络开销
L2 分布式缓存 Redis 300 秒 跨节点共享,毫秒级访问

防护机制:

  • 缓存穿透:Guava 布隆过滤器 + 空值缓存("NULL")

  • 缓存击穿:按 key 的互斥锁 + double-check

  • 缓存雪崩:Redis 过期时间加随机偏移(300s + 0~60s)

分片感知预热:调度中心下发任务时,会携带 preloadKeys,Worker 收到任务后提前加载这些数据到本地缓存。这样,任务真正执行时,数据已经在本地了。


9. 缓存迁移:节点宕机后,缓存怎么办?

当 Worker A 宕机时,调度中心会:

  • 从 CacheMigrationService 取出 Worker A 负责的热点 key 列表

  • 选择另一个健康的 Worker B

  • 发送 TYPE_CACHE_MIGRATE 消息,携带热点 key 列表

  • Worker B 收到后,从 Redis 批量加载这些 key 到本地 Caffeine

  • 这样,Worker A 的任务被转移到 Worker B 后,Worker B 的本地缓存已经“预热”好了,不会因为缓存缺失而产生大量数据库查询。


    10. 监控与可观测性

    系统提供两个简单的 HTTP 管理接口(端口 8081):

    • GET /workers:查看当前在线的 Worker 列表

    • GET /jobs/pending:查看待执行任务数量

    四、技术栈总结

    分类技术用途
    语言 Java 17 核心开发语言
    网络通信 Netty 4.1 长连接、Reactor 线程模型
    本地缓存 Caffeine 3.1 L1 缓存
    分布式缓存 Redis + Jedis 4.4 L2 缓存
    布隆过滤器 Guava 33.0 缓存穿透防护
    数据库 MySQL 8.0 + HikariCP 任务持久化
    构建工具 Maven 多模块管理
    测试 JUnit 4 编解码器单元测试

    五、项目收获与反思

    做对了什么

  • 从底层开始:没有直接使用 XXL-JOB 或 Quartz,而是从 Netty 开始构建。这让我真正理解了分布式系统的核心问题——网络通信、服务发现、故障转移。

  • 有数据验证:每个功能都做了压测或手动验证,不是“写完就结束了”。

  • 完整的错误处理:心跳超时、任务超时、节点宕机——每个“出问题”的场景都有应对逻辑。

  • 踩过的坑

  • Netty 粘包:最开始直接用 LengthFieldBasedFrameDecoder,但觉得没理解原理,于是自己写了一个 MessageDecoder。写完才真正明白“半包”是什么意思。

  • 一致性 Hash 的虚拟节点:一开始没有虚拟节点,节点数量变化时任务分配极度不均衡。加上 150 个虚拟节点后,分布才均匀。

  • 任务重复执行:调度中心重试时,Worker 会重复执行同一个任务。加 task_id 唯一索引和本地幂等检查后解决。


  • 六、后续计划

    • 时间轮调度器:替换当前的 ScheduledExecutorService,支持更精确的任务调度

    • Raft 一致性:实现调度中心集群,解决调度中心单点故障问题

    • 多路由策略:增加轮询、随机、加权等路由策略

    • 监控面板:提供更完善的任务执行历史和统计功能


    这个项目让我对分布式系统有了真正的理解——不是“知道某个框架怎么用”,而是“知道某个问题该怎么解决”。

    项目代码已开源:GitHub – yg33568/distributed-scheduler · GitHub

    如果你也在学习分布式系统,欢迎一起交流。

    赞(0)
    未经允许不得转载:171主机测评 » 分布式任务调度+多级缓存系统:从零实现一个自研调度框架
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址