news 2026/9/9 20:37:01

XXL-Job分片广播实战:亿级用户标签数据并行刷新方案

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
XXL-Job分片广播实战:亿级用户标签数据并行刷新方案

凌晨两点被电话叫醒是什么体验,我在接手那个用户标签刷新任务之后,连续体会了一个星期。任务本身不复杂:每天全量刷新用户标签表,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 -DskipTests

4.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写入的日志,但业务日志还是要靠服务器上的文件,所以日志文件按执行器实例分开落盘也是必要的。

6. 从"能跑"到"跑得好":分片任务的调优经验

6.1 分片数怎么定:

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/9 20:36:11

数据安全治理自动化框架:从数据测绘到响应闭环的落地指南

我上个月帮一家企业做数据安全治理现状摸底&#xff0c;拿到数据资产清单的时候愣了一下——Excel里堆了四千多张表&#xff0c;大部分没人说得清里面存的是什么数据、谁在访问、有没有出过库。这几乎是所有数据安全治理项目的常态&#xff1a;不是缺制度&#xff0c;不是缺工具…

作者头像 李华
网站建设 2026/9/9 20:36:02

用Matlab实现Ghil-Sellers能量平衡模型:从双稳态到气候突变模拟

我最近几个月一直在折腾一个看起来有点“复古”的模型——Ghil-Sellers能量平衡模型。说它复古&#xff0c;是因为这模型比现在动辄几十万行代码的大气环流模式&#xff08;GCM&#xff09;老了半个世纪&#xff0c;可它模拟出来的结果&#xff0c;却让人惊讶地“现代”&#x…

作者头像 李华
网站建设 2026/9/9 20:35:19

浏览器开发者工具 Network 面板实战:从请求分析到性能优化

我接触浏览器开发者工具这么多年&#xff0c;如果只让我选一个面板作为日常主力&#xff0c;我会毫不犹豫选 Network 面板。它就像是前端的侦察兵&#xff0c;所有发生在浏览器和服务器之间的数据往来&#xff0c;在这个面板里都无所遁形。不管是接口报错、页面加载慢、资源加载…

作者头像 李华
网站建设 2026/9/9 20:35:18

astcenc源码审计:移动端ASTC纹理压缩与GPU带宽优化

移动端项目做性能专项时&#xff0c;我几乎每次都会在RenderDoc里盯着纹理带宽那几栏待上半天。采样的纹理数量少说几十张&#xff0c;每张如果还是RGBA32直出&#xff0c;带宽压力说出来都是泪。ASTC&#xff08;Adaptive Scalable Texture Compression&#xff09;在这时候几…

作者头像 李华
网站建设 2026/9/9 20:33:51

一台Android手机同时运行多个Root环境:六套方案与实战

玩机圈里有个始终绕不开的问题&#xff1a;一台 Android 手机&#xff0c;能不能在同一个硬件上同时拥有多个互相独立的 Root 环境&#xff1f;不同场景需要不同的 Root 策略&#xff0c;或者想在一台备用机上同时运行“日用系统 Linux 容器 多系统测试环境”&#xff0c;却总…

作者头像 李华