凌晨两点被电话叫醒是什么体验,我在接手那个用户标签刷新任务之后,连续体会了一个星期。任务本身不复杂:每天全量刷新用户标签表,1.2亿行数据,需要关联订单表、登录日志、客服记录三个维度做聚合计算。最早是单机跑,执行时间稳定在一个半小时左右,隔三差五因为内存问题挂在半夜,值班电话比业务告警还勤。后来把方案改成XXL-Job的分片广播,数据按用户ID切成段,交给多个执行器并行处理,整体耗时从90分钟压到20多分钟,告警基本消失。这篇文章就把我从场景判断、方案设计、编码落地到本地部署验证、排障调优的完整过程写出来,给还在纠结"到底要不要用分片广播"的团队一个可参考的答案。
先说一个容易混淆的点:分片广播不是"把全量数据派发给每台机器,每台机器都跑一遍",而是把一个任务拆成多个分片,每台执行器只处理自己负责的那一段。理解了这个,后面所有设计和排障都顺了。
1. 超过单机能力的数据任务,为什么答案偏偏是分片广播
1.1 三种任务模型的对比:单机、轮询、分片广播
遇到海量数据定时任务,先别急着上框架,得先看清手里的牌。拿我那个1.2亿行的用户标签刷新来说,摆在面前的三条路,我都试过或者认真评估过。
第一种是单机任务。代码简单,逻辑直观,把聚合SQL一写,循环分批更新就完事。但它的天花板非常明显:一台机器的内存、CPU、数据库连接数都是有限的。1.2亿行数据关联三张业务表,光是查询和计算的内存压力,就经常把堆内存顶到上限。就算你分批处理避免OOM,执行时间也是线性往上走——数据翻倍,时间翻倍,没有意外。更麻烦的是,单机任务一旦进程崩溃,整个任务从头再来,前面跑的几个小时全部作废。
第二种是XXL-Job的轮询/一致性哈希路由。很多人误以为"轮询就是多机并行处理同一个任务",其实不是。路由策略解决的是"多个任务实例怎么分给执行器",不是"一个任务怎么拆成多份"。轮询模式下,调度中心每次只把任务交给一个执行器,这个执行器依然是单机全量跑一遍。如果三个执行器都在线,轮询的效果是"今天A跑全量,明天B跑全量"——它不是并行,是轮流被折磨。一致性哈希同理,本质还是单机处理。
第三种才是分片广播。调度中心在任务触发时,把当前所有在线执行器实例数统计出来,然后把这个任务广播给每一个实例,同时告诉每个实例"你是第几个,总共有几个"。每个执行器拿到自己的序号之后,只处理按这个序号划分的那一段数据。所有执行器同时开工,处理完成后各自上报。这才是真正把"单机跑不完"变成了"多机分着跑"。
1.2 什么场景值得用分片广播
我总结了一个自己的判断清单,命中三条以上基本就可以考虑分片广播:
- 数据量级上千万起步,单机跑批时间超过30分钟,或者存在OOM风险;
- 数据本身可以横向切分,比如按主键ID、时间范围、租户维度切段后互不依赖;
- 任务允许最终一致,不需要所有分片同时提交成功才对外提供服务;
- 业务逻辑可以做成幂等,同一段数据重跑一遍不会产生脏数据;
- 团队手里有两台以上的执行器资源,能承受并行带来的数据库压力。
1.3 不适合分片广播的场景
分片广播不是万能药。拿它处理实时性要求高的数据同步,或者强一致性的账务计算,就属于用错了工具。分片广播天然是"各分片独立提交、最后汇总"的模型,分片之间没有分布式事务,如果某个分片失败,其他分片已经提交的数据不会自动回滚。另外,如果任务本身就是单条记录级别的操作、跑一次只要几秒钟,为了分片广播引入多执行器和状态管理,成本反而大于收益。
2. 一次分片广播任务的完整生命周期:调度中心、执行器与业务代码的三方配合
2.1 先搞清楚XXL-Job的两个核心角色
要用好分片广播,必须理解XXL-Job的两个核心进程。
调度中心(xxl-job-admin)是大脑,负责维护定时任务配置、到了触发时间发起调度、记录调度日志、处理失败重试和告警。它不执行业务代码,只负责"叫醒"执行器。
执行器(xxl-job-executor)是干活的进程,真正的JobHandler业务代码跑在这里。一个执行器进程可以注册多个JobHandler,比如一个数据清洗的、一个报表生成的。执行器启动后会主动向调度中心注册自己的地址,调度中心通过在线列表知道当前有几个执行器可用。
分片广播的特殊之处在于:调度中心在触发任务的那一刻,会把所有在线执行器实例作为一个整体,把任务调用请求同时发给每一个实例,并在请求里携带两个参数——分片总数(shardTotal)和当前分片序号(shardIndex)。
2.2 分片参数在代码里怎么拿
在JobHandler里通过ShardingUtil可以拿到这两个参数,这是最经典的方式:
ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo(); int shardIndex = shardingVO.getIndex(); // 当前分片序号,从0开始 int shardTotal = shardingVO.getTotal(); // 分片总数,= 在线执行器数量新版本的XXL-Job更推荐用XxlJobContext,因为ShardingUtil在后续版本中有过标记废弃的调整:
XxlJobContext context = XxlJobContext.getXxlJobContext(); int shardIndex = context.getShardIndex(); int shardTotal = context.getShardTotal();这里有一个容易踩的坑:如果任务的路由策略不是"分片广播",getShardingVo()返回的是null,直接调用会空指针。所以代码里最好判空,或者在路由策略配置时就定死分片广播。
2.3 执行器动态伸缩时,分片参数怎么变
分片总数不是你在任务配置里写死的,而是调度中心在每次触发时根据在线执行器数量动态计算的。你启动了两个执行器,这次任务shardTotal就是2;过两天加了第三台机器,下次触发shardTotal自动变成3。这带来两个连锁反应:
一是任务触发瞬间不在线的执行器,这次不参与分片,它原本负责的数据段这次不会有人处理。如果恰好一台机器在凌晨服务重启,那它名下的数据段就会漏掉。这个风险要在设计上兜住,后面第五节详细说。
二是分片参数是"本次触发"快照,不是动态漂移的。处理过程中某台机器挂了,已经分出去的任务不会自动转移给其他机器,只能靠失败重试或人工干预。理解了这一点,就不会在监控上犯错——毕竟XXL-Job的失败重试是整任务维度,不是某一个分片维度。
2.4 常见误解:分片广播到底广播了什么
"广播"两个字特别容易让人误解。它广播的不是数据,而是触发信号加分片参数。每个执行器收到信号后自己去数据库捞属于自己那一段的数据。所以,实际的数据查询压力是分散在多个执行器节点上的,但如果你的分片逻辑写成了"每台机器都全表扫一遍,只取其中一部分",那数据库全表扫描的压力依然存在,只是由一台机器扛变成了N台机器一起扛。好的分片策略一定是从数据访问路径上就切开,让每个执行器通过索引范围扫描只碰自己那一段。
3. 实战:亿级用户标签表如何用分片广播做全量刷新
3.1 任务整体设计与分片策略选择
场景再具体一点。user_tag表大概1.2亿行,主键user_id,需要根据订单、登录、客服记录重新计算每个用户的标签,更新到tag_json字段。
分片维度看起来有三个方案可以选择,我分别评估过:
第一个方案是按user_id取模,也就是:
SELECT ... FROM user_tag WHERE user_id % #{shardTotal} = #{shardIndex}优点是实现简单、数据均匀,缺点是取模运算会让MySQL放弃主键索引,被迫全表扫描。1.2亿行的表全表扫描,即便每个分片只取其中1/N行,扫描成本依然全部压在数据库上,分片越多,数据库越痛苦。
第二个方案是按主键值范围切分。先查MIN(user_id)和MAX(user_id),把整个值域按分片数量等分,每个执行器只负责一段ID区间:
SELECT ... FROM user_tag WHERE user_id BETWEEN #{rangeStart} AND #{rangeEnd}这个方案能让查询走主键索引的range scan,每个节点只扫自己负责的那一段,数据库压力是真正的"N分之一的成本"。缺点是如果user_id分布极不均匀,可能出现某个分片的记录数远多于其他分片。
第三个方案是先统计再按数据量切分。比如先跑一个count按ID区间分组,摸清楚分布后,把数据量均匀地分到每个分片。解决倾斜问题,但要多跑一次统计查询,实现复杂度高一些。
对于用户ID这种自增主键场景,分布基本均匀,我最终选了方案二——ID范围切分。如果你们场景里的切分键分布很不均匀,可以先用方案三的思路做边界修正。
3.2 核心代码:分片参数获取与数据范围计算
JobHandler的核心逻辑分四步:拿到分片参数、计算当前分片负责的ID区间、分批拉取并处理、记录处理日志。代码骨架如下:
@Component public class UserTagRefreshJobHandler { @Resource private UserTagMapper userTagMapper; @Resource private BatchJobLogMapper batchJobLogMapper; @XxlJob("userTagRefreshJob") public void execute() throws Exception { // 1. 获取分片参数 int shardIndex = XxlJobContext.getXxlJobContext().getShardIndex(); int shardTotal = XxlJobContext.getXxlJobContext().getShardTotal(); XxlJobHelper.log("分片任务启动, shardIndex={}, shardTotal={}", shardIndex, shardTotal); // 2. 计算本分片负责的ID区间 Long minId = userTagMapper.selectMinUserId(); Long maxId = userTagMapper.selectMaxUserId(); if (minId == null || maxId == null) { XxlJobHelper.log("user_tag表为空, 任务结束"); return; } long chunkSize = (maxId - minId) / shardTotal + 1; long rangeStart = minId + chunkSize * shardIndex; long rangeEnd = (shardIndex == shardTotal - 1) ? maxId : Math.min(rangeStart + chunkSize - 1, maxId); XxlJobHelper.log("本分片处理ID区间: [{}, {}]", rangeStart, rangeEnd); // 3. 检查状态表,避免重复处理 BatchJobLog jobLog = batchJobLogMapper.selectByJobAndShard("userTagRefreshJob", shardIndex); if (jobLog != null && jobLog.getStatus() == 1) { XxlJobHelper.log("本分片已完成处理, 跳过, processCount={}", jobLog.getProcessCount()); return; } // 4. 分批处理 BatchJobLog newLog = createJobLog("userTagRefreshJob", shardIndex, rangeStart, rangeEnd); long processCount = processByRange(rangeStart, rangeEnd); finishJobLog(newLog, processCount); } }3.3 分批拉取与批量更新:不要一次性把所有数据load进内存
海量数据处理的死法是OOM,所以绝不能把整段ID区间的数据一次性查出来。每批拉一两千条,处理完提交,再拉下一批。我用的是游标步进的方式:
private long processByRange(Long rangeStart, Long rangeEnd) { long cursor = rangeStart; int batchSize = 2000; long processCount = 0; while (cursor <= rangeEnd) { List<UserTag> batch = userTagMapper.selectUserTagByRange(cursor, cursor + batchSize); if (batch.isEmpty()) { break; } // 根据订单、登录、客服记录重新计算标签 List<UserTag> updated = batch.stream() .map(this::recalculateTag) .collect(Collectors.toList()); // 批量更新,每批一个短事务 userTagMapper.batchUpdateTag(updated); processCount += batch.size(); cursor = cursor + batchSize; if (processCount % 10000 == 0) { XxlJobHelper.log("已处理记录数: {}", processCount); } } return processCount; }批量更新我建议用INSERT ... ON DUPLICATE KEY UPDATE或者UPDATE ... CASE WHEN拼SQL,不要一条一条update。单条更新在千万级数据上完全是灾难,连接交互次数直接打满。
3.4 幂等与断点续跑:批处理状态表
分片任务跑在分布式环境里,失败重试、网络抖动、节点重启都是常态。要让任务"失败了能重跑、重跑了不出错",必须设计幂等。我的做法是加一张批次状态表:
CREATE TABLE batch_job_log ( id BIGINT AUTO_INCREMENT PRIMARY KEY, job_name VARCHAR(64) NOT NULL, shard_index INT NOT NULL, range_start BIGINT NOT NULL, range_end BIGINT NOT NULL, status TINYINT NOT NULL DEFAULT 0 COMMENT '0-处理中 1-成功 2-失败', process_count BIGINT DEFAULT 0, start_time DATETIME, finish_time DATETIME, UNIQUE KEY uk_job_shard (job_name, shard_index) );每个分片开始前先查这个表,如果状态已经是成功,直接跳过。任务失败重试时,已经成功的分片不会重复处理,还没处理完的分片接着跑。配合主键ID的范围条件,天然支持断点续跑。比如某个分片处理到一半进程挂了,重试时会从状态表看到status=0,再从rangeStart开始跑,已处理的部分虽然会重新计算一遍,但由于标签计算是覆盖式更新,最终结果不会被破坏,这就是幂等的意义。
3.5 日志与异常处理
分片任务一定要在日志里带上shardIndex,否则出问题的时候根本不知道这条日志来自哪个执行器、哪个分片。另外,XXL-Job提供了XxlJobHelper.log方法,日志会同步到调度中心的后台,可以在admin界面直接查看,比翻服务器日志方便得多。业务处理过程中的异常要捕获并记录,但不要吞掉,否则分片被标记为成功,数据却漏处理了。
4. 分片任务在本地如何部署验证:从拉源码到双执行器并行
4.1 本地部署的意义
生产环境搭建XXL-Job有专门的运维流程,但本地部署一套是理解它运行机制最快的方式,也是验证分片广播效果的必要步骤。拉一套源码起来跑一遍,比自己看十篇文档都管用。下面以xxl-job 2.4.0为例,完整走一遍。
环境要求:JDK 8+、Maven 3.6+、MySQL 5.7+/8.0。先从GitHub拉取源码:
git clone https://github.com/xuxueli/xxl-job.git cd xxl-job mvn clean package -DskipTests4.2 启动调度中心
创建数据库并导入初始化脚本:
CREATE DATABASE IF NOT EXISTS xxl_job DEFAULT CHARACTER SET utf8mb4; USE xxl_job; SOURCE xxl-job/doc/db/tables_xxl_job.sql;修改admin模块的配置,文件在xxl-job-admin/src/main/resources/application.properties,重点确认这几项:
server.port=8080 spring.datasource.url=jdbc:mysql://127.0.0.1:3306/xxl_job?useSSL=false&serverTimezone=Asia/Shanghai spring.datasource.username=root spring.datasource.password=123456 xxl.job.accessToken=default_token然后启动:
java -jar xxl-job-admin/target/xxl-job-admin-2.4.0.jar浏览器访问 http://localhost:8080/xxl-job-admin,初始账号admin/123456。
4.3 执行器接入与双节点验证
在示例执行器工程xxl-job-executor-samples/xxl-job-executor-sample-springboot里,把我们上面写的UserTagRefreshJobHandler放进去。然后修改它的application.properties:
xxl.job.admin.addresses=http://localhost:8080/xxl-job-admin xxl.job.accessToken=default_token xxl.job.executor.appname=xxl-job-executor-sample xxl.job.executor.port=9999 xxl.job.executor.logpath=logs/xxl-job/jobhandler启动第一个执行器实例。然后再起一个实例,注意把执行器端口改掉,否则两个进程端口冲突:
# 第二个实例指定不同端口 java -jar xxl-job-executor-sample-springboot-2.4.0.jar --xxl.job.executor.port=9998两个执行器都启动后,去admin后台完成三件事:
- 执行器管理:新增执行器,AppName填xxl-job-executor-sample,注册方式选自动注册。稍等几秒,应该能看到两个在线实例的IP和端口。
- 任务管理:新增任务,配置项如下:
- 调度类型:CRON,比如0 0 2 * * ?(如果只想手动验证,可以选"无"然后手动执行一次)
- 运行模式:BEAN
- JobHandler:userTagRefreshJob
- 路由策略:分片广播
- 阻塞处理策略:单机串行
- 任务超时时间:0(不超时,后面调优会细说)
- 失败重试次数:1
- 在任务管理页面点击"执行一次",然后到调度日志里看结果。
如果一切正常,调度日志里会出现两条调度记录,分别对应两个执行器。每个执行器的日志里会打印:
分片任务启动, shardIndex=0, shardTotal=2 分片任务启动, shardIndex=1, shardTotal=2这就算跑通了。分片总数等于执行器数量,每个实例只处理自己负责的ID区间,整个任务并行完成。
4.4 本地验证时特别留意的一个细节
如果你只启动一个执行器,任务触发时shardTotal会是1,分片任务退化成单机全量任务,也能正常跑完。这其实是XXL-Job的一个容错特性:执行器少了多少分片,任务本身不会崩。但在生产环境,你要意识到"退化"不等于"没问题",节点数量直接影响横向扩展能力,监控上要盯执行器在线数量。
5. 分片广播实战中的常见问题与排查思路
5.1 数据重复处理:先怀疑路由策略,再查区间边界
我第一次把分片任务推到测试环境,对账发现同一条记录的update_time变了两次,立刻警觉起来。排查链路是这样的:
先看任务配置的路由策略是不是分片广播。如果配置的是轮询或一致性哈希,那每次触发只会有一个执行器跑全量任务,多台机器轮流跑,数据当然会重复处理。这是最典型的误配置。
再看阻塞处理策略。如果选了"覆盖之前调度",前一次任务还没跑完,下次调度会强制把前一次干掉再起新的,两个任务在时间上交错,业务上就可能出现重复或者部分覆盖。海量数据处理任务建议用"单机串行"。
最后查分片区间边界。我见过同事在计算rangeStart和rangeEnd时,两个相邻分片的区间重叠了,比如分片0的end是10000,分片1的start也是10000,那user_id=10000这条记录就被处理了两次。修复方式是统一用左闭右开区间:[start, end),下一个分片的start取上一个分片的end + 1。
5.2 数据倾斜:某个分片慢到拖垮整体
三个分片跑下来,两个10分钟完成,一个40分钟还在跑。这种倾斜问题在ID范围切分里很常见——如果切分键不是自增主键,或者业务数据在某个ID区间内密集堆积,必然会出现一个分片处理的数据量远大于其他分片。
排查时先把每个分片处理的记录数和耗时打出来对比。如果确实有倾斜,我有两个方向的解法:
第一,改切分策略。先从ID范围切分改成按数据量切分:先统计user_id在哪些区间密集,然后按"每段固定条数"来切边界。代价是每次任务开始前多跑一次count统计。
第二,保留ID范围切分,但把分片做细。比如只有3台执行器,但把ID范围切成9段,每个执行器分配3段。这样即使某一段特别多,也只会让某一个执行器多跑一段,而不会让整个任务的完成时间被一段极端数据拖死。
倾斜问题不是分片广播独有的,但分片广播让它的影响被放大了。日志里每个分片的处理耗时一定要记录下来,这是发现倾斜的第一手资料。
5.3 节点宕机导致数据漏处理
分片任务触发时,某台执行器恰好不在线,调度中心不会把任务分给它。这样一来,其他机器正常处理,但那一段ID区间没人管。更隐蔽的情况是任务跑到一半执行器宕机,admin显示调度失败,其它分片已经提交了,失败分片的数据就缺了一块。
我的兜底方案有两层。第一层是状态表幂等,任务重试时已成功的分片直接跳过,失败的分片重新处理。第二层是"整体校验",任务全部结束后校验batch_job_log里所有分片是否都是成功状态,并且各分片process_count之和是否等于理论上应处理的总行数。如果数量对不上,再针对具体区间补跑。
5.4 数据库连接池被打爆
分片任务并行度一上来,数据库连接池往往先扛不住。执行器的默认线程池是200,如果JobHandler内部还自己开了多线程,每个线程都从连接池拿连接,N个执行器同时打库,HikariCP默认的10个连接瞬间耗尽,任务开始一段时间后就会出现大量getConnection timeout。
解决办法是控制好并发度。JobHandler内部不要盲目开线程,先把并行度压在4到8之间;数据库连接池的maximum-pool-size调大到30到50;另外每个分片任务里的每批操作保持短事务,处理完一批立刻释放连接。如果你用Druid或者HikariCP,务必要在本地用两三个执行器同时压一下,实测连接池参数是否够用。
5.5 日志排查:按分片维度去追
分片任务出了bug,最怕的是所有实例日志混在一起,找不到谁是谁。我在代码里强制让每行日志带上shardIndex和当前处理的ID区间,排查的时候直接grep某个分片号,单独看它的处理链路。调度中心的"执行日志"功能可以查看每个执行器通过XxlJobHelper.log写入的日志,但业务日志还是要靠服务器上的文件,所以日志文件按执行器实例分开落盘也是必要的。