ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

XXL-JOB分片广播模式实战:原理、分片逻辑与生产避坑指南

XXL-JOB分片广播模式实战:原理、分片逻辑与生产避坑指南 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 及之前用的是ShardingUtilShardingUtil.ShardingVO shardingVO ShardingUtil.getShardingVo(); int index shardingVO.getIndex(); // 当前分片序号从 0 开始 int total shardingVO.getTotal(); // 分片总数新版本2.3.x 及之后推荐用XxlJobHelperint 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 列表实际场景中建议分批查不要一次全捞出来 ListLong allIds orderMapper.selectPendingIds(); // 过滤出属于当前分片的数据 ListLong 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 台机器分片总数是 4ID 为 100 的数据由分片 0 处理100 % 4 0。后来扩容到 5 台分片总数变成 5100 % 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) { ListOrder 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.addresseshttp://your-admin-host:8080/xxl-job-admin xxl.job.executor.appnameyour-app-name xxl.job.executor.port9999 xxl.job.executor.logpath/data/applogs/xxl-job/jobhandler xxl.job.executor.logretentiondays30然后在 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处理完状态更新为已同步。扩容后分片总数变成 5ID 为 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才能验证取模、范围切分、数据归属这些逻辑是否正确。我见过太多人单机测试通过、上线后数据重复的案例了。
返回列表