ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

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) { ListUserTag batch userTagMapper.selectUserTagByRange(cursor, cursor batchSize); if (batch.isEmpty()) { break; } // 根据订单、登录、客服记录重新计算标签 ListUserTag 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的范围条件天然支持断点续跑。比如某个分片处理到一半进程挂了重试时会从状态表看到status0再从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.port8080 spring.datasource.urljdbc:mysql://127.0.0.1:3306/xxl_job?useSSLfalseserverTimezoneAsia/Shanghai spring.datasource.usernameroot spring.datasource.password123456 xxl.job.accessTokendefault_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.propertiesxxl.job.admin.addresseshttp://localhost:8080/xxl-job-admin xxl.job.accessTokendefault_token xxl.job.executor.appnamexxl-job-executor-sample xxl.job.executor.port9999 xxl.job.executor.logpathlogs/xxl-job/jobhandler启动第一个执行器实例。然后再起一个实例注意把执行器端口改掉否则两个进程端口冲突# 第二个实例指定不同端口 java -jar xxl-job-executor-sample-springboot-2.4.0.jar --xxl.job.executor.port9998两个执行器都启动后去admin后台完成三件事执行器管理新增执行器AppName填xxl-job-executor-sample注册方式选自动注册。稍等几秒应该能看到两个在线实例的IP和端口。任务管理新增任务配置项如下调度类型CRON比如0 0 2 * * ?如果只想手动验证可以选无然后手动执行一次运行模式BEANJobHandleruserTagRefreshJob路由策略分片广播阻塞处理策略单机串行任务超时时间0不超时后面调优会细说失败重试次数1在任务管理页面点击执行一次然后到调度日志里看结果。如果一切正常调度日志里会出现两条调度记录分别对应两个执行器。每个执行器的日志里会打印分片任务启动, shardIndex0, shardTotal2 分片任务启动, shardIndex1, shardTotal2这就算跑通了。分片总数等于执行器数量每个实例只处理自己负责的ID区间整个任务并行完成。4.4 本地验证时特别留意的一个细节如果你只启动一个执行器任务触发时shardTotal会是1分片任务退化成单机全量任务也能正常跑完。这其实是XXL-Job的一个容错特性执行器少了多少分片任务本身不会崩。但在生产环境你要意识到退化不等于没问题节点数量直接影响横向扩展能力监控上要盯执行器在线数量。5. 分片广播实战中的常见问题与排查思路5.1 数据重复处理先怀疑路由策略再查区间边界我第一次把分片任务推到测试环境对账发现同一条记录的update_time变了两次立刻警觉起来。排查链路是这样的先看任务配置的路由策略是不是分片广播。如果配置的是轮询或一致性哈希那每次触发只会有一个执行器跑全量任务多台机器轮流跑数据当然会重复处理。这是最典型的误配置。再看阻塞处理策略。如果选了覆盖之前调度前一次任务还没跑完下次调度会强制把前一次干掉再起新的两个任务在时间上交错业务上就可能出现重复或者部分覆盖。海量数据处理任务建议用单机串行。最后查分片区间边界。我见过同事在计算rangeStart和rangeEnd时两个相邻分片的区间重叠了比如分片0的end是10000分片1的start也是10000那user_id10000这条记录就被处理了两次。修复方式是统一用左闭右开区间[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 分片数怎么定
返回列表