关注:CodingTechWork

引言

 &esmp;在大数据量处理场景中,我们常常会遇到这样的挑战:单个执行节点无法在有限时间内完成庞大数据集的处理任务。比如需要导出数百万条用户数据、批量处理千万级日志文件等。传统单机执行模式不仅效率低下,还容易因超时或内存溢出导致任务失败。今天我们就来探讨如何利用XXL-Job的分片广播功能解决这一难题。

需求场景:海量用户数据导出

假设我们有一个实际业务需求:导出所有注册用户的详细信息为Excel文件。我们的用户表有2000万条记录,且每行用户数据包含多个关联字段(用户信息、订单统计、行为数据等)。

面临的挑战:

  • 单次查询2000万条数据会导致数据库压力巨大,可能拖垮生产库
  • 内存中处理如此大的数据集容易导致OOM(内存溢出)
  • 导出过程耗时过长,单机执行可能超过任务超时时间
  • 任务执行过程中断后需要重头开始,无法断点续传

解决方案设计:
我们将使用XXL-Job的分片广播功能,将2000万条数据划分为多个分片,由多个执行器实例并行处理,每个实例只处理分配给自己的数据片段。

XXL-Job分片广播原理详解

基本概念

分片广播是XXL-Job提供的一种分布式任务调度模式,在这种模式下:

  • 调度中心的一次调度会广播触发集群中所有执行器
  • 每个执行器收到任务请求后,会获取相同的分片参数
  • 执行器根据分片参数(当前分片索引、总分片数)处理对应的数据子集

执行流程

好的,根据您描述的步骤,可以生成一个清晰的技术流程图。以下是使用 Mermaid 语法绘制的流程图,它直观地展示了 XXL-Job 分片广播的执行流程。

“调度中心触发任务”

“广播通知所有在线执行器”

“每个执行器并行执行”

“执行相同的 JobHandler”

“在 Handler 内获取分片参数
(index, total)”

“根据分片参数计算
本节点需处理的数据范围”

“处理指定范围的数据”

“任务完成”

流程解读:

  1. 调度触发:任务由 XXL-Job 调度中心定时或手动触发。
  2. 广播通知:调度中心采用广播路由策略,将任务信号同时下发到所有注册在线的执行器(Executor)。
  3. 并行执行:所有执行器接收并开始执行同一个任务定义(JobHandler)。
  4. 获取上下文:在每个执行器的 JobHandler 代码中,通过 ShardingUtil 工具获取本次执行的分片参数,即当前分片索引(index)总分片数(total)
  5. 计算数据分片:业务逻辑根据 indextotal,计算出当前执行器实例应该处理的数据子集(例如,将总数据ID范围等分,每个执行器处理其中一段)。
  6. 处理数据:执行器仅对自己负责的数据分片进行业务处理(如查询、计算、导出等),实现真正的分布式并行处理

分片参数说明

执行器获取到的分片参数包含:

  • index:当前分片序号(从0开始)
  • total:总分片数
  • 每个执行器的index是唯一的,total是相同的

与普通分片的区别

特性分片广播普通分片
触发方式所有执行器同时触发仅触发一个执行器
参数传递每个执行器获取完整分片参数单个执行器处理所有分片
适用场景真正并行处理单执行器串行处理分片

XXL-Job配置方式

调度中心配置

  1. 添加执行器

    • 登录XXL-Job调度中心
    • 进入"执行器管理"页面
    • 添加执行器,设置AppName(如:user-export-executor)
  2. 创建任务

    任务描述:用户数据导出任务
    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个执行器实例同时在线
  • 每个实例的任务执行日志独立
  • 任务整体执行时间大幅缩短
  • 可以实时查看每个分片的执行进度
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