从0到1搭建分布式定时任务平台:XXL-Job 原理拆解+手写精简版

引言

上个月的月结对账又双叒出问题了。凌晨 2 点的对账任务跑了 40 分钟,账单数据多了一倍——排查发现,两个订单服务实例各跑了一遍对账。单机时代用 @Scheduled 的代码,在部署第二个实例的那天起就埋下了这颗雷:没有分布式锁、没有分片、没有失败重试、没有告警,任务挂了只能靠第二天报表对不上才发现。

我们要的能力其实很明确:

需求@Scheduled 的现状期望
多实例不重复执行❌ 每个实例都会跑集群下同一任务只跑一次
大数据量提速❌ 单线程慢分片并行:10 个实例各处理 1/10 数据
任务失败感知❌ 日志里默默报错自动重试 + 告警通知
动态调整执行时间❌ 改 cron 要重启控制台改完立即生效
执行记录可查❌ 无记录每次执行的耗时/结果/日志留痕
手动触发/停止❌ 不支持控制台一键执行、终止

这些需求拼在一起,就是一个"分布式定时任务平台"。业界答案是 XXL-Job——它优雅地解决了上面所有问题,而且实现思路并不神秘。这篇文章先拆解 XXL-Job 的核心架构(调度中心做什么、执行器做什么、双方怎么协作),然后手写一个精简版:Quartz 调度 + HTTP 回调执行 + 数据库日志 + 失败重试 + 邮件告警,几百行代码跑通分布式任务平台的核心闭环。


一、XXL-Job 核心架构拆解

1.1 整体架构:调度与执行分离

XXL-Job 的架构核心是调度中心和执行器的分离——这是理解它所有设计的钥匙:

┌─────────────────────────────────────────────────────────────┐
│                     调度中心 (admin)                          │
│                                                             │
│  ┌──────────┐   ┌──────────────┐   ┌───────────────────┐   │
│  │ 执行器管理 │   │ 任务管理      │   │ 调度模块           │   │
│  │ (注册表)  │   │ (jobinfo 表)  │   │ (quartz + 线程池)  │   │
│  └──────────┘   └──────────────┘   └───────────────────┘   │
│  ┌──────────┐   ┌──────────────┐                           │
│  │ 日志管理  │   │ 告警通知      │                           │
│  │ (log 表) │   │ (邮件/钉钉)   │                           │
│  └──────────┘   └──────────────┘                           │
└─────────────────────────────────────────────────────────────┘
        │ ①注册(自动发现)          ▲ ③调度触发(HTTP REST)
        │                         │
        ▼ ②心跳续约               │
┌─────────────────────────────────────────────────────────────┐
│                  执行器集群 (executor)                        │
│                                                             │
│  ┌────────────────┐   ┌────────────────┐                   │
│  │ 执行器实例-1     │   │ 执行器实例-2     │                   │
│  │ 内嵌Server      │   │ 内嵌Server      │                   │
│  │ ├注册线程        │   │ ├注册线程        │                   │
│  │ ├回调线程        │   │ ├回调线程        │                   │
│  │ └JobHandler×N  │   │ └JobHandler×N  │                   │
│  └────────────────┘   └────────────────┘                   │
└─────────────────────────────────────────────────────────────┘

角色分工

角色职责类比
调度中心任务配置管理、到点触发、选择执行器实例、收集结果、失败重试、告警"指挥官":决定谁、什么时间、干什么
执行器接收调度请求、在自己的线程池里跑 JobHandler、上报执行结果"士兵":只管执行,不管什么时候执行

这个分离解决了什么:任务逻辑在业务应用里(执行器),调度能力集中管理(调度中心)。业务应用不关心调度细节,调度中心不关心任务内容——互不侵入。同时调度中心是无状态可水平扩展的(配合 DB 锁),执行器天然是集群的,两端都不存在单点瓶颈。

1.2 关键机制一:执行器自动注册

执行器启动后如何被调度中心发现?XXL-Job 的方案是心跳注册 + 表存储 + 定时清理

① 执行器启动:内嵌一个 HTTP Server(默认 9999 端口)
② 注册线程每 30 秒向调度中心发送注册请求:
   POST /api/registry  { registryGroup, registryKey(应用名), registryValue(ip:port) }
③ 调度中心写入 job_group / registry 表
④ 调度中心的 JobRegistryMonitorHelper 每 30 秒扫描:
   90 秒没心跳的实例 → 判定死亡,从注册表删除
⑤ 执行器优雅停机:主动发送 registryRemove,立即摘除

对应的表结构(简化):

-- 执行器实例注册表
CREATE TABLE xxl_job_registry (
    registry_group  varchar(50) NOT NULL,   -- 分组(EXECUTOR)
    registry_key    varchar(64)  NOT NULL,   -- 执行器 appName
    registry_value  varchar(255) NOT NULL,   -- 实例地址 ip:port
    update_time     datetime NOT NULL,       -- 最后心跳时间(判活依据)
    PRIMARY KEY (registry_group, registry_key, registry_value)
);

-- 执行器分组(任务路由的目标)
CREATE TABLE xxl_job_group (
    id          int PRIMARY KEY AUTO_INCREMENT,
    app_name    varchar(64) NOT NULL,        -- 执行器标识
    title       varchar(12) NOT NULL,        -- 展示名
    address_type tinyint NOT NULL,           -- 0自动注册 1手动录入
    address_list text                         -- 实例地址列表
);

