XXL-Job 分布式定时任务
概述
XXL-Job 是一个轻量级分布式任务调度平台,核心设计理念是「调度中心 + 执行器」分离架构。
核心特性
- 调度中心统一管理任务、执行器注册、调度日志
- 执行器以 Spring Boot Starter 集成到业务服务中
- 支持分片广播、动态创建、任务依赖、失败告警
一、架构
┌──────────────────────────────────────────────┐
│ 调度中心 (xxl-job-admin) │
│ 任务管理 / 执行器管理 / 调度日志 / 报警通知 │
└────────────────┬─────────────────────────────┘
│ HTTP 调度(RESTful)
┌────────────┼────────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ 执行器 A │ │ 执行器 B │ │ 执行器 C │
│ ─────── │ │ ─────── │ │ ─────── │
│ order-svc│ │ pay-svc │ │ report-svc│
└──────────┘ └──────────┘ └──────────┘二、Spring Boot Starter 集成
2.1 依赖
xml
<dependency>
<groupId>com.xuxueli</groupId>
<artifactId>xxl-job-core</artifactId>
<version>2.4.1</version>
</dependency>2.2 配置
yaml
xxl:
job:
admin:
addresses: http://localhost:8080/xxl-job-admin
accessToken: default_token
executor:
appname: order-service
address:
ip:
port: 9999
logpath: /data/applogs/xxl-job/jobhandler
logretentiondays: 302.3 执行器配置类
java
@Configuration
public class XxlJobConfig {
@Value("${xxl.job.admin.addresses}")
private String adminAddresses;
@Value("${xxl.job.accessToken}")
private String accessToken;
@Value("${xxl.job.executor.appname}")
private String appname;
@Value("${xxl.job.executor.port}")
private int port;
@Bean
public XxlJobSpringExecutor xxlJobExecutor() {
XxlJobSpringExecutor executor = new XxlJobSpringExecutor();
executor.setAdminAddresses(adminAddresses);
executor.setAppname(appname);
executor.setPort(port);
executor.setAccessToken(accessToken);
executor.setLogRetentionDays(30);
return executor;
}
}三、任务开发
3.1 简单任务
java
@Component
public class OrderJobHandler {
@XxlJob("orderTimeoutCancel")
public ReturnT<String> orderTimeoutCancel(String param) {
XxlJobHelper.log("开始处理超时订单取消任务");
List<Order> timeoutOrders = orderService.findTimeoutOrders();
for (Order order : timeoutOrders) {
orderService.cancel(order.getId(), "超时自动取消");
XxlJobHelper.log("取消订单: {}", order.getOrderNo());
}
XxlJobHelper.log("任务完成,共取消 {} 个订单", timeoutOrders.size());
return ReturnT.SUCCESS;
}
}3.2 分片任务
java
@Component
public class ShardingJobHandler {
@XxlJob("dataSyncSharding")
public ReturnT<String> dataSyncSharding(String param) {
// 获取分片信息
int shardIndex = XxlJobHelper.getShardIndex(); // 当前分片索引
int shardTotal = XxlJobHelper.getShardTotal(); // 总分片数
List<User> users = userService.selectByPage(shardIndex, shardTotal, 5000);
for (User user : users) {
syncUser(user);
XxlJobHelper.log("同步用户: {}", user.getId());
}
return ReturnT.SUCCESS;
}
}3.3 任务参数传递
java
@XxlJob("dynamicReport")
public ReturnT<String> generateReport(String param) {
// param 是调度中心传入的 JSON 字符串
ReportParam reportParam = JSON.parseObject(param, ReportParam.class);
String date = reportParam.getDate();
String type = reportParam.getType();
XxlJobHelper.log("生成报表: date={}, type={}", date, type);
reportService.generate(date, type);
return ReturnT.SUCCESS;
}四、任务监听与回调
4.1 任务执行结果
java
@XxlJob("paymentCallback")
public ReturnT<String> paymentCallback(String param) {
try {
// 模拟处理
int total = processPayment();
XxlJobHelper.log("处理支付回调 {} 笔", total);
// 返回 SUCCESS
return ReturnT.SUCCESS;
} catch (Exception e) {
XxlJobHelper.log(e);
// 返回 FAIL → 调度中心会触发告警
return ReturnT.FAIL;
}
}4.2 失败告警
XXL-Job 内置了邮件告警,同时支持自定义告警:
java
@Component
public class MyJobAlarm implements JobAlarm {
@Override
public boolean doAlarm(JobInfo info, JobLog log) {
String message = String.format(
"任务 [%s] 执行失败,日志ID: %d,执行时间: %s",
info.getJobDesc(), log.getId(), log.getTriggerTime()
);
// 发送钉钉/企业微信/短信告警
dingTalkService.send(message);
return true;
}
}五、实战:数据每日同步
5.1 场景
数据平台每天需要从 MySQL 同步 2000 万订单数据到 ClickHouse,要求:
- 分片并行处理,提高效率
- 失败自动告警
- 完整的执行日志
5.2 实现
java
@Component
public class DataSyncJobHandler {
@Autowired
private OrderDao orderDao;
@Autowired
private ClickHouseDao clickHouseDao;
private static final int PAGE_SIZE = 5000;
@XxlJob("orderSyncToClickHouse")
public ReturnT<String> syncOrders(String param) {
int shardIndex = XxlJobHelper.getShardIndex();
int shardTotal = XxlJobHelper.getShardTotal();
LocalDate syncDate = param != null ? LocalDate.parse(param) : LocalDate.now().minusDays(1);
long startId = 0;
int totalCount = 0;
boolean hasMore = true;
while (hasMore) {
List<Order> orders = orderDao.selectByPage(
startId, shardIndex, shardTotal, PAGE_SIZE, syncDate
);
if (orders.isEmpty()) {
hasMore = false;
} else {
clickHouseDao.batchInsert(orders);
totalCount += orders.size();
startId = orders.get(orders.size() - 1).getId();
XxlJobHelper.log("分片 {} 已同步 {} 条", shardIndex, totalCount);
}
}
XxlJobHelper.log("分片 {} 同步完成,共 {} 条", shardIndex, totalCount);
return ReturnT.SUCCESS;
}
}5.3 调度中心配置
| 配置项 | 值 |
|---|---|
| 任务名称 | orderSyncToClickHouse |
| Cron | 0 0 3 * * ?(每天凌晨 3 点) |
| 路由策略 | 分片广播(Sharding) |
| 阻塞策略 | 串行(SERIAL_EXECUTION) |
| 失败策略 | 失败转移(FAILOVER) |
六、总结
| 知识点 | 说明 |
|---|---|
| 架构 | 调度中心 + 执行器,HTTP 通信 |
| 任务 | @XxlJob("name") 注解,方法级别定义 |
| 分片 | getShardIndex() / getShardTotal() 分片广播 |
| 日志 | XxlJobHelper.log() 输出到调度中心 |
| 失败告警 | 内置邮件 + JobAlarm 接口自定义告警 |
| 与 @Scheduled 对比 | 支持集群、分片、失败转移、可视化控制台 |
参考链接: