XXL-Job | 基于分片广播的海量数据处理实践
·
关注:CodingTechWork
引言
&esmp;在大数据量处理场景中,我们常常会遇到这样的挑战:单个执行节点无法在有限时间内完成庞大数据集的处理任务。比如需要导出数百万条用户数据、批量处理千万级日志文件等。传统单机执行模式不仅效率低下,还容易因超时或内存溢出导致任务失败。今天我们就来探讨如何利用XXL-Job的分片广播功能解决这一难题。
需求场景:海量用户数据导出
假设我们有一个实际业务需求:导出所有注册用户的详细信息为Excel文件。我们的用户表有2000万条记录,且每行用户数据包含多个关联字段(用户信息、订单统计、行为数据等)。
面临的挑战:
- 单次查询2000万条数据会导致数据库压力巨大,可能拖垮生产库
- 内存中处理如此大的数据集容易导致OOM(内存溢出)
- 导出过程耗时过长,单机执行可能超过任务超时时间
- 任务执行过程中断后需要重头开始,无法断点续传
解决方案设计:
我们将使用XXL-Job的分片广播功能,将2000万条数据划分为多个分片,由多个执行器实例并行处理,每个实例只处理分配给自己的数据片段。
XXL-Job分片广播原理详解
基本概念
分片广播是XXL-Job提供的一种分布式任务调度模式,在这种模式下:
- 调度中心的一次调度会广播触发集群中所有执行器
- 每个执行器收到任务请求后,会获取相同的分片参数
- 执行器根据分片参数(当前分片索引、总分片数)处理对应的数据子集
执行流程
好的,根据您描述的步骤,可以生成一个清晰的技术流程图。以下是使用 Mermaid 语法绘制的流程图,它直观地展示了 XXL-Job 分片广播的执行流程。
流程解读:
- 调度触发:任务由 XXL-Job 调度中心定时或手动触发。
- 广播通知:调度中心采用广播路由策略,将任务信号同时下发到所有注册在线的执行器(
Executor)。 - 并行执行:所有执行器接收并开始执行同一个任务定义(
JobHandler)。 - 获取上下文:在每个执行器的
JobHandler代码中,通过ShardingUtil工具获取本次执行的分片参数,即当前分片索引(index) 和总分片数(total)。 - 计算数据分片:业务逻辑根据
index和total,计算出当前执行器实例应该处理的数据子集(例如,将总数据ID范围等分,每个执行器处理其中一段)。 - 处理数据:执行器仅对自己负责的数据分片进行业务处理(如查询、计算、导出等),实现真正的分布式并行处理。
分片参数说明
执行器获取到的分片参数包含:
index:当前分片序号(从0开始)total:总分片数- 每个执行器的
index是唯一的,total是相同的
与普通分片的区别
| 特性 | 分片广播 | 普通分片 |
|---|---|---|
| 触发方式 | 所有执行器同时触发 | 仅触发一个执行器 |
| 参数传递 | 每个执行器获取完整分片参数 | 单个执行器处理所有分片 |
| 适用场景 | 真正并行处理 | 单执行器串行处理分片 |
XXL-Job配置方式
调度中心配置
-
添加执行器
- 登录XXL-Job调度中心
- 进入"执行器管理"页面
- 添加执行器,设置AppName(如:user-export-executor)
-
创建任务
任务描述:用户数据导出任务 JobHandler:userExportJobHandler 路由策略:分片广播 Cron表达式:0 0 2 * * ? # 每天凌晨2点执行 任务参数:可传递日期范围等自定义参数 阻塞处理策略:单机串行 任务超时时间:3600(单位:秒) 失败重试次数:3
执行器配置
pom.xml依赖:
<dependency>
<groupId>com.xuxueli</groupId>
<artifactId>xxl-job-core</artifactId>
<version>2.3.1</version>
</dependency>
application.yml配置:
xxl:
job:
admin:
addresses: http://xxl-job-admin:8080/xxl-job-admin
executor:
appname: user-export-executor
ip:
port: 9999
logpath: /data/applogs/xxl-job/jobhandler
logretentiondays: 30
accessToken:
执行器配置类:
@Configuration
public class XxlJobConfig {
@Value("${xxl.job.admin.addresses}")
private String adminAddresses;
@Value("${xxl.job.executor.appname}")
private String appName;
@Value("${xxl.job.executor.ip}")
private String ip;
@Value("${xxl.job.executor.port}")
private int port;
@Value("${xxl.job.accessToken}")
private String accessToken;
@Value("${xxl.job.executor.logpath}")
private String logPath;
@Value("${xxl.job.executor.logretentiondays}")
private int logRetentionDays;
@Bean
public XxlJobSpringExecutor xxlJobExecutor() {
XxlJobSpringExecutor xxlJobSpringExecutor = new XxlJobSpringExecutor();
xxlJobSpringExecutor.setAdminAddresses(adminAddresses);
xxlJobSpringExecutor.setAppname(appName);
xxlJobSpringExecutor.setIp(ip);
xxlJobSpringExecutor.setPort(port);
xxlJobSpringExecutor.setAccessToken(accessToken);
xxlJobSpringExecutor.setLogPath(logPath);
xxlJobSpringExecutor.setLogRetentionDays(logRetentionDays);
return xxlJobSpringExecutor;
}
}
Java实现逻辑
分片策略设计
对于2000万用户数据,我们设计以下分片策略:
- 按用户ID范围进行分片
- 每个分片处理固定数量的用户(如:每个分片50万条)
- 总分片数 = ceil(总记录数 / 每分片记录数)
JobHandler实现
@Component
public class UserExportJobHandler {
private static final Logger logger = LoggerFactory.getLogger(UserExportJobHandler.class);
@XxlJob("userExportJobHandler")
public void execute() throws Exception {
// 1. 获取分片参数
ShardingUtil.ShardingVO shardingVO = ShardingUtil.getShardingVo();
// 当前分片索引
int shardIndex = shardingVO.getIndex();
// 总分片数
int shardTotal = shardingVO.getTotal();
logger.info("分片参数:当前分片索引={}, 总分片数={}", shardIndex, shardTotal);
// 2. 获取任务参数(可从调度中心传递)
String param = XxlJobHelper.getJobParam();
DateRange dateRange = parseParam(param);
// 3. 计算当前分片的数据范围
long totalCount = getUserTotalCount(dateRange);
long pageSize = calculatePageSize(totalCount, shardTotal);
long startId = calculateStartId(shardIndex, pageSize);
long endId = calculateEndId(shardIndex, pageSize, totalCount);
logger.info("数据范围:总记录数={}, 当前分片处理ID范围={}~{}",
totalCount, startId, endId);
// 4. 分页查询并处理数据
processUserDataByRange(dateRange, startId, endId, shardIndex);
// 5. 任务执行结果
XxlJobHelper.handleSuccess("用户导出任务执行成功,分片索引:" + shardIndex);
}
/**
* 计算每个分片的数据量
*/
private long calculatePageSize(long totalCount, int shardTotal) {
long pageSize = totalCount / shardTotal;
if (totalCount % shardTotal != 0) {
// 向上取整
pageSize += 1;
}
return pageSize;
}
/**
* 计算起始ID
*/
private long calculateStartId(int shardIndex, long pageSize) {
return shardIndex * pageSize;
}
/**
* 计算结束ID
*/
private long calculateEndId(int shardIndex, long pageSize, long totalCount) {
long endId = (shardIndex + 1) * pageSize - 1;
return Math.min(endId, totalCount - 1);
}
/**
* 按范围处理用户数据
*/
private void processUserDataByRange(DateRange dateRange, long startId,
long endId, int shardIndex) {
// 每次查询1000条
int pageSize = 1000;
long currentId = startId;
while (currentId <= endId) {
// 分页查询用户数据
List<UserDTO> userList = userDao.findUsersByIdRange(
dateRange, currentId, Math.min(currentId + pageSize - 1, endId)
);
if (CollectionUtils.isEmpty(userList)) {
break;
}
// 处理用户数据(如:生成Excel行)
processUserBatch(userList, shardIndex);
// 更新进度
currentId += pageSize;
updateProgress(startId, endId, currentId, shardIndex);
}
// 合并生成最终文件
mergeShardFiles(shardIndex);
}
/**
* 处理一批用户数据
*/
private void processUserBatch(List<UserDTO> userList, int shardIndex) {
List<UserExportVO> exportData = new ArrayList<>();
for (UserDTO user : userList) {
UserExportVO exportVO = convertToExportVO(user);
// 补充关联数据
enrichUserData(exportVO);
exportData.add(exportVO);
}
// 写入分片临时文件
writeToShardFile(exportData, shardIndex);
logger.info("已处理{}条用户数据,分片索引:{}", userList.size(), shardIndex);
}
/**
* 更新处理进度
*/
private void updateProgress(long startId, long endId,
long currentId, int shardIndex) {
long total = endId - startId + 1;
long processed = currentId - startId;
double progress = (double) processed / total * 100;
logger.info("分片{}处理进度:{:.2f}% ({}/{})",
shardIndex, progress, processed, total);
// 可记录到数据库供监控查看
saveProgressToDB(shardIndex, progress);
}
/**
* 合并分片文件
*/
private void mergeShardFiles(int shardIndex) {
// 如果是最后一个分片,负责合并所有分片文件
if (isLastShard(shardIndex)) {
logger.info("最后分片{}开始合并文件", shardIndex);
// 等待其他分片完成
waitForOtherShards();
// 合并所有临时文件
List<File> shardFiles = collectAllShardFiles();
File finalFile = mergeToFinalFile(shardFiles);
// 上传到文件服务器或OSS
uploadToStorage(finalFile);
// 清理临时文件
cleanTempFiles(shardFiles);
logger.info("文件合并完成,最终文件:{}", finalFile.getName());
}
}
}
辅助服务类
@Mapper
public interface UserExportMapper {
/**
* 获取用户总数
*/
@Select("""
SELECT COUNT(*)
FROM user
WHERE create_time BETWEEN #{startTime} AND #{endTime}
""")
Long selectUserTotalCount(@Param("startTime") Date startTime,
@Param("endTime") Date endTime);
/**
* 根据ID范围查询用户详细信息
*/
@Select("""
SELECT
u.id,
u.username,
u.email,
u.create_time as createTime,
COALESCE(o.order_count, 0) AS orderCount,
COALESCE(o.total_amount, 0) AS totalAmount
FROM user u
LEFT JOIN (
SELECT
user_id,
COUNT(*) AS order_count,
SUM(amount) AS total_amount
FROM orders
WHERE create_time BETWEEN #{startTime} AND #{endTime}
GROUP BY user_id
) o ON u.id = o.user_id
WHERE u.id BETWEEN #{startId} AND #{endId}
AND u.create_time BETWEEN #{startTime} AND #{endTime}
ORDER BY u.id
LIMIT #{limit}
""")
List<UserDTO> selectUsersByIdRange(@Param("startTime") Date startTime,
@Param("endTime") Date endTime,
@Param("startId") Long startId,
@Param("endId") Long endId,
@Param("limit") Long limit);
}
运行结果分析
执行日志示例
2024-01-15 02:00:00 [UserExportJobHandler] INFO - 分片参数:当前分片索引=0,总分片数=4
2024-01-15 02:00:00 [UserExportJobHandler] INFO - 数据范围:总记录数=20000000,当前分片处理ID范围=0~4999999
2024-01-15 02:00:05 [UserExportJobHandler] INFO - 已处理1000条用户数据,分片索引:0
2024-01-15 02:00:10 [UserExportJobHandler] INFO - 分片0处理进度:0.10% (5000/5000000)
2024-01-15 02:30:45 [UserExportJobHandler] INFO - 分片0处理进度:99.50% (4975000/5000000)
2024-01-15 02:32:10 [UserExportJobHandler] INFO - 最后分片3开始合并文件
2024-01-15 02:35:20 [UserExportJobHandler] INFO - 文件合并完成,最终文件:user_export_20240115.xlsx
调度中心监控视图
在XXL-Job调度中心,我们可以看到:
- 4个执行器实例同时在线
- 每个实例的任务执行日志独立
- 任务整体执行时间大幅缩短
- 可以实时查看每个分片的执行进度
更多推荐



所有评论(0)