设计要点:注册表的主键是 (group, key, value) 三列联合——同一应用的多实例是"多行"而不是"覆盖",心跳就是 update_time 的刷新,判活就是 update_time > now() - 90s 的过滤。用最朴素的表结构实现了服务发现,没有引入 Zookeeper/Nacos 的复杂度。

1.3 关键机制二:Quartz 集群调度 + DB 行锁

调度中心怎么保证"到点触发一次且只触发一次",尤其在调度中心自己也部署多实例的情况下?答案是 Quartz 集群模式 + 数据库行锁

调度中心实例-1 ─┐
调度中心实例-2 ─┼── 共享同一个数据库(Quqrtz 表 + xxl_job 表)
调度中心实例-3 ─┘

Quartz 集群行为:
① 每个实例的 quartz 定时扫 qrtz_triggers 表找"到点待触发的 trigger"
② 执行前先 SELECT ... FOR UPDATE 锁住 qrtz_locks 表的 STATE_ACCESS 行
③ 拿到锁的实例触发本次调度,其他实例本次扫空
④ 释放锁,下一轮继续竞争

→ 效果:多实例部署时,同一 trigger 只会被一个实例触发

关键代码路径(简化):XXL-Job 在 Quartz 触发后的 JobThread 里做了一次二次确认——先查任务信息、路由策略选执行器实例,再 HTTP 调用执行器。触发链路:

Quartz 触发(集群保证一次)
    ↓
XxlJobTrigger.trigger(jobId)
    ├─ 查任务配置(jobinfo)
    ├─ 路由策略选出目标执行器实例
    │   (第一个/轮询/随机/一致性HASH/分片广播/故障转移...)
    └─ ExecutorBizClient.run(...)
        POST http://executor:9999/run
        { jobId, executorParams, shardIndex, shardTotal, logId }

为什么二次调度不用 MQ? 这是很多人问的问题。XXL-Job 用 HTTP 直连触发而不是发消息,换来的是调用语义的确定性:调度失败当场感知(连接超时/业务拒绝),可以立即记录日志并走重试逻辑;MQ 的投递语义(at-least-once + 消费延迟)反而让"这次到底触发没有"变得难判断。调度场景要的是"确定的触发"而不是"最终会到达"

1.4 关键机制三:路由策略与分片广播

调度中心拿到执行器实例列表后,如何选择目标?8 种路由策略各有适用场景:

策略行为适用场景
第一个固定选第一个实例单实例执行即可,无需负载
轮询依次分配无状态任务均摊负载
随机随机选简单均摊
一致性 HASH按 jobId 哈希到实例同任务倾向固定实例(利用本地缓存)
故障转移心跳探测选第一个活着的实例高可用优先
忙碌转移逐个探测,选不忙的实例任务执行时间长,避免堆积
分片广播全量实例都触发,带 shardIndex/shardTotal大数据量并行处理
分片+故障转移分片基础上排除死实例大数据量 + 高可用

分片广播是 XXL-Job 最有价值的策略,它把"任务怎么拆"交给了业务代码:

/**
 * 分片任务示例:100 万条对账数据,10 个实例各处理 1/10
 * shardTotal=10:调度中心广播时告知总分片数
 * shardIndex:每个实例收到的编号(0~9)
 */
@XxlJob("reconciliationShardingJob")
public void reconciliationShardingJob() {
    int shardIndex = XxlJobHelper.getShardIndex();   // 当前实例的分片号
    int shardTotal = XxlJobHelper.getShardTotal();   // 总分片数

    // 按 id 取模分片:每个实例只处理自己分片的数据
    List<Bill> bills = billMapper.selectUnsettled(shardIndex, shardTotal);
    // SQL: WHERE MOD(id, #{shardTotal}) = #{shardIndex} AND status='UNSETTLED'

    for (Bill bill : bills) {
        reconcile(bill);
    }
    XxlJobHelper.handleSuccess("分片" + shardIndex + "处理" + bills.size() + "条");
}

注意分片和负载均衡的本质区别:轮询是"这次任务谁跑",分片是"这次任务大家各跑一段"。单次任务数据量大到单实例处理不完时,分片是唯一解。

1.5 关键机制四:执行结果回调与失败重试

执行器跑完任务不直接"返回结果"就完事——异步回调是它的另一个关键设计:

① 调度中心触发任务 → 创建 logId → HTTP 调执行器 /run(异步提交到执行器线程池)
② /run 立即返回 200(表示"已接收",不是"已执行完")
③ 执行器线程池里跑 JobHandler(可能跑几分钟)
④ 执行完成后,回调线程 POST 调度中心 /api/callback { logId, code, msg }
⑤ 调度中心更新日志状态、判断是否重试、是否告警

为什么要异步回调而不是同步等待执行完? 任务可能跑 30 分钟,HTTP 同步等待既占用连接又不可靠(超时后结果丢失)。异步回调把"接收任务"和"执行任务"解耦,结果通过独立通道可靠送达。

失败重试的关键规则(源码里的判定逻辑):

配置含义注意点
失败重试次数 N回调状态为失败时,重新调度,最多再触发 N 次重试是"重新走完整调度链路":重新选实例、重新路由,不是原实例原地重跑
重试与超时单次执行超过"任务超时时间"被判定失败 → 同样触发重试超时 + 重试要一起评估,避免 10 分钟任务 × 5 次重试放大故障
重试与告警每次失败都会触发告警重试次数别贪多,一般 1~3 次足够

