XXL-Job 源码阅读
XXL-Job 是使用最广的分布式任务调度框架之一,中心化架构:调度中心(Admin)负责触发,执行器(Executor)负责执行。本文从源码拆解调度触发、路由、分片、故障转移与日志回传的实现,核心代码位于 xxl-job-admin 与 xxl-job-core 两个模块。
模块与核心类
XXL-Job 源码模块:
├─ xxl-job-admin:调度中心(Spring Boot 应用)
│ ├─ controller:任务管理、执行器管理、日志管理接口
│ ├─ core.thread:调度线程 JobScheduleHelper
│ ├─ core.trigger:任务触发器 JobTrigger
│ └─ dao:任务、日志、执行器注册信息持久化
└─ xxl-job-core:执行器 SDK(内嵌到业务应用)
├─ thread:JobLogFileCleanThread、JobThread 等
├─ handler:任务处理器(IJobHandler)
└─ glue:Glue 模式(动态代码执行)核心接口
java
// 执行器侧:任务处理器抽象
public abstract class IJobHandler {
public abstract void execute() throws Exception;
}
// 内置处理器:方法模式(Spring Bean 方法)
public class MethodJobHandler extends IJobHandler {
private final Object target; // 目标 Bean
private final Method method; // 目标方法
private final Method initMethod; // 初始化方法
private final Method destroyMethod; // 销毁方法
}一、执行器注册与发现
执行器启动注册
执行器启动流程:
1. XxlJobExecutor 初始化(Spring Boot 自动装配)
2. 加载执行器配置(appname、ip、port)
3. 启动嵌入式 Netty 服务(接收调度请求)
4. 启动 ExecutorRegistryThread(注册线程)
5. 向调度中心发起注册(每 30 秒一次心跳)注册线程
java
public class ExecutorRegistryThread {
public void start() {
registryThread = new Thread(() -> {
while (!toStop) {
// 1. 组装注册参数:registryParam(appname, ip, port)
RegistryParam registryParam = new RegistryParam(
XxlJobExecutor.getAdminConfig().getAppname(),
XxlJobExecutor.getAdminConfig().getAddress());
// 2. 遍历所有配置的调度中心地址
for (AdminBiz adminBiz : XxlJobExecutor.getAdminBizList()) {
try {
adminBiz.registry(registryParam); // 注册到调度中心
} catch (Exception e) {
logger.error(e.getMessage(), e);
}
}
// 3. 每 30 秒循环一次
TimeUnit.SECONDS.sleep(30);
}
});
registryThread.start();
}
}调度中心侧处理
java
// AdminBiz.registry:调度中心接收注册
public class AdminBizImpl implements AdminBiz {
@Override
public ReturnT<String> registry(RegistryParam registryParam) {
// 1. 查执行器是否存在(按 appname)
// 2. 不存在则创建执行器记录
// 3. 更新执行器地址列表(注册信息)
return ReturnT.SUCCESS;
}
}二、任务触发(核心)
调度中心每 5 秒扫描一次任务表,找出到期的任务并触发。
JobScheduleHelper
java
public class JobScheduleHelper {
// 调度线程:每 5 秒扫描一次
public void start() {
// 轮询线程
scheduleThread = new Thread(() -> {
while (!toStop) {
// 1. 查下一批待调度的任务
List<XxlJobInfo> scheduleList = xxlJobInfoDao
.scheduleJobQuery(nowTime, loadOnce);
for (XxlJobInfo jobInfo : scheduleList) {
// 2. 计算下次触发时间
// 3. 将任务放入触发队列 ringData
// 4. 更新任务的下次触发时间
}
}
});
// 触发线程:消费 ringData 队列,执行触发
ringThread = new Thread(() -> {
while (!toStop) {
// 取出到期的任务 → 调用 JobTrigger.trigger(...)
}
});
}
}触发流程
JobTrigger.trigger 全流程:
1. 组装触发参数(jobId、triggerType、glueType)
2. 根据任务配置的路由策略选择执行器
3. 若配置了阻塞处理策略 → 处理阻塞(单机串行/丢弃后续调度)
4. 发起远程调度:EmbedServer 发送调度请求
5. 记录调度日志(XxlJobLog)调度请求发送
java
public class EmbedServer {
// 执行器侧 Netty 服务收到调度请求
private class EmbedHttpServerHandler extends SimpleChannelInboundHandler<FullHttpRequest> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, FullHttpRequest request) {
// 解析 URL:/run → 执行任务
// 反序列化 TriggerParam
// 交给 JobThread 执行
}
}
}三、路由策略
路由策略在调度中心侧计算"这次任务派发给哪个执行器"。
路由策略接口
java
public abstract class ExecutorRouter {
public abstract ReturnT<String> route(TriggerParam triggerParam,
List<String> addressList);
}
// 具体策略示例:轮询
public class ExecutorRouteRound extends ExecutorRouter {
private static ConcurrentMap<Integer, AtomicInteger> routeCountEachJob =
new ConcurrentHashMap<>();
@Override
public ReturnT<String> route(TriggerParam triggerParam, List<String> addressList) {
// 按 jobId 取计数器 → 自增 → 取模得到地址
int index = routeCount.incrementAndGet() % addressList.size();
return new ReturnT<>(addressList.get(index));
}
}一致性哈希路由
java
public class ExecutorRouteConsistentHash extends ExecutorRouter {
// 对每个执行器地址生成多个虚拟节点
private static int VIRTUAL_NODE_NUM = 5;
public String hashJob(int jobId, List<String> addressList) {
TreeMap<Long, String> addressRing = new TreeMap<>();
for (String address : addressList) {
for (int i = 0; i < VIRTUAL_NODE_NUM; i++) {
// 虚拟节点:address + "-" + i
long addressHash = hash(address + i);
addressRing.put(addressHash, address);
}
}
// 按 jobId 哈希找环上最近的节点
long jobHash = hash(String.valueOf(jobId));
SortedMap<Long, String> lastRing = addressRing.tailMap(jobHash);
if (!lastRing.isEmpty()) {
return lastRing.get(lastRing.firstKey());
}
return addressRing.firstEntry().getValue();
}
}路由策略清单
内置路由策略:
├─ 第一个 / 最后一个:固定地址
├─ 轮询:Round(计数器取模)
├─ 随机:Random
├─ 一致性哈希:ConsistentHash(同 jobId 固定到同一执行器)
├─ 最不经常使用:LFU(记录调用频率,选频率最低的)
├─ 最近最久未使用:LRU
├─ 故障转移:Failover(先 ping 探测,选可用的)
├─ 忙碌转移:Busyover(选当前空闲的执行器)
└─ 分片广播:Sharding(所有执行器都执行)四、分片广播
分片参数传递
分片广播触发时:
├─ 所有执行器都会收到调度请求
├─ 每个请求携带分片参数:
│ ├─ shardingParam:形如 "0/3"(当前分片/总分片)
│ └─ 由调度中心按执行器列表序号生成
└─ 各执行器按自己的分片号处理对应数据业务侧使用
java
@XxlJob("orderShardingJob")
public void orderShardingJob() {
// 从上下文拿分片参数
ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo();
int shardIndex = shardingVO.getIndex(); // 当前分片
int shardTotal = shardingVO.getTotal(); // 总分片
// 只处理自己分片的数据(如订单号取模 = shardIndex)
List<Order> orders = orderDao.queryByShard(shardIndex, shardTotal);
for (Order order : orders) {
process(order);
}
}分片场景示例
4 个执行器,1 个分片任务:
├─ 执行器1:shardIndex=0 → 处理订单号 % 4 == 0
├─ 执行器2:shardIndex=1 → 处理订单号 % 4 == 1
├─ 执行器3:shardIndex=2 → 处理订单号 % 4 == 2
└─ 执行器4:shardIndex=3 → 处理订单号 % 4 == 3
效果:单任务拆成 4 份并行处理,吞吐提升 4 倍五、执行器侧执行
JobThread
java
public class JobThread extends Thread {
// 接收调度请求(生产)
public ReturnT<String> pushTriggerQueue(TriggerParam triggerParam) {
// 1. 判断阻塞策略:若线程忙且配置了丢弃 → 直接返回
// 2. 放入阻塞队列 triggerQueue
// 3. 若线程未启动 → 启动线程
}
// 消费执行(消费者)
@Override
public void run() {
while (!toStop) {
// 1. 从队列取任务
TriggerParam triggerParam = triggerQueue.poll();
// 2. 获取任务处理器(缓存)
IJobHandler jobHandler = XxlJobExecutor.loadJobHandler(...);
// 3. 执行业务逻辑,记录执行日志
// 4. 回传执行结果到调度中心
}
}
}阻塞处理策略
阻塞处理策略(任务排队还是丢弃):
├─ 单机串行:调度请求放入队列,逐个执行
├─ 丢弃后续调度:线程忙时直接丢弃新调度
└─ 覆盖之前调度:终止旧任务,执行新任务六、日志回传与查看
日志写入
执行器侧写日志:
├─ XxlJobFileAppender:日志追加到本地文件
│ └─ 按任务 id 分目录:/logs/{jobId}/{logId}.log
└─ 日志行包含时间戳、执行线程、级别
日志回传:
├─ 执行完成后调用 callback(回传执行结果)
├─ 调度中心按需拉取日志文件内容(日志查询接口)
└─ 控制台远程查看执行器上的日志日志清理
java
// 执行器侧日志清理线程(JobLogFileCleanThread)
// 按配置的保留天数清理过期日志文件
public class JobLogFileCleanThread {
private static long logRetentionDays = 30; // 可配置
private void cleanLogFile() {
// 扫描日志目录,删除超过保留天数的文件
}
}七、故障转移
故障转移流程(路由策略 = 故障转移):
1. 取到候选执行器地址列表
2. 逐个地址做 HTTP ping(GET /beat 或发送调度探测)
3. 第一个成功的地址 → 派发任务
4. 全部失败 → 返回失败并告警
其他保障:
├─ 执行器心跳超时(90 秒)→ 调度中心标记为失效
├─ 失效执行器不再接收调度
└─ 任务失败支持重试(调度失败重试次数配置)八、Glue 模式(动态代码)
Glue 模式:调度中心在线编写代码,执行器动态加载执行
├─ GLUE_GROOVY:Groovy 脚本(源码存调度中心 DB)
├─ GLUE_SHELL / GLUE_PYTHON / GLUE_PHP:脚本语言
└─ GLUE_JAVA:Java 代码在线编辑
原理:
├─ 执行器按需从调度中心拉取源码
├─ 本地编译/解释执行(Groovy 用 GroovyClassLoader)
└─ 适合"小逻辑快速迭代"的运维场景关键链路图
调度中心 执行器
│ │
├─ JobScheduleHelper 扫描到期任务 ──┤
├─ JobTrigger.trigger │
├─ ExecutorRouter 路由选执行器 │
├─ EmbedServer 发送调度请求 ────────▶│ Netty 接收
│ ├─ JobThread 消费
│ ├─ IJobHandler.execute
│◀─────────── 结果回调 ─────────────┤
├─ 写调度日志 │
└─ 告警/重试 │总结
XXL-Job 的源码骨架一句话概括:调度中心定时扫描、路由派发;执行器注册心跳、接收执行、回传结果。读懂 JobScheduleHelper 的触发线程、ExecutorRouter 的路由策略、JobThread 的生产消费模型,就掌握了它的核心。分片广播通过参数传递实现数据并行,故障转移靠心跳与探测保证可靠性,日志回传与 Glue 模式则解决了运维问题。