从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_code 与 handle_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
});
}
}
这里的三个设计决策值得注意:
- 乐观锁抢占触发权替代 Quartz 集群行锁——精简版的扫描式调度,用
version乐观锁达到同样的"集群不重复触发"效果; - 触发与执行两阶段记录——
trigger_code记"有没有派下去",handle_code记"有没有跑成功",失败重试的判定才准确; - 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 次 │
│ 告警 │ 最终失败邮件/自定义扩展 │
└──────────────┴────────────────────────────────────────────┘
手写精简版的代码量账
| 模块 | 核心类 | 行数(约) |
|---|---|---|
| 注册中心 | JobRegistryService | 60 |
| 调度器 | JobScheduler + RouteStrategy | 180 |
| 回调/重试/告警 | CallbackService + RetryService + AlertService | 150 |
| 执行器 | MiniJobExecutor | 200 |
| 表结构 + 业务接入 | 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 的分片广播在你业务里单次最多分了多少片?评论区聊聊。
参考资料
- XXL-Job GitHub 仓库
- XXL-Job 官方文档
- Quartz 集群模式官方文档
- Cron 表达式语法(Spring/Quartz 差异说明)
- PowerJob(新一代分布式调度框架)
- Spring Boot @Scheduled 文档
标题:从0到1搭建分布式定时任务平台:XXL-Job 原理拆解+手写精简版
作者:jiangyi
地址:http://www.jiangyi.space/articles/2026/09/04/1788098721234.html
公众号:服务端技术精选
- 引言
- 一、XXL-Job 核心架构拆解
- 1.1 整体架构:调度与执行分离
- 1.2 关键机制一:执行器自动注册
- 1.3 关键机制二:Quartz 集群调度 + DB 行锁
- 1.4 关键机制三:路由策略与分片广播
- 1.5 关键机制四:执行结果回调与失败重试
- 1.6 XXL-Job 架构总结:一张图看懂全链路
- 二、手写精简版:整体设计
- 2.1 目标与边界
- 2.2 项目结构
- 2.3 数据库设计(4 张表)
- 三、手写精简版:调度中心实现
- 3.1 执行器注册与判活
- 3.2 调度入口:到点触发 + 路由 + 触发结果落库
- 3.3 路由策略
- 3.4 接收回调 + 失败重试 + 告警
- 3.5 执行器:内嵌 Server + 注册线程 + 任务执行
- 四、验证:全链路跑通
- 4.1 验证步骤
- 4.2 三类典型问题在精简版里的表现
- 五、常见问题
- 5.1 精简版用扫描 + 乐观锁,和 XXL-Job 的 Quartz 集群差在哪?
- 5.2 回调丢了怎么办(执行成功但调度中心没收到)?
- 5.3 为什么调度中心不做成"推任务到 MQ"?
- 5.4 XXL-Job 的任务为什么必须幂等?怎么做到?
- 5.5 分片广播时某个实例挂了,它负责的分片数据谁处理?
- 5.6 精简版和 XXL-Job 的差距清单(生产前必须补的)
- 六、总结
- XXL-Job 核心机制速查卡
- 手写精简版的代码量账
- 一句话
- 给团队的建议
- 参考资料
评论