1.6 XXL-Job 架构总结:一张图看懂全链路

业务接入        调度中心                    执行器集群
──────         ─────────                  ─────────
@XxlJob        quartz(集群锁) → 触发
声明Handler                     ──路由策略──►  选定实例
                                              内嵌Server接收
               logId 创建                       ↓
                              ◄──注册/心跳──  注册线程(30s)
                               执行任务
               ◄───────────────────────────  异步回调结果
               更新日志
               ├─失败→重试(重新调度≤N次)
               └─失败→告警(邮件/钉钉)

到这里,XXL-Job 的核心秘密已经拆完:执行器注册(表+心跳)、集群调度(Quartz+行锁)、路由策略(含分片广播)、异步回调、失败重试、告警。下面手写的精简版,就是把这条链路以最小代码量复刻一遍。


二、手写精简版:整体设计

2.1 目标与边界

复刻的核心链路:注册发现 → 到点调度 → 路由选择 → HTTP 触发 → 异步执行 → 结果回调 → 日志落库 → 失败重试 → 邮件告警。精简掉的:控制台前端(用表+接口代替)、Token 鉴权、GLUE 模式、子任务依赖。

2.2 项目结构

mini-job/
├── mini-job-admin/                 # 调度中心
│   └── src/main/java/space/jiangyi/minijob/admin/
│       ├── MiniJobAdminApplication.java
│       ├── config/
│       │   ├── MailConfig.java            # 邮件配置
│       │   └── AdminProperties.java       # 端口/超时/重试配置
│       ├── core/
│       │   ├── JobRegistryService.java    # 执行器注册/心跳/摘除/判活
│       │   ├── JobScheduler.java          # 调度入口:到点触发 + 路由
│       │   ├── RouteStrategy.java         # 路由策略(轮询/随机/第一个/分片广播)
│       │   ├── ExecutorClient.java        # HTTP 调用执行器
│       │   ├── CallbackService.java       # 接收执行结果回调
│       │   ├── RetryService.java          # 失败重试调度
│       │   └── AlertService.java          # 邮件告警
│       └── web/
│           └── AdminController.java       # 注册/回调 API + 任务管理接口
├── mini-job-executor/              # 执行器(嵌入到业务应用)
│   └── src/main/java/space/jiangyi/minijob/executor/
│       ├── MiniJobExecutor.java           # 执行器核心:内嵌Server+注册线程
│       ├── ExecutorProperties.java        # admin地址/appName/端口
│       └── TaskContext.java               # 任务上下文
└── demo-app/                       # 演示业务应用
    └── src/main/java/space/jiangyi/minijob/demo/
        ├── DemoApplication.java
        └── tasks/
            ├── SimpleTask.java            # 普通任务
            └── ShardTask.java             # 分片任务

2.3 数据库设计(4 张表)

-- ① 任务定义表
CREATE TABLE mj_job_info (
    id              bigint PRIMARY KEY AUTO_INCREMENT,
    job_name        varchar(64)  NOT NULL,
    cron            varchar(32)  NOT NULL,          -- cron 表达式
    handler_name    varchar(64)  NOT NULL,          -- 执行器里的任务处理器名
    executor_params varchar(255) DEFAULT '',        -- 任务参数
    route_type      varchar(16)  NOT NULL,          -- ROUND/RANDOM/FIRST/SHARDING
    shard_total     int          DEFAULT 1,         -- 分片总数(分片广播用)
    misfire_policy  varchar(16)  DEFAULT 'FIRE_NOW',-- 错过触发策略
    retry_count     int          DEFAULT 0,         -- 失败重试次数
    timeout_seconds int          DEFAULT 300,       -- 单次执行超时
    trigger_status  tinyint      DEFAULT 0,         -- 0停 1启动
    version         int          DEFAULT 0,         -- 乐观锁(防止并发改配置)
    create_time     datetime DEFAULT CURRENT_TIMESTAMP,
    update_time     datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);

-- ② 执行日志表(每次触发一条)
CREATE TABLE mj_job_log (
    id              bigint PRIMARY KEY AUTO_INCREMENT,
    job_id          bigint NOT NULL,
    handler_name    varchar(64) NOT NULL,
    executor_addr   varchar(64),                    -- 实际执行实例
    trigger_time    datetime NOT NULL,
    trigger_code    tinyint DEFAULT 0,              -- 0触发中 1成功 2失败
    handle_code     tinyint DEFAULT 0,              -- 0执行中 1成功 2失败
    handle_msg      text,                           -- 执行结果/异常摘要
    retry_index     int DEFAULT 0,                  -- 第几次尝试
    timeout_deadline datetime,                      -- 超时判定线
    KEY idx_job_trigger (job_id, trigger_time)
);

-- ③ 执行器注册表
CREATE TABLE mj_registry (
    app_name     varchar(64) NOT NULL,
    instance     varchar(64) NOT NULL,              -- ip:port
    update_time  datetime NOT NULL,
    PRIMARY KEY (app_name, instance)
);

-- ④ 路由轮询游标表(轮询策略记录"上次选到谁")
CREATE TABLE mj_route_cursor (
    job_id       bigint PRIMARY KEY,
    cursor       int DEFAULT 0
);

日志表的 trigger_codehandle_code 分开,对应"触发是否成功"和"执行是否成功"两个阶段——这是理解异步回调模型的关键:调度中心先记录触发结果,执行器回调后再补执行结果。


