1. 为什么分片广播模式值得单独拿出来讲
做过分布式任务调度的朋友大概率都遇到过这样的场景:一张订单表里有几千万条待处理记录,单机跑批处理要跑几个小时,业务方催得急,机器却闲着一大半。这时候你自然会想到——能不能让多台机器同时干活,每台机器负责一部分数据?XXL-JOB 的分片广播模式就是为解决这类问题而生的。
我第一次在生产环境用分片广播,是因为一个对账任务。每天凌晨要处理前一天的对账数据,单机跑下来接近四十分钟,随着业务量增长迟早要出问题。改成五台机器分片跑之后,整体耗时压到了八分钟左右。但这个过程并不是改个配置就完事了,中间踩了不少坑,比如分片参数理解偏差导致数据重复处理、分片数远大于执行器数量导致部分分片空跑、任务执行时间超过调度周期引发连锁问题等等。
这篇文章我会把 XXL-JOB 分片广播模式从底层原理到生产实战完整拆一遍。不管你是刚接触 XXL-JOB 的新手,还是已经在用但没深究过分片机制的开发者,都能从中拿到可以直接落地的东西。核心关键词包括XXL-JOB、分片广播模式、分布式任务调度,代码示例以Java为主。读完你至少能搞清楚三件事:分片广播到底怎么分的、分片参数怎么用才不出错、生产环境有哪些坑必须提前防。
2. 分片广播模式的核心原理拆解
2.1 从路由策略说起:分片广播在调度链路中的位置
XXL-JOB 的调度中心(admin)向执行器(executor)下发任务时,有一个关键环节叫路由策略。路由策略决定了这次调度请求发给哪些执行器。常见的策略有第一个、最后一个、轮询、随机、一致性HASH、最不经常使用、最近最久未使用、故障转移、忙碌转移等,而分片广播是其中比较特殊的一种。
特殊在哪?其他路由策略本质上都是选一台执行器来执行任务,而分片广播是向所有在线执行器广播调度请求,每个执行器都会收到这次调度。但收到之后不是每个执行器都完整跑一遍业务逻辑,而是每个执行器拿到一个属于自己的分片序号和分片总数,根据这两个参数决定自己该处理哪部分数据。
你可以这样理解:调度中心是包工头,执行器是工人。普通路由策略是包工头挑一个工人去干活;分片广播是包工头把活拆成 N 份,通知所有工人来领,每个工人领一份。至于怎么拆、怎么领,靠的就是分片参数。
2.2 分片参数是怎么传递的:ShardingUtil 与 XxlJobHelper
执行器收到调度请求后,怎么拿到自己的分片序号和分片总数?XXL-JOB 提供了两种方式。
早期版本(2.2.x 及之前)用的是ShardingUtil:
ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo(); int index = shardingVO.getIndex(); // 当前分片序号,从 0 开始 int total = shardingVO.getTotal(); // 分片总数新版本(2.3.x 及之后)推荐用XxlJobHelper:
int index = XxlJobHelper.getShardIndex(); // 当前分片序号,从 0 开始 int total = XxlJobHelper.getShardTotal(); // 分片总数这两个工具类底层都是从调度请求的参数中解析出来的。调度中心在触发任务时,会把分片信息塞进任务参数里,执行器解析后放到 ThreadLocal 中,业务代码直接取就行。
注意:分片序号是从 0 开始的,不是从 1 开始。这个细节看起来不起眼,但在写取模逻辑的时候如果搞错了,会导致第一个分片永远拿不到数据或者数据错位。
2.3 分片总数到底等于多少:一个容易搞混的关键点
很多人以为分片总数是自己在配置里指定的,比如我想分 10 片就配 10。实际上在 XXL-JOB 的分片广播模式下,分片总数默认等于当前在线执行器的数量。也就是说,如果你部署了 5 台执行器,分片总数就是 5,每台执行器拿到一个 0 到 4 之间的序号。
这个设计有它的合理性:每台机器处理一份,天然负载均衡。但也带来一个问题——如果执行器数量动态变化(比如扩容或者某台机器挂了),分片总数会跟着变。今天 5 台机器分 5 片,明天扩容到 8 台就分 8 片。如果你的分片逻辑写得不够健壮,就可能出现数据重复处理或者遗漏。
那能不能固定分片数?可以,但需要绕一下。常见做法是不依赖执行器数量作为分片总数,而是自己在任务参数里指定一个固定的分片数,然后结合执行器的 IP 或序号做二次分配。这个后面在实战部分会详细讲。
2.4 分片广播的调度流程:一次完整的链路追踪
把整个流程串起来看,一次分片广播调度大致经历这几个步骤:
- 调度中心根据 Cron 表达式触发任务,查询当前在线的执行器列表。
- 调度中心向所有在线执行器广播调度请求,请求中携带分片总数(等于执行器数量)。
- 每个执行器收到请求后,根据自身在列表中的位置确定分片序号。
- 执行器将分片序号和分片总数放入 ThreadLocal,供业务代码通过
XxlJobHelper获取。 - 业务代码根据分片参数,从数据源中捞出属于自己的那部分数据并处理。
- 各执行器独立完成处理,向调度中心上报执行结果。
这里有个细节值得注意:分片序号的分配是在调度中心侧完成的,执行器只是被动接收。调度中心怎么知道哪台执行器对应哪个序号?它是按照执行器注册到调度中心的顺序来分配的。这个顺序在运行期间可能变化,所以不要假设某个 IP 永远对应某个固定的分片序号。
3. 分片逻辑怎么写才不出错
3.1 最常见的取模分片:原理与代码模板
分片逻辑最常用的就是取模。假设你有一批数据,每条数据有一个自增 ID 或者可以排序的字段,用 ID 对分片总数取模,余数等于当前分片序号的记录就归你处理。
@XxlJob("shardingJobHandler") public void shardingJob() { int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); // 查询所有待处理数据的 ID 列表(实际场景中建议分批查,不要一次全捞出来) List<Long> allIds = orderMapper.selectPendingIds(); // 过滤出属于当前分片的数据 List<Long> myIds = allIds.stream() .filter(id -> id % shardTotal == shardIndex) .collect(Collectors.toList()); // 处理属于自己的数据 for (Long id : myIds) { processOrder(id); } XxlJobHelper.log("分片 {}/{} 处理完成,共处理 {} 条", shardIndex, shardTotal, myIds.size()); }这段代码看起来简单,但有几个地方需要留意。第一,allIds如果数据量很大,一次性查出来会撑爆内存,实际生产建议用分页查询或者游标查询。第二,取模运算要求 ID 是数字类型,如果业务主键是字符串(比如 UUID),需要先转成数字哈希值再取模。第三,id % shardTotal的结果范围是 0 到 shardTotal-1,正好对应分片序号,这个没问题。
3.2 取模分片的隐患:数据倾斜与执行器数量变化
取模分片最大的问题是数据倾斜。如果 ID 不是均匀分布的,比如某些 ID 段的数据特别密集,那么对应的分片就会处理特别多的数据,其他分片早早跑完闲着。我遇到过一种情况:订单 ID 是按时间递增的,最近一个月的订单 ID 集中在某个区间,结果负责那个区间的分片跑了二十分钟,其他分片两分钟就结束了。
另一个隐患是执行器数量变化。假设你原本 4 台机器,分片总数是 4,ID 为 100 的数据由分片 0 处理(100 % 4 = 0)。后来扩容到 5 台,分片总数变成 5,100 % 5 = 0,还是分片 0 处理,看起来没问题。但 ID 为 101 的数据,原来 101 % 4 = 1 由分片 1 处理,现在 101 % 5 = 1 还是分片 1。再试一个:ID 为 103,原来 103 % 4 = 3 由分片 3 处理,现在 103 % 5 = 3 还是分片 3。好像变化不大?那是因为我挑的数字比较巧。换一个:ID 为 102,原来 102 % 4 = 2,现在 102 % 5 = 2,也没变。再换:ID 为 104,原来 104 % 4 = 0,现在 104 % 5 = 4,分片变了。
所以执行器数量变化会导致部分数据的分片归属发生变化。如果任务是一次性的(比如处理完就标记状态),问题不大;但如果是周期性的、依赖上一次处理结果的,就可能出问题。解决办法后面会讲。
3.3 更稳健的分片方式:按数据段切分而非取模
为了避免取模带来的倾斜和数量变化问题,我后来改用按数据段切分的方式。思路是:先查出待处理数据的最小 ID 和最大 ID,然后按分片总数把 ID 范围均分成 N 段,每个分片处理自己那一段。
@XxlJob("rangeShardingJobHandler") public void rangeShardingJob() { int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal(); Long minId = orderMapper.selectMinPendingId(); Long maxId = orderMapper.selectMaxPendingId(); if (minId == null || maxId == null) { XxlJobHelper.log("没有待处理数据"); return; } long rangeSize = (maxId - minId + 1) / shardTotal; long startId = minId + shardIndex * rangeSize; long endId = (shardIndex == shardTotal - 1) ? maxId : startId + rangeSize - 1; XxlJobHelper.log("分片 {}/{} 处理 ID 范围 [{}, {}]", shardIndex, shardTotal, startId, endId); // 分批查询并处理 int pageSize = 500; long cursor = startId; while (cursor <= endId) { List<Order> orders = orderMapper.selectByRange(cursor, endId, pageSize); if (orders.isEmpty()) break; for (Order order : orders) { processOrder(order); } cursor = orders.get(orders.size() - 1).getId() + 1; } }这种方式的好处是每个分片处理的数据量大致均匀(前提是 ID 分布均匀),而且不依赖取模运算,执行器数量变化时只需要重新计算范围即可。缺点是如果 ID 有空洞(比如删除了很多数据),某些分片可能实际处理的数据很少。但总体来说比取模稳健得多。
3.4 分片数固定 vs 动态:如何根据业务场景选择
前面提到分片总数默认等于执行器数量。这在大多数场景下够用,但有两种情况需要固定分片数:
第一种是执行器数量少于期望的并行度。比如你只有 2 台机器,但希望分成 10 片来跑,充分利用每台机器的多线程能力。这时候可以在任务参数里指定分片数为 10,然后每台执行器内部再起线程池处理多个分片。
第二种是执行器数量频繁变化。比如用了弹性伸缩,机器数量忽多忽少。固定分片数可以避免分片归属频繁变动导致的数据问题。
固定分片数的实现方式通常是:在任务参数中传入一个fixedShardTotal,业务代码读取这个值作为分片总数,然后结合执行器的 IP 哈希或者序号做二次分配。具体代码这里不展开,核心思路就是“调度中心的分片总数”和“业务逻辑的分片总数”解耦。
4. 生产环境实战:从配置到上线
4.1 执行器配置与分片参数获取的完整示例
先看执行器侧的配置。在application.properties或application.yml中配置调度中心地址和执行器信息:
xxl.job.admin.addresses=http://your-admin-host:8080/xxl-job-admin xxl.job.executor.appname=your-app-name xxl.job.executor.port=9999 xxl.job.executor.logpath=/data/applogs/xxl-job/jobhandler xxl.job.executor.logretentiondays=30然后在 Spring 配置类中注册执行器:
@Configuration public class XxlJobConfig { @Value("${xxl.job.admin.addresses}") private String adminAddresses; @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.setLogPath("/data/applogs/xxl-job/jobhandler"); executor.setLogRetentionDays(30); return executor; } }任务处理器就是前面展示的@XxlJob注解方法。在调度中心新建任务时,路由策略选择分片广播,Cron 表达式按业务需求配置,任务参数可以留空(如果不需要固定分片数)。
4.2 分片任务的日志排查:XxlJobHelper.log 的正确用法
分片任务最头疼的问题之一是排查。5 台机器同时跑,每台机器的日志分散在不同文件里,出了问题怎么定位?XXL-JOB 提供了XxlJobHelper.log()方法,它会把日志写到调度中心可以查看的执行日志中。
XxlJobHelper.log("分片 {}/{} 开始处理,数据范围 [{}, {}]", shardIndex, shardTotal, startId, endId);这个日志的好处是,在调度中心的任务日志页面可以直接看到每个分片的执行情况,不用一台台机器去翻日志文件。但要注意,XxlJobHelper.log()写日志是有性能开销的,不要在循环里每条数据都写,建议按批次写或者只写关键节点。
实操心得:我通常会在分片任务开始时打一条日志记录分片参数,结束时打一条记录处理条数和耗时。中间如果出错,用 try-catch 捕获后通过
XxlJobHelper.log()记录异常堆栈。这样在调度中心就能看到完整的执行轨迹。
4.3 分片任务超时与阻塞:调度周期怎么设才合理
分片广播模式下,所有执行器是并行跑的,整体耗时取决于最慢的那个分片。如果某个分片因为数据倾斜跑了很久,而调度周期又比较短,就会出现上一次任务还没跑完、下一次调度又来了的情况。
XXL-JOB 默认的阻塞处理策略是单机串行,意思是同一台执行器上的同一个任务,如果上一次还没执行完,下一次调度会排队等待。这个策略在分片场景下可能导致任务堆积。另一种策略是丢弃后续调度,即上一次没跑完就跳过这一次。还有一种覆盖之前调度,直接终止上一次执行。
我的建议是:分片任务的调度周期要留足余量,至少是平均执行时间的 3 倍以上。同时开启任务超时配置,比如设置超时时间为 30 分钟,超过就自动失败,避免无限阻塞。另外,在业务代码里加一个执行时长监控,如果发现某次执行明显变慢,及时告警。
4.4 动态扩容场景下的分片处理:一个真实案例
前面提到执行器数量变化会导致分片归属变化。我遇到过一个真实案例:一个数据同步任务,原本 3 台执行器,分片总数 3。某天运维扩容到 5 台,分片总数变成 5。结果发现有一部分数据被重复同步了。
原因是什么?任务逻辑是“查询状态为待同步的数据,同步后更新状态”。扩容前,ID 为 100 的数据由分片 1 处理(100 % 3 = 1),处理完状态更新为已同步。扩容后,分片总数变成 5,ID 为 100 的数据变成由分片 0 处理(100 % 5 = 0)。但此时 ID 为 100 的数据状态已经是已同步,按理说不应该被再次查询出来。问题出在查询条件上——查询的是“状态为待同步”,但更新状态和查询之间有时间窗口,扩容恰好发生在这个窗口内,导致数据被两个分片同时捞到。
解决办法有两个:一是用数据库行锁或者乐观锁保证幂等,二是固定分片数避免归属变化。我最后选了第二种,在任务参数里固定分片数为 3,扩容到 5 台后,每台执行器根据 IP 哈希决定自己处理哪几个分片,这样分片归属不变,问题解决。
5. 常见问题与排查技巧实录
5.1 分片任务只在一台机器上执行:路由策略配错了吗
这是新手最常遇到的问题:明明部署了多台执行器,任务却只在一台机器上跑。排查思路如下:
| 排查项 | 检查方法 | 常见原因 |
|---|---|---|
| 路由策略 | 调度中心任务编辑页查看 | 误选为“第一个”或其他单机策略 |
| 执行器在线状态 | 调度中心执行器管理页查看 | 其他执行器未注册成功或已下线 |
| 执行器 appname | 各机器配置文件对比 | appname 不一致导致注册到不同分组 |
| 网络连通性 | 执行器日志查看注册结果 | 执行器无法访问调度中心 |
我遇到过一种隐蔽情况:两台执行器的 appname 配得一样,但其中一台的端口被占用,启动时自动换了端口,导致注册信息混乱。后来在启动脚本里加了端口检查才解决。
5.2 分片数据重复处理:取模逻辑的边界陷阱
数据重复处理通常有几个原因。一是前面说的执行器数量变化导致分片归属变化。二是取模逻辑写错,比如用了id % shardTotal == shardIndex + 1这种偏移。三是数据查询没有加状态过滤,导致已经处理过的数据被再次捞出来。
排查方法:在分片任务开始时,把分片参数和查询条件都打到日志里。然后手动验证几条数据,看它们是否只被一个分片处理。如果发现重复,先检查取模公式,再检查执行器数量是否变化过,最后检查数据状态更新是否及时。
避坑技巧:在数据表上加一个
shard_mark字段,记录这条数据被哪个分片处理过。虽然有点冗余,但排查问题时非常有用。
5.3 分片任务执行时间过长:如何定位慢分片
分片任务整体耗时取决于最慢的分片。定位慢分片的方法是在每个分片结束时记录耗时,然后在调度中心日志里对比。如果发现某个分片 consistently 比其他分片慢很多,大概率是数据倾斜。
解决数据倾斜的思路:如果用的是取模分片,考虑改成范围分片;如果已经是范围分片,检查 ID 分布是否均匀,必要时按数据量而非 ID 范围来切分。另一种思路是动态分片——先统计每个分片待处理的数据量,然后让数据量少的分片“支援”数据量多的分片。这个实现比较复杂,一般场景用不上。
5.4 执行器扩容后分片数不对:注册中心缓存问题
执行器扩容后,调度中心需要感知到新的执行器上线。XXL-JOB 的执行器注册是心跳机制,默认 30 秒一次。扩容后如果立即触发任务,可能调度中心还没感知到新执行器,分片总数还是旧的。
解决办法:扩容后等一个心跳周期再触发任务,或者在调度中心手动刷新执行器列表。另外,如果执行器下线,调度中心也要等心跳超时才会摘除,这期间分片总数可能包含已下线的执行器,导致部分分片没有执行器处理。所以分片任务最好配置失败重试和告警。
5.5 分片任务与数据库连接池:并发查询的隐藏风险
分片广播模式下,多台执行器同时查询数据库,如果每台执行器还起了多线程,数据库连接池可能瞬间被打满。我遇到过执行器报“无法获取数据库连接”的错误,排查后发现是分片任务并发度太高。
解决办法:控制每台执行器的并发线程数,数据库连接池大小要大于“执行器数量 × 每执行器并发线程数”。另外,查询尽量走索引,避免全表扫描导致锁表。
6. 分片广播模式的适用边界与替代方案
分片广播不是万能的。它适合数据可以水平切分、各分片之间无依赖、处理结果可以独立提交的场景。比如批量数据处理、对账、报表生成、数据同步等。
不适合的场景包括:任务之间有严格的先后顺序依赖、需要全局聚合结果、数据无法切分(比如必须全量加载到内存计算)。这些场景用分片广播反而会增加复杂度。
如果分片广播不适用,可以考虑的替代方案有:用消息队列做任务分发,每个消费者处理一部分;或者用 MapReduce 式的两阶段处理,先分片计算再汇总。XXL-JOB 本身也支持子任务和依赖配置,可以组合使用。
我在实际项目中的体会是,分片广播最大的价值不是“快”,而是“可扩展”。单机跑十分钟的任务,分片后可能只快两三倍(因为还有调度开销和数据切分开销),但当数据量增长十倍时,你只需要加机器就行,不用改代码。这种弹性才是它真正的意义。
最后分享一个小技巧:分片任务的测试不要只在单机上测。单机测试时分片总数是 1,很多分片逻辑的边界问题暴露不出来。至少起两个执行器实例,把分片总数变成 2,才能验证取模、范围切分、数据归属这些逻辑是否正确。我见过太多人单机测试通过、上线后数据重复的案例了。