1. 为什么单机定时任务在集群里一定会出问题
很多团队第一次把应用从单机部署改成多副本集群时,都会遇到同一个尴尬场景:原本在单机跑得好好的定时任务,突然开始重复执行。比如每天凌晨两点给用户发结算邮件,结果三台机器同时触发,用户收到了三封一模一样的邮件;又比如库存对账任务,三个实例同时跑,数据被反复覆盖,最后账目全乱。
这个问题的根源其实很朴素。@Scheduled这类注解是绑定在 JVM 进程上的,它只认自己所在的这台机器,根本不知道外面还有几个兄弟实例。你部署了三个副本,就等于把同一个闹钟买了三个,到点当然一起响。有人会说,那我加个 Redis 分布式锁不就行了?锁确实能解决"同一时刻只有一个实例执行"的问题,但它解决不了另一个更隐蔽的需求:我想让三个实例一起干活,每个只干一部分。
举个真实例子。我做过一个批量推送项目,需要给两百万用户发站内信。单机跑一轮要四十多分钟,业务方嫌慢。这时候分布式锁反而成了阻碍——它强制串行,三台机器只有一台在忙,另外两台干看着。真正需要的是一种"把任务拆开、分给所有实例并行处理"的机制。这就是 XXL-JOB 分片广播模式要解决的核心问题:让一次调度触发所有执行器,同时给每个执行器一个编号,让它们各自认领属于自己的那份数据。
理解了这个出发点,后面所有的配置、代码、坑点才有意义。分片广播不是"更高级的定时任务",它是"分布式并行计算"在任务调度领域的一个具体落地。你把它当成 MapReduce 里那个 Map 阶段就对了——调度中心负责切分,执行器负责各自处理分片,最后汇总结果。
2. 分片广播的底层机制:调度中心和执行器到底怎么配合
2.1 一次广播调度的完整链路
要搞懂分片广播,得先看清楚一次调度请求从发出到执行完毕,中间经过了哪些环节。XXL-JOB 的架构里有两个核心角色:调度中心(admin)和执行器(executor)。普通任务的路由策略是"选一台",而分片广播的路由策略是"全选"。
具体链路是这样的:调度中心到点触发任务,根据路由策略发现当前是SHARDING_BROADCAST,于是它不会只挑一台执行器,而是把注册在案的所有执行器地址全部拉出来,逐个发起调度请求。注意,这里是逐个发起,不是发一条消息让执行器自己抢。每个执行器收到的请求里,都带着两个关键参数:分片总数(shardTotal)和当前分片序号(shardIndex)。
执行器收到请求后,把这两个参数塞进任务方法的入参里。你的业务代码通过XxlJobHelper.getShardIndex()和XxlJobHelper.getShardTotal()就能拿到。分片序号从 0 开始,比如三台机器,序号就是 0、1、2。业务代码拿到序号后,用取模的方式决定"我该处理哪些数据"。
这里有个容易被忽略的细节:分片序号和执行器实例不是永久绑定的。今天序号 0 可能是 A 机器,明天 A 机器下线了,序号 0 就落到 B 机器头上。所以你的业务逻辑绝对不能依赖"序号 0 一定是某台特定机器"这种假设,只能依赖"序号 0 处理 id 取模等于 0 的数据"这种纯计算逻辑。
2.2 分片参数是怎么传进业务方法的
很多人第一次写分片任务,会习惯性地在方法签名里加参数,结果发现拿不到值。XXL-JOB 的新版本(2.2.0 之后)推荐用XxlJobHelper这个工具类来获取上下文,而不是靠方法参数注入。
@XxlJob("shardingJobHandler") public void shardingJobHandler() throws Exception { // 获取分片序号和总分片数 int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); XxlJobHelper.log("当前分片序号:{},分片总数:{}", shardIndex, shardTotal); // 业务逻辑:只处理属于自己分片的数据 List<Long> userIds = userMapper.selectByShard(shardIndex, shardTotal); for (Long userId : userIds) { processUser(userId); } }对应的 SQL 通常长这样:
SELECT id, name FROM user WHERE MOD(id, #{shardTotal}) = #{shardIndex} LIMIT 1000用MOD(id, shardTotal) = shardIndex这个条件,就能保证每条数据只会被一个分片捞到,不会重复也不会遗漏。这是分片广播最经典的用法,也是面试里被问得最多的点。
2.3 分片总数到底等于几
这是新手最容易踩的坑。分片总数不等于你配置的机器数量,而是等于当前在线执行器实例的数量。调度中心在发起广播前,会实时查询执行器注册表,拿到当前活着的实例列表,然后shardTotal就等于这个列表的长度。
这意味着什么?意味着你的分片数是动态的。早上三台机器在线,shardTotal就是 3;中午扩容到五台,shardTotal就变成 5。业务代码必须能适应这种变化,不能写死。我见过有人把分片逻辑写成if (shardIndex == 0) { 处理前半部分 } else { 处理后半部分 },结果扩容到三台后直接崩了——第三台机器不知道该干啥。
正确的做法永远是用取模运算动态计算,让分片数和数据切分规则解耦。这样无论实例怎么增减,数据都能被正确覆盖。
3. 从零搭一个分片广播任务:配置、代码、验证
3.1 调度中心和执行器的版本对齐
动手之前先确认版本。XXL-JOB 在 2.1.0 之后对分片广播的支持才比较完善,2.2.0 引入了XxlJobHelper,2.3.0 之后对注册中心做了优化。我建议直接用 2.3.x 或 2.4.x,老版本在实例上下线时容易出现分片数计算不准的问题。
调度中心(xxl-job-admin)和执行器(xxl-job-core)的版本必须一致,这是硬性要求。我踩过一次坑:admin 用的 2.3.0,执行器依赖写成了 2.2.0,结果分片参数传过去是 null,排查了大半天才发现是版本不匹配。Maven 里这样对齐:
<dependency> <groupId>com.xuxueli</groupId> <artifactId>xxl-job-core</artifactId> <version>2.4.0</version> </dependency>3.2 执行器配置文件的关键项
执行器的application.properties里,和分片广播相关的配置其实不多,但每一项都不能错:
# 调度中心地址,多个用逗号分隔 xxl.job.admin.addresses=http://127.0.0.1:8080/xxl-job-admin # 执行器注册方式,推荐自动注册 xxl.job.executor.address= xxl.job.executor.ip= xxl.job.executor.port=9999 xxl.job.executor.logpath=/data/applogs/xxl-job/jobhandler xxl.job.executor.logretentiondays=30 # 执行器AppName,调度中心靠这个找到执行器 xxl.job.executor.appname=xxl-job-executor-sharding重点说appname。调度中心在配置任务时,要选择"执行器",这个执行器就是靠appname关联的。如果你部署了三个实例,它们的appname必须完全相同,这样调度中心才会认为它们是同一个执行器下的三个节点,广播时才会把三个都算进去。如果appname写得不一致,调度中心会当成三个独立的执行器,分片逻辑就乱了。
port这一项在多实例部署时要注意,如果三台机器在不同主机上,端口可以都用 9999;如果在同一台机器上跑多个实例(本地测试常见),端口必须错开,否则后启动的会绑定失败。
3.3 调度中心里怎么配这个任务
登录调度中心,新建任务,几个关键配置项:
| 配置项 | 值 | 说明 |
|---|---|---|
| 路由策略 | 分片广播 | 核心,选错就不是广播了 |
| 运行模式 | BEAN | 用注解方式注册的处理器 |
| JobHandler | shardingJobHandler | 和代码里@XxlJob的值对应 |
| 阻塞处理策略 | 丢弃后续调度 | 分片任务通常不希望堆积 |
| 调度过期策略 | 忽略 | 避免补跑导致数据重复 |
路由策略选"分片广播"是整件事的开关。选成"轮询"或"第一个",任务只会在一台机器上跑,分片参数虽然也会传,但shardTotal永远是 1,等于没分片。
阻塞处理策略我一般选"丢弃后续调度"。因为分片任务往往是批处理,一轮没跑完又来一轮,容易造成数据竞争。如果你的任务必须串行,那就得靠业务层面的幂等来兜底。
3.4 本地模拟多实例验证
本地验证分片广播,最省事的办法是启动多个执行器实例,改端口和日志路径:
java -jar executor.jar --server.port=8081 --xxl.job.executor.port=9991 java -jar executor.jar --server.port=8082 --xxl.job.executor.port=9992 java -jar executor.jar --server.port=8083 --xxl.job.executor.port=9993三个实例的appname保持一致。启动后去调度中心的"执行器管理"页面,应该能看到这个执行器下面挂着三个在线节点。这时候手动触发一次任务,看日志:三个实例的日志里应该分别打印出shardIndex=0/1/2,shardTotal=3。如果只看到一个实例在跑,八成是appname不一致或者路由策略没选对。
提示:本地测试时如果发现分片数不对,先检查执行器注册列表里到底有几个在线节点。调度中心的分片数是实时算的,注册表不准,分片就一定不准。
4. 分片逻辑写不对,等于白配
4.1 取模分片的正确姿势
分片广播最容易出问题的地方不在配置,而在业务代码里的数据切分逻辑。取模是最常用的方式,但写法有讲究。
// 推荐:用分片总数和序号做取模 int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); // 查询属于当前分片的数据 List<Task> tasks = taskMapper.selectSharding(shardIndex, shardTotal);对应的 MyBatis 映射:
<select id="selectSharding" resultType="Task"> SELECT * FROM task WHERE status = 0 AND MOD(id, #{shardTotal}) = #{shardIndex} ORDER BY id LIMIT #{pageSize} </select>这里有个性能陷阱:MOD(id, shardTotal)在数据量大时会导致全表扫描,因为函数作用在列上,索引失效。如果表有几百万行,这个查询会非常慢。优化思路有两个:一是用id的范围分片代替取模,二是提前算好分片字段并建索引。
范围分片的写法:
SELECT * FROM task WHERE status = 0 AND id % #{shardTotal} = #{shardIndex} AND id > #{lastMaxId} ORDER BY id LIMIT #{pageSize}配合游标(lastMaxId)翻页,既能走索引,又能避免深分页。这是我在实际项目里验证过的方案,两百万数据的分片查询从十几秒降到几百毫秒。
4.2 分片数变化时的数据一致性
前面说过,分片数是动态的。三台机器时MOD(id, 3),扩容到五台变成MOD(id, 5)。如果任务跑到一半扩容了,会发生什么?
假设第一轮用 3 个分片处理了 id 1 到 1000,第二轮扩容到 5 个分片,MOD(id, 5)的切分方式和MOD(id, 3)完全不同。原本 id=3 的数据在第一轮属于分片 0,第二轮可能属于分片 3。如果任务没有幂等保护,这条数据会被处理两次。
解决办法有两个层面。任务层面,给每条数据的处理加上状态标记,处理过的打上processed=1,查询时过滤掉。调度层面,尽量避免在任务执行期间扩容,或者把分片任务设计成"每次全量重算"而不是"增量累加"。我个人的经验是,批处理类任务用状态标记最稳妥,因为扩容是运维的常规操作,你没法保证它不在任务执行时发生。
4.3 分片任务里的日志和监控
分片任务出问题时,排查比单机任务麻烦得多,因为日志散在多个实例上。XXL-JOB 提供了XxlJobHelper.log()方法,它会把日志回传到调度中心,在"调度日志"页面能看到每个分片的执行情况。
XxlJobHelper.log("分片 {} 开始处理,待处理数量:{}", shardIndex, tasks.size()); // ... 处理逻辑 XxlJobHelper.log("分片 {} 处理完成,成功:{},失败:{}", shardIndex, success, fail);这个日志是分片任务排查的生命线。我习惯在每个分片的开头和结尾都打一条,这样在调度中心一眼就能看出哪个分片没跑、哪个分片卡住了。如果某个分片的日志一直不出现,说明那台执行器可能掉线了,或者任务分发时没收到请求。
另外,XxlJobHelper.log()的日志有长度限制,默认单条日志不能太长。如果要打大量内容,建议只打摘要,详细日志写到本地文件。
5. 那些文档里不会写的坑
5.1 执行器掉线导致的分片空洞
这是分片广播最隐蔽的问题。假设三台机器,分片 0、1、2。任务执行到一半,分片 1 的机器突然挂了。会发生什么?
调度中心在发起调度时,分片数是基于调度那一刻的在线实例算的。如果机器是在任务执行过程中挂的,那分片 1 的数据就没人处理了,而且调度中心不会自动把分片 1 重新分配给其他机器。结果就是这部分数据被漏掉,直到下一次调度才会被重新捞起来。
如果任务对实时性要求高,这个空洞是不能接受的。应对方案是补偿机制:任务执行完后,检查每个分片是否都上报了完成状态,没上报的分片由调度中心或某个协调者重新触发。XXL-JOB 本身不提供这个能力,需要自己在业务层实现。我通常会在任务表里加一个shard_status字段,每个分片完成后写入自己的状态,然后有一个独立的巡检任务定期检查有没有"超时未完成"的分片。
5.2 分片序号为负数的诡异情况
有次线上排查,发现日志里打印出shardIndex=-1。查了半天,原因是执行器版本和调度中心版本不一致,老版本执行器在拿不到分片参数时默认返回 -1。这种问题不会报错,只会让业务逻辑静默地处理错误的数据范围。
防御性写法:
int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); if (shardIndex < 0 || shardTotal <= 0 || shardIndex >= shardTotal) { XxlJobHelper.log("分片参数异常,shardIndex={}, shardTotal={}", shardIndex, shardTotal); return; }这段校验看起来多余,但能帮你快速定位版本不匹配的问题。我现在的习惯是每个分片任务开头都加这段,成本很低,收益很高。
5.3 分片任务和事务的冲突
分片任务通常是批处理,一批处理几百上千条。如果整个方法包在一个大事务里,一旦某条数据出错,整批回滚,前面处理成功的也白干了。而且大事务会长时间占用数据库连接,分片数一多,连接池直接被打满。
正确的做法是分批提交,每处理 N 条提交一次事务。可以用编程式事务:
int batchSize = 100; List<Task> tasks = taskMapper.selectSharding(shardIndex, shardTotal); for (int i = 0; i < tasks.size(); i += batchSize) { List<Task> batch = tasks.subList(i, Math.min(i + batchSize, tasks.size())); transactionTemplate.execute(status -> { for (Task task : batch) { processTask(task); } return null; }); }这样单批失败只影响这一批,前面的成果保留。配合状态标记,下一轮调度可以继续处理失败的批次。
5.4 分片数远大于数据量时的空转
如果数据只有 100 条,但你有 50 台机器,分片数就是 50。大部分分片查出来是空的,白白浪费调度资源。这种情况在小数据量、大集群的场景下很常见。
优化思路是限制分片数的上限。可以在业务代码里判断,如果数据量小于某个阈值,就只让分片 0 处理全部数据:
long totalCount = taskMapper.countPending(); if (totalCount < 1000) { // 数据量小,只让分片 0 处理 if (shardIndex != 0) { return; } // 分片 0 处理全部 processAll(); } else { // 正常分片处理 processSharding(shardIndex, shardTotal); }这个判断逻辑让分片机制在小数据量时退化成单机执行,避免无谓的空转。阈值设多少取决于你的单条处理耗时和集群规模,我一般设在 1000 到 5000 之间。
6. 分片广播和其他方案的取舍
6.1 和消息队列削峰的对比
有人会问,批量处理为什么不用消息队列?把两百万用户丢进 MQ,多个消费者并行消费,不也能达到并行处理的效果吗?
两者确实有重叠,但适用场景不同。MQ 适合流式、实时、无状态的处理,消息发出去就不管了,消费失败靠重试机制。分片广播适合批量、定时、有状态的处理,比如每天凌晨对账、每月生成报表。这类任务需要知道"总共处理了多少、成功多少、失败多少",需要能重新触发整个批次,这些是 MQ 不擅长的。
我的经验是:如果是用户触发的实时任务,用 MQ;如果是定时批处理,用分片广播。两者不是替代关系,很多系统里是共存的。
6.2 和 ElasticJob 的分片对比
ElasticJob 也支持分片,而且它的分片是持久化的——分片信息存在注册中心,实例上下线时会自动重新分片。XXL-JOB 的分片是每次调度时实时计算的,不持久化。
这个差异导致 ElasticJob 在实例频繁上下线时更稳定,因为它有重新分片的机制。但 ElasticJob 的运维复杂度更高,需要额外的注册中心(ZooKeeper 或 Nacos)。XXL-JOB 胜在轻量,调度中心自带注册功能,部署简单。
选型建议:如果你的集群规模不大(十几台以内),实例变动不频繁,XXL-JOB 足够用。如果集群规模大、弹性伸缩频繁,ElasticJob 的持久化分片更合适。我做过一个上百实例的项目,最后选了 ElasticJob,就是因为 XXL-JOB 在实例频繁上下线时分片数抖动太厉害。
6.3 分片广播的适用边界
不是所有任务都适合分片广播。判断标准很简单:任务能否被拆成互不依赖的独立子任务。能拆,就用分片;不能拆,就老老实实用单机加锁。
比如"生成全站统计报表"这种任务,它需要汇总所有数据,拆开反而要合并结果,用分片就是自找麻烦。而"给每个用户发通知"这种任务,用户之间互不依赖,天然适合分片。
还有一个边界是数据倾斜。如果数据分布不均匀,比如 id 取模后某个分片的数据量是其他分片的十倍,那这个分片就会成为瓶颈,其他分片早早跑完干等着。解决办法是换分片键,用更均匀的字段(比如用户 id 的哈希值)代替自增 id。这个坑我在一个订单处理项目里踩过,订单 id 是按时间递增的,取模后最近的分片数据量爆炸,后来改成按用户 id 哈希才解决。
7. 几个实战中的参数调优经验
7.1 调度超时时间怎么设
调度中心里有个"任务超时时间"配置,默认是 0(不限制)。分片任务一定要设这个值,否则某个分片卡死,整个任务永远不结束。
设多少合适?取决于你的单批处理耗时。我的经验公式是:超时时间 = 单批处理耗时 × 3 + 30秒。留三倍余量是为了应对数据量波动,加 30 秒是给网络和调度留缓冲。比如单批处理 2 分钟,超时设 7 分钟左右。
超时后 XXL-JOB 会中断任务,但注意,它只是标记任务失败,不会真的 kill 掉执行线程。所以业务代码里要有响应中断的逻辑,比如在循环里检查Thread.currentThread().isInterrupted()。
7.2 失败重试的坑
XXL-JOB 支持失败重试,配置项是"失败重试次数"。分片任务开重试要谨慎,因为重试是整个任务重试,不是只重试失败的分片。如果三个分片里只有一个失败,重试会把三个分片全部重跑一遍,已经成功的分片会重复处理。
如果任务不是幂等的,这个重试就是灾难。我的建议是:分片任务默认关闭重试,靠业务层的状态标记和补偿任务来处理失败。如果一定要开重试,确保业务逻辑幂等。
7.3 日志保留天数的权衡
xxl.job.executor.logretentiondays控制执行器本地日志的保留天数,默认 30 天。分片任务日志量大,如果集群规模大,30 天可能占满磁盘。我一般设成 7 天,因为调度中心已经存了关键日志,本地日志主要用于详细排查,7 天足够覆盖大部分问题。
调度中心的日志也有清理策略,在 admin 的配置文件里,默认保留 30 天。这个可以按需调整,但别设太短,否则出了问题想查历史日志都查不到。
8. 一个完整的分片任务模板
把前面所有经验揉在一起,给一个可以直接抄的模板:
@Component public class ShardingJobHandler { @Autowired private TaskMapper taskMapper; @Autowired private TransactionTemplate transactionTemplate; private static final int BATCH_SIZE = 100; private static final long SMALL_DATA_THRESHOLD = 1000; @XxlJob("shardingJobHandler") public void execute() { int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); // 参数校验 if (shardIndex < 0 || shardTotal <= 0 || shardIndex >= shardTotal) { XxlJobHelper.log("分片参数异常:index={}, total={}", shardIndex, shardTotal); XxlJobHelper.handleFail("分片参数异常"); return; } // 小数据量退化为单机 long totalCount = taskMapper.countPending(); if (totalCount < SMALL_DATA_THRESHOLD && shardIndex != 0) { XxlJobHelper.log("数据量小,分片 {} 跳过", shardIndex); return; } XxlJobHelper.log("分片 {} 开始,总数 {},待处理 {}", shardIndex, shardTotal, totalCount); long lastMaxId = 0; int successCount = 0; int failCount = 0; while (true) { // 游标分页,避免深分页 List<Task> tasks = taskMapper.selectShardingByCursor( shardIndex, shardTotal, lastMaxId, BATCH_SIZE); if (tasks.isEmpty()) { break; } for (Task task : tasks) { try { transactionTemplate.execute(status -> { processTask(task); return null; }); successCount++; } catch (Exception e) { failCount++; XxlJobHelper.log("处理失败:id={}, error={}", task.getId(), e.getMessage()); } lastMaxId = Math.max(lastMaxId, task.getId()); } } XxlJobHelper.log("分片 {} 完成,成功 {},失败 {}", shardIndex, successCount, failCount); if (failCount > 0) { XxlJobHelper.handleFail("存在失败记录:" + failCount); } } private void processTask(Task task) { // 业务处理逻辑 } }这个模板覆盖了参数校验、小数据退化、游标分页、分批事务、失败统计几个关键点。你可以根据自己的业务替换processTask和查询 SQL,其余部分基本不用改。
9. 排查分片问题的固定套路
分片任务出问题,排查顺序我总结成一套固定流程,照着走基本能定位:
第一步,看调度中心的调度日志。确认任务是否触发了、触发了几个分片、每个分片的执行状态。如果只触发了一个分片,问题在路由策略或执行器注册。
第二步,看执行器注册列表。确认在线实例数和预期是否一致。实例数不对,分片数就不对。
第三步,看每个分片的业务日志。通过XxlJobHelper.log()回传的日志,确认每个分片处理的数据范围。如果某个分片处理的数据为空,检查取模逻辑。
第四步,看数据是否重复或遗漏。用 SQL 统计每个分片处理的数据量,加起来是否等于总数。不等就说明分片逻辑有问题。
第五步,看是否有分片空洞。检查是否有分片没有上报完成状态,这通常意味着执行器掉线或任务超时。
这套流程我在多个项目里用过,大部分分片问题都能在前三步定位。第四步和第五步主要用于排查数据一致性问题,需要结合业务表的状态字段来分析。
分片广播这个机制,配置本身不复杂,难的是业务逻辑的正确性和边界情况的处理。把取模逻辑写对、把分片数变化考虑进去、把失败补偿做好,基本就能稳定运行。剩下的就是根据实际数据量和集群规模调参数,这部分没有标准答案,得靠实际跑几轮来摸索。