三、手写精简版:调度中心实现

3.1 执行器注册与判活

/**
 * 执行器注册中心:心跳注册 + 定时判活
 * 逻辑对齐 XXL-Job:30 秒心跳、90 秒超时判死
 */
@Service
@RequiredArgsConstructor
public class JobRegistryService {

    private final JdbcTemplate jdbc;
    private final ScheduledExecutorService cleaner =
            Executors.newSingleThreadScheduledExecutor();

    /** 执行器调用:注册实例(每次心跳都会调,等效幂等 upsert) */
    public void registry(String appName, String instance) {
        jdbc.update("""
            INSERT INTO mj_registry(app_name, instance, update_time)
            VALUES(?, ?, NOW())
            ON DUPLICATE KEY UPDATE update_time = NOW()
            """, appName, instance);
    }

    /** 执行器优雅停机调用:主动摘除 */
    public void registryRemove(String appName, String instance) {
        jdbc.update("DELETE FROM mj_registry WHERE app_name=? AND instance=?",
                appName, instance);
    }

    /** 查询某执行器组的存活实例列表(路由的输入) */
    public List<String> aliveInstances(String appName) {
        // 90 秒内有心跳的实例视为存活
        return jdbc.queryForList("""
            SELECT instance FROM mj_registry
            WHERE app_name = ? AND update_time > DATE_SUB(NOW(), INTERVAL 90 SECOND)
            ORDER BY instance
            """, String.class, appName);
    }

    /** 每 30 秒清理死实例:DB 兜底,防止"假活"记录堆积 */
    @PostConstruct
    public void startCleaner() {
        cleaner.scheduleAtFixedRate(() -> {
            try {
                jdbc.update("""
                    DELETE FROM mj_registry
                    WHERE update_time < DATE_SUB(NOW(), INTERVAL 90 SECOND)
                    """);
            } catch (Exception e) {
                log.warn("registry cleaner error", e);
            }
        }, 30, 30, TimeUnit.SECONDS);
    }
}

3.2 调度入口:到点触发 + 路由 + 触发结果落库

/**
 * 调度器:扫描到期任务 → 路由选实例 → HTTP 触发 → 记录触发日志
 * 集群部署时靠 DB 乐观锁抢占"触发权"(对齐 XXL-Job 的行锁思想)
 */
@Service
@RequiredArgsConstructor
@Slf4j
public class JobScheduler {

    private final JdbcTemplate jdbc;
    private final JobRegistryService registryService;
    private final ExecutorClient executorClient;
    private final RetryService retryService;

    /**
     * 每秒扫描一次到期任务(生产建议用 Quartz 集群模式替换这里的扫描)
     * 精简版用"秒级扫描 + 乐观锁抢占"实现集群下不重复触发
     */
    @Scheduled(fixedDelay = 1000)
    public void scanAndTrigger() {
        // 1. 找出到期任务:启用中 && cron 到点(简化:用 next_fire_time 判断)
        List<Map<String, Object>> dueJobs = jdbc.queryForList("""
            SELECT id, job_name, handler_name, executor_params, route_type,
                   shard_total, retry_count, timeout_seconds, version
            FROM mj_job_info
            WHERE trigger_status = 1
              AND next_fire_time <= NOW()
            """);

        for (Map<String, Object> job : dueJobs) {
            long jobId = ((Number) job.get("id")).longValue();
            int version = ((Number) job.get("version")).intValue();
            try {
                triggerOnce(jobId, job, version);
            } catch (Exception e) {
                log.error("trigger job {} error", jobId, e);
            }
        }
    }

    /**
     * 触发一次任务:乐观锁抢占触发权 + 推进 next_fire_time
     * 只有一个调度中心实例能 UPDATE 成功 → 集群下不重复触发
     */
    private void triggerOnce(long jobId, Map<String, Object> job, int version) {
        String cron = (String) job.get("cron");
        Date nextFire = new CronExpression(cron).getNextValidTimeAfter(new Date());

        // 2. 乐观锁抢占:UPDATE 成功才继续触发
        int updated = jdbc.update("""
            UPDATE mj_job_info
            SET next_fire_time = ?, version = version + 1
            WHERE id = ? AND version = ?
            """, nextFire, jobId, version);
        if (updated == 0) {
            return;   // 被其他调度中心实例抢走了
        }

        // 3. 路由选择执行器实例
        String routeType = (String) job.get("route_type");
        int shardTotal = ((Number) job.get("shard_total")).intValue();

        // 4. 创建执行日志(触发阶段)
        long logId = insertLog(jobId, job);

        // 5. 按策略触发
        if ("SHARDING".equals(routeType)) {
            // 分片广播:全量实例各自触发,带分片参数
            List<String> instances = registryService.aliveInstances(appNameOf(jobId));
            if (instances.isEmpty()) {
                markTriggerFail(logId, "无存活执行器实例");
                return;
            }
            for (int i = 0; i < instances.size(); i++) {
                doDispatch(jobId, job, logId, instances.get(i), i, shardTotal);
            }
        } else {
            // 普通路由:选一个实例
            String addr = route(routeType, jobId, appNameOf(jobId));
            if (addr == null) {
                markTriggerFail(logId, "无存活执行器实例");
                return;
            }
            doDispatch(jobId, job, logId, addr, 0, 1);
        }
    }

