XXL-JOB分片广播灵活控制执行节点数
阿昌 Java小菜鸡

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。

那要做的事情很简单:

  1. 通过任务参数传入一个 maxNode(最大执行节点数),比如 8
  2. 在 JobHandler 里判断:如果当前分片序号 >= maxNode,则跳过,不执行
  3. 对于需要执行的节点,把 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,让分片查询正确
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 抽象类,把上面那套判断逻辑抽到了基类里。

业务方只需要:

  1. 继承 ShardingJobHandler 而不是 IJobHandler
  2. 实现 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()) {
// 当前分片序号 >= maxNode,跳过
if (ShardingUtil.getShardingVo().getIndex() + 1 > maxNode) {
log.info("分片设置了最大Node数:{} 总分片数: {} 当前分片为第:{} 跳过执行",
maxNode,
ShardingUtil.getShardingVo().getTotal(),
ShardingUtil.getShardingVo().getIndex() + 1);
return ReturnT.SUCCESS;
} else {
// 修正 shardTotal 后回调子类
ShardingUtil.getShardingVo().setTotal(maxNode);
return executeWithSharding(param);
}
}
}
return executeWithSharding(param);
}

/**
* 分片执行——子类只需关注业务逻辑,shardTotal 已自动修正
*/
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


五、注意事项

  1. 参数格式:内部封装版要求参数是 JSON,如果老任务的参数是纯数字字符串,改造时要做兼容处理。
  2. maxNode 范围maxNode 必须在 [1, shardTotal] 之间,超出范围直接报错。
  3. 分片查询 SQL:关键在 MOD(id, shardTotal) = shardIndex,如果用其他分片字段,确保分布均匀。
  4. 跳过的节点返回 SUCCESS:千万别返回 FAIL,否则 XXL-JOB 会触发重试,导致日志刷屏。

六、总结

这篇文章介绍了两种利用 XXL-JOB 分片广播灵活控制执行节点数的方式:

方式 适用场景 复杂度
原生 API 实现 项目没有统一封装,或不想引入内部依赖 每个 JobHandler 需重复编写
ShardingJobHandler 封装 引入 对应包 后即可使用 业务代码极简

核心思路都一样:

通过任务参数指定 maxNode,结合分片广播的路由机制,跳过不需要执行的节点,同时修正 shardTotal 以保证分片查询正确。

这个方案很轻量,不需要改 XXL-JOB 源码,也没有引入额外中间件,纯靠任务参数和分片序号判断就能实现。

 请作者喝咖啡