XXL-JOB 分片广播灵活控制执行节点数
Hi,我是阿昌,今天学习记录下 XXL-JOB 分片广播中一个很实用的技巧——灵活控制执行节点数量。
一、问题背景
一个真实场景:
线上部署了 20 个消费服务节点,但是数据库只有 8 个实例。
希望对这 8 个库分别用一个节点消费,其余 12 个节点保持空闲,不参与任务处理。
也就是说,虽然 XXL-JOB 的路由策略「分片广播」会把任务广播到所有注册的 executor 节点,但当前场景并不希望所有节点都干活。
这个需求其实很常见:
- 节点数多于数据库分片数,想做到 1:1 对应
- 灰度发布时,只想让部分节点执行新逻辑
- 压测时,动态扩大或缩小执行节点数量
核心诉求:总节点数是 N,但我只想让其中 M 个节点真正执行任务(M ≤ N)。
![image]()
二、核心思路
XXL-JOB 的分片广播会给每个 executor 分配一个分片序号(shardIndex)和分片总数(shardTotal),比如 20 个节点,序号分别是 0~19,总数是 20。
那要做的事情很简单:
- 通过任务参数传入一个
maxNode(最大执行节点数),比如 8
- 在 JobHandler 里判断:如果当前分片序号
>= maxNode,则跳过,不执行
- 对于需要执行的节点,把
shardTotal 改为 maxNode,然后用 mod(id, shardTotal) = shardIndex 做分片查询
这样一来,只有序号 07 的节点会执行,序号 819 的节点直接返回成功,全程无报错、无重试。
下面分别介绍两种实现方式。
三、方案一:原生 XXL-JOB 实现
用 XXL-JOB 原生 API 实现
JobHandler 代码
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38
| @XxlJob("demoJobHandler") public void execute() { String param = XxlJobHelper.getJobParam(); if (StringUtils.isBlank(param)) { XxlJobHelper.log("任务参数为空"); XxlJobHelper.handleFail(); return; }
int executeNodeNum = Integer.valueOf(param);
int shardIndex = XxlJobHelper.getShardIndex(); int shardTotal = XxlJobHelper.getShardTotal();
if (executeNodeNum <= 0 || executeNodeNum > shardTotal) { XxlJobHelper.log("执行任务节点数量取值范围[1, 节点总数]"); XxlJobHelper.handleFail(); return; }
if (shardIndex > (executeNodeNum - 1)) { XxlJobHelper.log("当前分片 {} 无需执行", shardIndex); XxlJobHelper.handleSuccess(); return; }
shardTotal = executeNodeNum;
process(shardIndex, shardTotal);
XxlJobHelper.handleSuccess(); }
|
分片查询 SQL 示例
1 2 3 4 5
| SELECT field1, field2 FROM table_name WHERE ... AND MOD(id, #{shardTotal}) = #{shardIndex} ORDER BY id LIMIT #{rows};
|
利用 MOD(id, shardTotal) 取模,保证不同分片处理的数据互不重叠。
优缺点
- 优点:零依赖,纯 XXL-JOB API
- 缺点:每个 JobHandler 都要写一遍判断逻辑,参数校验也要自己处理
四、方案二: ShardingJobHandler 封装
基于 XXL-JOB 封装了一个 ShardingJobHandler 抽象类,把上面那套判断逻辑抽到了基类里。
业务方只需要:
- 继承
ShardingJobHandler 而不是 IJobHandler
- 实现
executeWithSharding 方法——这个方法被回调时,shardTotal 已经是修正后的值了
核心源码
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43
| public abstract class ShardingJobHandler extends TraceIJobHandler {
private static final String MAX_NODE = "maxNode";
public ShardingJobHandler(ApplicationContext beanFactory) { super(beanFactory); }
@Override protected ReturnT<String> executeWithTrace(String param) throws Exception { JSONObject jsonObject; try { jsonObject = JSON.parseObject(param); } catch (Exception e) { log.error("分片参数设置必须为json字符串"); return new ReturnT<>(ReturnT.FAIL_CODE, "分片策略任务参数设置必须为json字符串"); }
if (ShardingUtil.getShardingVo().getTotal() > 1 && jsonObject != null) { Integer maxNode = jsonObject.getInteger(MAX_NODE); if (maxNode != null && maxNode < ShardingUtil.getShardingVo().getTotal()) { if (ShardingUtil.getShardingVo().getIndex() + 1 > maxNode) { log.info("分片设置了最大Node数:{} 总分片数: {} 当前分片为第:{} 跳过执行", maxNode, ShardingUtil.getShardingVo().getTotal(), ShardingUtil.getShardingVo().getIndex() + 1); return ReturnT.SUCCESS; } else { ShardingUtil.getShardingVo().setTotal(maxNode); return executeWithSharding(param); } } } return executeWithSharding(param); }
protected abstract ReturnT<String> executeWithSharding(String param) throws Exception; }
|
业务使用示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
| @Component public class InventoryUploadJob extends ShardingJobHandler {
@Override protected ReturnT<String> executeWithSharding(String param) throws Exception { int shardIndex = ShardingUtil.getShardingVo().getIndex(); int shardTotal = ShardingUtil.getShardingVo().getTotal();
List<Data> dataList = queryByShard(shardIndex, shardTotal); processData(dataList);
return ReturnT.SUCCESS; } }
|
简洁多了,业务代码完全不用关心”要不要跳过”和”修正 shardTotal”这些事。
使用指南
- 依赖:项目需要引入
qianniu-core 包
- JobHandler:继承
ShardingJobHandler,实现 executeWithSharding 即可
- 参数格式:分片广播任务的传参必须为 JSON 字符串,例子:
{"maxNode": 8}
- 注意兼容:改代码时注意,原来传普通字符串的,改用 JSON 会不兼容
配置示例
XXL-JOB 调度中心配置:
- 路由策略:分片广播
- 任务参数:
{"maxNode": 1} —— 表示最多只允许 1 个 pod 执行
![image]()
当某个分片被跳过时,日志输出:
1
| 分片设置了最大Node数:1 总分片数: 2 当前分片为第:2 跳过执行
|
![image]()
五、注意事项
- 参数格式:内部封装版要求参数是 JSON,如果老任务的参数是纯数字字符串,改造时要做兼容处理。
- maxNode 范围:
maxNode 必须在 [1, shardTotal] 之间,超出范围直接报错。
- 分片查询 SQL:关键在
MOD(id, shardTotal) = shardIndex,如果用其他分片字段,确保分布均匀。
- 跳过的节点返回 SUCCESS:千万别返回 FAIL,否则 XXL-JOB 会触发重试,导致日志刷屏。
六、总结
这篇文章介绍了两种利用 XXL-JOB 分片广播灵活控制执行节点数的方式:
| 方式 |
适用场景 |
复杂度 |
| 原生 API 实现 |
项目没有统一封装,或不想引入内部依赖 |
每个 JobHandler 需重复编写 |
| ShardingJobHandler 封装 |
引入 对应包 后即可使用 |
业务代码极简 |
核心思路都一样:
通过任务参数指定 maxNode,结合分片广播的路由机制,跳过不需要执行的节点,同时修正 shardTotal 以保证分片查询正确。
这个方案很轻量,不需要改 XXL-JOB 源码,也没有引入额外中间件,纯靠任务参数和分片序号判断就能实现。