    /** 实际下发:异步 HTTP 调用,不等执行结果(结果走回调) */
    private void doDispatch(long jobId, Map<String, Object> job, long logId,
                            String addr, int shardIndex, int shardTotal) {
        jdbc.update("UPDATE mj_job_log SET executor_addr=? WHERE id=?",
                addr, logId);

        executorClient.run(addr, DispatchRequest.builder()
                .logId(logId)
                .jobId(jobId)
                .handlerName((String) job.get("handler_name"))
                .params((String) job.get("executor_params"))
                .shardIndex(shardIndex)
                .shardTotal(shardTotal)
                .timeoutSeconds(((Number) job.get("timeout_seconds")).intValue())
                .build()
        ).whenComplete((resp, err) -> {
            if (err != null || !resp.isSuccess()) {
                // 触发失败:连接不上/执行器拒绝 → 直接判失败,走重试
                markTriggerFail(logId,
                        err != null ? err.getMessage() : resp.getErrorMsg());
            }
            // 触发成功不在此处结束——等执行器异步回调 handle_result
        });
    }
}

这里的三个设计决策值得注意:

  1. 乐观锁抢占触发权替代 Quartz 集群行锁——精简版的扫描式调度,用 version 乐观锁达到同样的"集群不重复触发"效果;
  2. 触发与执行两阶段记录——trigger_code 记"有没有派下去",handle_code 记"有没有跑成功",失败重试的判定才准确;
  3. HTTP 异步触发——whenComplete 只处理"派发失败",不等待任务执行完成。

3.3 路由策略

/**
 * 路由策略:精简版实现 4 种(第一个/随机/轮询/分片广播在调度器里处理)
 * 扩展一致性 HASH/故障转移:对照 XXL-Job 的 ExecutorRouteStrategyEnum
 */
@Component
@RequiredArgsConstructor
public class RouteStrategy {

    private final JobRegistryService registryService;
    private final JdbcTemplate jdbc;
    private final Random random = new Random();

    public String route(String routeType, long jobId, String appName) {
        List<String> instances = registryService.aliveInstances(appName);
        if (instances.isEmpty()) {
            return null;
        }
        return switch (routeType) {
            case "FIRST"  -> instances.get(0);
            case "RANDOM" -> instances.get(random.nextInt(instances.size()));
            case "ROUND"  -> roundRobin(jobId, instances);
            default       -> instances.get(0);
        };
    }

    /** 轮询:游标存 DB,多调度中心实例下也一致 */
    private String roundRobin(long jobId, List<String> instances) {
        int cursor = jdbc.queryForObject(
                "SELECT cursor FROM mj_route_cursor WHERE job_id=? FOR UPDATE",
                Integer.class, jobId);
        String target = instances.get(cursor % instances.size());
        jdbc.update("""
            INSERT INTO mj_route_cursor(job_id, cursor) VALUES(?, ?)
            ON DUPLICATE KEY UPDATE cursor = ?
            """, jobId, cursor + 1, cursor + 1);
        return target;
    }
}

3.4 接收回调 + 失败重试 + 告警

/**
 * 回调处理:更新执行结果 → 判定重试 → 触发告警
 */
@Service
@RequiredArgsConstructor
@Slf4j
public class CallbackService {

    private final JdbcTemplate jdbc;
    private final RetryService retryService;
    private final AlertService alertService;

    /**
     * 执行器回调:POST /api/callback { logId, handleCode, handleMsg }
     */
    public void callback(long logId, int handleCode, String handleMsg) {
        int updated = jdbc.update("""
            UPDATE mj_job_log
            SET handle_code = ?, handle_msg = ?, handle_time = NOW()
            WHERE id = ? AND handle_code = 0
            """, handleCode, truncate(handleMsg, 2000), logId);
        if (updated == 0) {
            return;   // 已被回调过(超时判定线兜底也回调过),幂等保护
        }

        Map<String, Object> logRow = jdbc.queryForMap(
                "SELECT job_id, retry_index FROM mj_job_log WHERE id=?", logId);
        long jobId = ((Number) logRow.get("job_id")).longValue();
        int retryIndex = ((Number) logRow.get("retry_index")).intValue();

        if (handleCode == 1) {
            return;   // 执行成功,链路结束
        }

        // 执行失败:判断还能不能重试
        int maxRetry = jdbc.queryForObject(
                "SELECT retry_count FROM mj_job_info WHERE id=?", Integer.class, jobId);
        if (retryIndex < maxRetry) {
            retryService.scheduleRetry(jobId, retryIndex + 1);
        } else {
            // 重试用尽:最终失败,告警
            alertService.alertFinalFail(jobId, logId, handleMsg);
        }
    }
}

/**
 * 失败重试:延迟后重新走"路由+触发"链路(对齐 XXL-Job:重试=重新调度)
 * 精简版用延迟队列实现,XXL-Job 用"重新触发 quartz"实现,思想一致
 */
@Service
@RequiredArgsConstructor
@Slf4j
public class RetryService {

    private final JdbcTemplate jdbc;
    private final DelayQueue<RetryTask> queue = new DelayQueue<>();
    private Thread worker;

    @PostConstruct
    public void start() {
        worker = new Thread(() -> {
            while (true) {
                try {
                    RetryTask task = queue.take();
                    doRetry(task);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    return;
                } catch (Exception e) {
                    log.error("retry error", e);
                }
            }
        }, "mj-retry");
        worker.setDaemon(true);
        worker.start();
    }

    public void scheduleRetry(long jobId, int retryIndex) {
        // 重试间隔递增:1s → 2s → 4s(指数退避,避免故障期轰炸)
        long delayMs = (long) Math.pow(2, retryIndex - 1) * 1000;
        queue.offer(new RetryTask(jobId, retryIndex, System.currentTimeMillis() + delayMs));
    }

    private void doRetry(RetryTask task) {
        // 插入一条新的执行日志(retry_index 递增),重新走触发链路
        long newLogId = insertRetryLog(task.getJobId(), task.getRetryIndex());
        // ...复用 JobScheduler 的触发逻辑(省略:调 triggerOnce 的公共部分)
        log.info("job {} 第 {} 次重试已下发, logId={}",
                task.getJobId(), task.getRetryIndex(), newLogId);
    }
}

/**
 * 邮件告警:最终失败时通知(精简版纯邮件,XXL-Job 支持邮件+用户自定义扩展)
 */
@Service
@RequiredArgsConstructor
public class AlertService {

    private final JavaMailSender mailSender;
    private final AdminProperties props;
    private final JdbcTemplate jdbc;

    public void alertFinalFail(long jobId, long logId, String errorMsg) {
        Map<String, Object> job = jdbc.queryForMap(
                "SELECT job_name FROM mj_job_info WHERE id=?", jobId);

        String subject = String.format("[MiniJob告警] 任务最终失败: %s", job.get("job_name"));
        String content = String.format("""
                任务ID: %d
                日志ID: %d
                失败信息: %s
                时间: %s
                请及时登录控制台查看执行日志。
                """, jobId, logId, errorMsg, LocalDateTime.now());

        SimpleMailMessage mail = new SimpleMailMessage();
        mail.setFrom(props.getAlertFrom());
        mail.setTo(props.getAlertTo().split(","));
        mail.setSubject(subject);
        mail.setText(content);
        mailSender.send(mail);
    }
}

3.5 执行器:内嵌 Server + 注册线程 + 任务执行

/**
 * 执行器核心:业务应用引入 jar 后,一个 @Bean 搞定接入
 * 职责对齐 XXL-Job 执行器:
 *   ① 内嵌 HTTP Server 接收 /run
 *   ② 注册线程 30 秒心跳
 *   ③ 独立线程池执行任务(不占用业务线程)
 *   ④ 执行完异步回调调度中心
 */
@Slf4j
public class MiniJobExecutor {

    private final ExecutorProperties props;
    private final ExecutorService workerPool =
            Executors.newFixedThreadPool(16);          // 任务执行线程池
    private final ScheduledExecutorService registryThread =
            Executors.newSingleThreadScheduledExecutor();
    private final Map<String, Method> handlerMap = new ConcurrentHashMap<>();
    private final HttpClient httpClient = HttpClient.newHttpClient();
    private volatile boolean running = true;

    public MiniJobExecutor(ExecutorProperties props, Map<String, Object> handlerBeans) {
        this.props = props;
        // 扫描业务容器里的 @MiniJobTask Bean,注册 handler
        handlerBeans.forEach((beanName, bean) -> {
            for (Method m : bean.getClass().getDeclaredMethods()) {
                MiniJobTask anno = m.getAnnotation(MiniJobTask.class);
                if (anno != null) {
                    handlerMap.put(anno.value(), methodOf(bean, m));
                }
            }
        });
    }

    @PostConstruct
    public void start() {
        // ① 心跳注册线程:每 30 秒向调度中心注册
        registryThread.scheduleAtFixedRate(this::doRegistry, 0, 30, TimeUnit.SECONDS);

        // ② 内嵌 HTTP Server:接收调度触发(用 JDK HttpServer,零依赖)
        try {
            HttpServer server = HttpServer.create(
                    new InetSocketAddress(props.getPort()), 0);
            server.createContext("/run", this::handleRun);
            server.setExecutor(Executors.newFixedThreadPool(8));
            server.start();
            log.info("[MiniJob] executor started on port {}", props.getPort());
        } catch (IOException e) {
            throw new IllegalStateException("executor server start failed", e);
        }
    }

    /** 处理调度触发:校验 handler → 提交线程池 → 立即应答"已接收" */
    private void handleRun(HttpExchange exchange) {
        DispatchRequest req = parse(exchange);
        // 立即应答:200 表示"已接收",执行结果走异步回调
        respond(exchange, 200, "{\"code\":200}");

        workerPool.submit(() -> executeAndCallback(req));
    }

    /** 执行任务并回调结果(分片参数透传给业务方法) */
    private void executeAndCallback(DispatchRequest req) {
        TaskContext ctx = TaskContext.builder()
                .jobId(req.getJobId())
                .params(req.getParams())
                .shardIndex(req.getShardIndex())
                .shardTotal(req.getShardTotal())
                .build();
        int code; String msg;
        try {
            Method handler = handlerMap.get(req.getHandlerName());
            if (handler == null) {
                throw new IllegalStateException("handler 不存在: " + req.getHandlerName());
            }
            handler.invoke(handler.getBean(), ctx);
            code = 1; msg = "success";
        } catch (Exception e) {
            code = 2;
            msg = Throwables.getStackTraceAsString(e);
            log.error("[MiniJob] task {} execute failed", req.getHandlerName(), e);
        }
        // 异步回调调度中心
        callback(req.getLogId(), code, msg);
    }

