XXL-JOB架构设计与调度中心部署
XXL-JOB是大众点评开源的分布式任务调度平台,核心设计原理是调度与执行分离——调度中心负责定时触发任务,执行器负责实际业务逻辑。相比Quartz和Elastic-Job,XXL-JOB的优势在于运维界面友好、动态路由策略丰富、支持任务分片广播,在微服务架构的后端开发中被广泛使用。
架构组件说明:调度中心(xxl-job-admin)是独立部署的Web应用,负责任务管理、日志收集和调度触发;执行器(xxl-job-executor)嵌入在业务应用中,接收调度中心的请求执行具体任务。两者通过HTTP协议通信。
调度中心部署(Docker方式):
# 创建数据库和表
mysql -u root -p -e "CREATE DATABASE xxl_job DEFAULT CHARACTER SET utf8mb4;"
# 导入初始化SQL
wget https://github.com/xuxueli/xxl-job/raw/2.4.1/doc/db/tables_xxl_job.sql
mysql -u root -p xxl_job < tables_xxl_job.sql
# Docker部署调度中心
docker run -d --name xxl-job-admin -p 8080:8080 -e PARAMS="--spring.datasource.url=jdbc:mysql://mysql-host:3306/xxl_job?Unicode=true&characterEncoding=UTF-8 --spring.datasource.username=root --spring.datasource.password=your_password --xxl.job.accessToken=your_token" -v /data/logs/xxl-job:/data/applogs/xxl-job xuxueli/xxl-job-admin:2.4.1
# 验证: 访问 http://localhost:8080/xxl-job-admin
# 默认账号 admin / 123456
执行器配置与任务开发
Spring Boot项目集成执行器配置:
<!-- pom.xml -->
<dependency>
<groupId>com.xuxueli</groupId>
<artifactId>xxl-job-core</artifactId>
<version>2.4.1</version>
</dependency>
# application.yml
xxl:
job:
admin:
addresses: http://xxl-job-admin:8080/xxl-job-admin
accessToken: your_token
executor:
appname: payment-job-executor
address:
ip:
port: 9999
logpath: /data/logs/xxl-job/jobhandler
logretentiondays: 30
// 执行器配置类
@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;
@Value("${xxl.job.executor.logpath}")
private String logPath;
@Bean
public XxlJobSpringExecutor xxlJobExecutor() {
XxlJobSpringExecutor executor = new XxlJobSpringExecutor();
executor.setAdminAddresses(adminAddresses);
executor.setAccessToken(accessToken);
executor.setAppname(appname);
executor.setPort(port);
executor.setLogPath(logPath);
executor.setLogRetentionDays(30);
return executor;
}
}
任务Handler开发:
@Component
public class PaymentJobHandler {
@Autowired
private PaymentService paymentService;
/**
* 简单任务:对账文件生成
*/
@XxlJob("reconciliationJob")
public void reconciliationJob() {
XxlJobHelper.log("对账任务开始执行");
try {
String date = XxlJobHelper.getJobParam();
if (StringUtils.isBlank(date)) {
date = LocalDate.now().minusDays(1)
.format(DateTimeFormatter.ISO_DATE);
}
XxlJobHelper.log("处理日期: {}", date);
int total = paymentService.generateReconciliation(date);
XxlJobHelper.log("对账完成,共处理 {} 笔交易", total);
XxlJobHelper.handleSuccess("对账成功: " + total + "笔");
} catch (Exception e) {
XxlJobHelper.log("对账任务异常: {}", e.getMessage());
XxlJobHelper.handleFail("对账失败: " + e.getMessage());
}
}
/**
* 分片广播任务:并行处理大量数据
*/
@XxlJob("shardingPaymentJob")
public void shardingPaymentJob() {
int shardIndex = XxlJobHelper.getShardIndex();
int shardTotal = XxlJobHelper.getShardTotal();
XxlJobHelper.log("分片参数: 当前第{}片, 共{}片",
shardIndex, shardTotal);
List<Payment> payments = paymentService
.findByShard(shardIndex, shardTotal, 1000);
int successCount = 0;
int failCount = 0;
for (Payment payment : payments) {
try {
paymentService.processPayment(payment);
successCount++;
} catch (Exception e) {
failCount++;
XxlJobHelper.log("处理失败: orderId={}, error={}",
payment.getOrderId(), e.getMessage());
}
}
XxlJobHelper.log("分片{}处理完成: 成功{}, 失败{}",
shardIndex, successCount, failCount);
XxlJobHelper.handleSuccess();
}
}
任务路由策略与分片广播
XXL-JOB内置多种路由策略,适用于不同场景:
FIRST // 第一个: 固定选择第一个执行器
LAST // 最后一个: 固定选择最后一个执行器
ROUND // 轮询: 依次轮询所有执行器
RANDOM // 随机: 随机选择一个执行器
CONSISTENT_HASH // 一致性HASH: 同一任务每次路由到同一执行器
LEAST_FREQUENTLY_USED // 最不经常使用
LEAST_RECENTLY_USED // 最近最久未使用
FAILOVER // 故障转移: 依次检测,跳过不可用执行器
BUSYOVER // 忙碌转移: 依次检测,跳过忙碌执行器
SHARDING_BROADCAST // 分片广播: 广播给所有执行器
分片广播是多实例并行处理的核心。假设有4个执行器实例,调度中心广播任务时,每个实例收到不同的shardIndex(0-3)。数据分片通常采用取模方式:
-- 数据分片查询SQL
-- MySQL方案: 利用MOD函数
SELECT * FROM payment_orders
WHERE MOD(id, #{shardTotal}) = #{shardIndex}
AND status = 'pending'
LIMIT 1000;
-- Elasticsearch方案: 利用routing
SearchResponse response = client.prepareSearch("payment_orders")
.setRouting(String.valueOf(shardIndex))
.setQuery(QueryBuilders.termQuery("status", "pending"))
.setSize(1000)
.get();
任务监控与告警配置
// 自定义告警处理器
@Component
public class CustomJobAlarm implements JobAlarm {
@Autowired
private DingTalkService dingTalkService;
@Override
public boolean doAlarm(JobInfo info, JobAlarmTrigger trigger) {
String content = String.format(
"【XXL-JOB任务告警】
任务: %s
状态: %s
触发时间: %s
" +
"执行参数: %s
执行结果: %s",
info.getJobDesc(),
trigger.getHandleCode() == 200 ? "成功" : "失败",
new Date(trigger.getTriggerTime().getTime()),
trigger.getExecutorParam(),
trigger.getHandleMsg()
);
dingTalkService.sendMarkdown("任务调度告警", content);
if (trigger.getHandleCode() != 200) {
sendAlertEmail(info.getJobDesc(), content);
}
return true;
}
}
在调度中心界面配置任务时,关键调度参数包括:Cron表达式定义触发时间(如0 0 2 * * ?表示每天凌晨2点),路由策略选择(简单任务用FAILOVER,大数据量任务用SHARDING_BROADCAST),阻塞处理策略(SERIAL_EXECUTION串行、DISCARD_LATER丢弃后续、COVER_EARLY覆盖之前),任务超时时间(单位秒,0表示不限制),失败重试次数(0表示不重试)。
执行器集群水平扩容时,新实例自动注册到调度中心,无需手动配置。分片广播任务会在下次触发时自动包含新实例,数据分片范围自动调整。注意事项:分片任务在扩缩容瞬间可能存在数据重复处理,需要在业务逻辑中实现幂等性保证。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/xxljob-fen-bu-shi-ding-shi-ren-wu-shi-zhan-ji-qun-bu-shu-yu/