    private void callback(long logId, int code, String msg) {
        String body = String.format(
                "{\"logId\":%d,\"handleCode\":%d,\"handleMsg\":%s}",
                logId, code, JSON.quote(truncate(msg, 2000)));
        HttpRequest request = HttpRequest.newBuilder()
                .uri(URI.create(props.getAdminUrl() + "/api/callback"))
                .header("Content-Type", "application/json")
                .POST(HttpRequest.BodyPublishers.ofString(body))
                .build();
        try {
            httpClient.send(request, HttpResponse.BodyHandlers.ofString());
        } catch (Exception e) {
            // 回调失败:本地重试队列兜底(省略:入内存队列延迟重投)
            log.error("[MiniJob] callback failed, logId={}", logId, e);
        }
    }

    private void doRegistry() {
        if (!running) return;
        String body = String.format(
                "{\"appName\":\"%s\",\"instance\":\"%s:%d\"}",
                props.getAppName(), localIp(), props.getPort());
        try {
            httpClient.send(HttpRequest.newBuilder()
                    .uri(URI.create(props.getAdminUrl() + "/api/registry"))
                    .header("Content-Type", "application/json")
                    .POST(HttpRequest.BodyPublishers.ofString(body))
                    .build(), HttpResponse.BodyHandlers.ofString());
        } catch (Exception e) {
            log.warn("[MiniJob] registry heartbeat failed: {}", e.getMessage());
        }
    }

    @PreDestroy
    public void stop() {
        running = false;
        // 优雅停机:主动摘除注册(对齐 XXL-Job registryRemove)
        try {
            httpClient.send(HttpRequest.newBuilder()
                    .uri(URI.create(props.getAdminUrl() + "/api/registryRemove"))
                    .header("Content-Type", "application/json")
                    .POST(HttpRequest.BodyPublishers.ofString(
                            String.format("{\"appName\":\"%s\",\"instance\":\"%s:%d\"}",
                                    props.getAppName(), localIp(), props.getPort())))
                    .build(), HttpResponse.BodyHandlers.ofString());
        } catch (Exception ignored) {
        }
        workerPool.shutdown();
    }
}

业务侧的使用方式——声明任务只需一个注解:

/**
 * 业务任务:普通任务 + 分片任务
 * 接入方只写业务逻辑,注册/调度/重试/告警全部由 MiniJob 接管
 */
@Component
public class SettlementTasks {

    @MiniJobTask("dailySettlement")
    public void dailySettlement(TaskContext ctx) {
        // ctx.getParams():调度中心传的任务参数
        settlementService.settle(LocalDate.now());
    }

    @MiniJobTask("reconciliationSharding")
    public void reconciliationSharding(TaskContext ctx) {
        // 分片任务:从上下文拿分片参数,处理自己分片的数据
        List<Bill> bills = billMapper.selectUnsettled(
                ctx.getShardIndex(), ctx.getShardTotal());
        bills.forEach(billService::reconcile);
    }
}

// application.yml
// mini-job:
//   admin-url: http://mini-job-admin:8080
//   app-name: order-service
//   port: 9999

四、验证:全链路跑通

4.1 验证步骤

# 1. 起两个业务实例(模拟集群)
java -jar demo-app.jar --server.port=8081
java -jar demo-app.jar --server.port=8082

# 2. 控制台(调 DB 接口)配置任务:
#    handler_name=reconciliationSharding, route_type=SHARDING, shard_total=2
#    cron=0 0 2 * * ?(每天凌晨 2 点)

# 3. 查注册表:两个实例已自动注册
SELECT * FROM mj_registry;
-- order-service | 10.0.0.1:9999 | 2024-08-22 10:00:30
-- order-service | 10.0.0.2:9999 | 2024-08-22 10:00:32

# 4. 手动触发,观察日志表
SELECT id, executor_addr, trigger_code, handle_code, retry_index, handle_msg
FROM mj_job_log ORDER BY id DESC LIMIT 5;
-- id=101 | 10.0.0.1:9999 | 1 | 1 | 0 | success(分片0,处理 50123 条)
-- id=102 | 10.0.0.2:9999 | 1 | 1 | 0 | success(分片1,处理 49877 条)

4.2 三类典型问题在精简版里的表现

问题场景精简版的行为依据机制
一个实例宕机心跳 90 秒超时被清理,路由不再选中它注册判活 + 路由回退
任务执行中实例被杀回调缺失 → 超时判定线(timeout_deadline)兜底标记失败 → 重试两阶段状态 + 超时兜底
调度中心重启乐观锁推进的 next_fire_time 已持久化,重启后继续正常调度DB 为唯一事实源
对账数据量翻倍shard_total 从 2 改 4,4 个实例并行处理分片广播

五、常见问题

5.1 精简版用扫描 + 乐观锁,和 XXL-Job 的 Quartz 集群差在哪?

功能上等价(集群不重复触发),差异在精度和成熟度:Quartz 集群用 DB 行锁 + trigger 表管理,秒级精度、misfire 策略完善、经过大量生产验证;精简版的秒级扫描在任务量上千后扫描效率下降,且 cron 解析、misfire(调度中心宕机期间错过的触发)要自己处理。生产建议直接用 Quartz 集群模式做触发层,本精简版的意义在于看清整条链路

5.2 回调丢了怎么办(执行成功但调度中心没收到)?

两层兜底:① 超时判定线——日志表里的 timeout_deadline,调度中心有个后台任务扫描"执行中但超过 deadline"的日志,标记为"未知失败"并走重试;② 执行器回调重试——回调失败入本地重试队列延迟重投。极端情况下(重试 + 回调都失败)会出现"实际成功但记录失败",重试导致任务跑两次——所以重试的任务必须幂等,这是分布式任务的铁律。

5.3 为什么调度中心不做成"推任务到 MQ"?

见 1.3 的分析:调度要的是"确定的触发"语义。MQ 引入投递延迟、消费不可控、至少一次的重复,把"这次触发成功没有"变成异步不可知。HTTP 直连的失败当场可见,重试和告警的判定链路才清晰。MQ 更适合任务"派发后状态无关"的场景,调度不属于这类。

5.4 XXL-Job 的任务为什么必须幂等?怎么做到?

三个来源都可能造成重复执行:重试机制、超时误判(实际在跑但被标失败再重试)、人工手动触发与调度触发撞车。幂等的常用手段:数据库唯一约束(对账任务以"对账日期"为唯一键,重复插入直接冲突跳过)、Redis 分布式锁(setnx 任务锁 + 执行完释放)、状态机(任务只允许从"待处理"到"处理中",已处理的直接跳过)。

5.5 分片广播时某个实例挂了,它负责的分片数据谁处理?

这是分片广播的固有缺口:shardIndex 是按触发瞬间的存活实例数分配的,实例中途挂掉,它的分片本次没人处理。应对方式:① 调度中心超时兜底发现"分片 2 的日志一直执行中"→ 告警人工介入;② 任务层面做补偿扫描——每次分片任务跑完后,由一个单实例任务(路由=第一个)扫描"遗漏未处理的数据"兜底;③ XXL-Job 的分片+故障转移策略在触发时排除死实例,但同样防不住"执行中挂掉"。

5.6 精简版和 XXL-Job 的差距清单(生产前必须补的)

差距点精简版现状生产要求
调度精度与 misfire秒级扫描,missed 触发简单处理Quartz 集群 + misfire 策略
安全API 无鉴权AccessToken / 双向认证
控制台无前端,靠接口任务 CRUD / 实时日志 tail / 用户权限
日志DB 存摘要执行日志文件 + 在线查看 + 滚动清理
依赖任务子任务串(任务 A 成功后触发 B)
高可用回调重试在内存,重启丢失回调持久化 + 调度中心无状态化验证

六、总结

XXL-Job 核心机制速查卡

┌──────────────┬────────────────────────────────────────────┐
│ 机制          │ 实现                                        │
├──────────────┼────────────────────────────────────────────┤
│ 服务发现      │ 执行器 30s 心跳注册 → 注册表 → 90s 判活清理    │
│ 集群调度      │ Quartz 集群模式 + DB 行锁(多 admin 不重复触发)│
│ 路由策略      │ 8 种:第一个/轮询/随机/一致性HASH/故障转移/     │
│              │ 忙碌转移/分片广播(+故障转移)                    │
│ 触发语义      │ HTTP 直连(确定的触发),不走 MQ                │
│ 执行模型      │ 异步两阶段:/run 立即应答 + 完成后独立回调       │
│ 失败处理      │ 重试=重新走完整调度链路(非原地重跑),≤N 次     │
│ 告警          │ 最终失败邮件/自定义扩展                        │
└──────────────┴────────────────────────────────────────────┘

手写精简版的代码量账

模块核心类行数(约)
注册中心JobRegistryService60
调度器JobScheduler + RouteStrategy180
回调/重试/告警CallbackService + RetryService + AlertService150
执行器MiniJobExecutor200
表结构 + 业务接入4 张表 + @MiniJobTask 注解100
合计 ~700 行

700 行跑通了分布式任务平台的核心闭环——不是说明 XXL-Job 简单,而是说明它的架构分层足够干净:每一层(发现/调度/执行/回调/重试/告警)的职责边界清晰,理解了边界,精简复刻就是水到渠成的事。

一句话

XXL-Job 的全部秘密就是"调度与执行分离"这六个字:调度中心用注册表+心跳管"谁活着",用 Quartz+行锁管"什么时候触发且只触发一次",用路由策略管"派给谁";执行器用内嵌 Server 接单、线程池干活、异步回调交差——剩下的重试、告警、日志,都是这条主链路上的状态判定。手写一遍的意义不是造轮子,而是把每个组件的"为什么存在"变成自己的知识。

给团队的建议

阶段建议
还在用 @Scheduled 多实例部署立即评估,重复执行是定时炸弹
需要分布式任务直接上 XXL-Job,不要自研(本精简版用于理解原理)
已用 XXL-Job大数据量任务改分片广播;所有重试任务检查幂等
任务量 > 千级关注 admin 扫描压力,评估升级 PowerJob(支持工作流/Map-Reduce 模型)
自研平台精简版基础上按 5.6 差距清单补齐,重点是安全与日志

互动话题:你们的定时任务用什么方案?@Scheduled 多实例重复执行的雷踩过吗?XXL-Job 的分片广播在你业务里单次最多分了多少片?评论区聊聊。


参考资料


标题:从0到1搭建分布式定时任务平台:XXL-Job 原理拆解+手写精简版
作者:jiangyi
地址:http://www.jiangyi.space/articles/2026/09/04/1788098721234.html
公众号:服务端技术精选
    评论
    0 评论
avatar

取消