ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Spring Boot实战:构建高并发任务分组调度系统

Spring Boot实战:构建高并发任务分组调度系统 最近在开发一个需要处理复杂业务逻辑和团队协作的项目时遇到了一个典型问题如何将一个大任务高效、清晰地拆分成多个子任务并分配给不同的执行单元或团队并行处理同时还要能直观地展示进度和结果。这让我想起了很多综艺节目里“分组对抗”的赛制设计其核心思想与软件开发中的任务分解、模块化、并行处理有异曲同工之妙。本文将以一个模拟的“团队竞技任务处理系统”为例完整拆解从需求分析、架构设计、核心代码实现到部署演示的全过程。我们将借鉴“分组对抗”的直观逻辑构建一个后端服务实现任务的动态分组、独立执行、进度追踪和结果汇总。无论你是想学习Spring Boot项目实战、理解多线程与任务调度还是需要为你的系统设计一个灵活的作业分发框架这篇文章都能提供一套可直接复用的代码和设计思路。1. 背景与核心概念任务分组执行的工程价值在软件工程中尤其是后台服务、数据处理平台或自动化测试框架里我们常常面临批量作业的处理需求。例如数据处理需要处理来自不同数据源的百万条记录。压力测试需要模拟成千上万个并发用户访问。报表生成需要为数十个部门生成每日运营报表。如果将所有任务塞进一个巨大的循环里串行执行效率低下且一个任务的失败可能导致整个流程中断。更优的解决方案是“分而治之”任务拆分将大任务拆解为多个独立或弱关联的子任务。分组/分片根据某种策略如资源类型、数据特征、优先级将子任务分组。并行执行每个分组由一个独立的执行单元线程、进程、服务实例处理。结果聚合收集各分组的执行结果汇总成最终输出。这恰好对应了“分组对抗”的模型总任务一公表演 - 分组黑马队、白马队 - 队员子任务 - 表演执行 - 评分结果汇总。我们的项目将模拟这一过程核心概念包括主任务Mission最高层级的任务单元包含总体描述、所有子任务以及最终状态。任务组Team子任务的集合拥有独立的执行上下文和进度跟踪。例如“黑马队”和“白马队”。子任务Task最小的可执行单元如“队员A完成舞蹈部分”。执行器Executor负责实际执行子任务的组件通常映射到一个线程或异步任务。调度器Scheduler负责任务的拆分、分组指派以及执行进度的监控。通过构建这样一个系统我们可以深入理解并发编程、任务调度、状态管理以及如何设计清晰的服务边界。2. 环境准备与版本说明本项目是一个标准的Spring Boot后端应用使用Maven进行构建。以下是开发环境建议操作系统Windows 10/11, macOS, 或 Linux (Ubuntu 20.04)。本文演示环境为macOS。Java SDKJDK 11或JDK 17(LTS版本)。推荐使用JDK 17以获得更好的性能和新特性。本文使用OpenJDK 17.0.10。构建工具Apache Maven 3.6。请确保mvn -v命令能正确输出版本信息。集成开发环境IDEIntelliJ IDEA (推荐), Eclipse 或 VS Code。本文使用 IntelliJ IDEA 2023.3。项目管理本文会提供完整的pom.xml和项目结构。版本兼容性说明Spring Boot 2.7.x 是一个长期支持版本与JDK 17兼容性好社区资源丰富。我们选择此版本进行演示。实际项目中请根据公司技术栈和稳定性要求选择合适的版本。3. 核心原理与项目架构设计在动手写代码之前我们先设计系统的核心类和它们之间的关系。这有助于理解后续的代码实现。3.1 领域模型设计我们定义以下几个核心领域对象Mission主任务属性ID、名称、描述、状态待开始、进行中、已完成、失败、创建时间、结束时间。关联包含多个Team。Team任务组属性ID、名称如“黑马队”、所属Mission ID、状态、进度百分比。关联包含多个Task。Task子任务属性ID、描述、预计耗时秒、实际状态待分配、执行中、成功、失败、执行结果消息、开始时间、结束时间。关联属于一个Team。TaskExecutor任务执行器接口定义execute(Task task)方法。这是策略模式的应用允许我们灵活替换不同的任务执行逻辑。3.2 系统架构与流程系统采用典型的应用分层结构Controller层提供RESTful API用于创建任务、查询进度等。Service层核心业务逻辑包括任务拆分、分组、调度执行。Repository层数据持久化这里为了简化使用并发安全的ConcurrentHashMap模拟实际项目可替换为MySQL、Redis等。Executor层具体任务的执行实现模拟耗时操作。核心执行流程用户通过API创建一个新的Mission并指定总任务描述。MissionService根据规则例如简单平均分配将Mission拆分成若干Task并将这些Task分配到两个Team黑马队、白马队中。TeamScheduler启动为每个Team创建一个独立的线程池或使用CompletableFuture并行执行该队内的所有Task。每个Task由对应的TaskExecutor执行并更新自己的状态和结果。Team监控其下所有Task的状态计算整体进度。Mission监控所有Team的状态当所有Team完成时标记自身为完成。用户可以通过API实时查询Mission和Team的进度。4. 完整实战构建团队竞技任务处理系统接下来我们一步步实现这个系统。4.1 创建项目结构与依赖首先使用 Spring Initializr 或IDE创建Spring Boot项目。pom.xml 关键依赖?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version !-- 使用稳定的2.7.x版本 -- relativePath/ /parent groupIdcom.example/groupId artifactIdteam-task-system/artifactId version0.0.1-SNAPSHOT/version nameteam-task-system/name descriptionDemo project for team-based task processing/description properties java.version17/java.version /properties dependencies !-- Web支持提供REST API -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- 简化JSON处理 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-json/artifactId /dependency !-- 参数校验 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-validation/artifactId /dependency !-- 测试 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency !-- Lombok 简化Getter/Setter等代码 -- dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency /dependencies build plugins plugin groupIdorg.springframework.boot/groupId artifactIdspring-boot-maven-plugin/artifactId configuration excludes exclude groupIdorg.projectlombok/groupId artifactIdlombok/artifactId /exclude /excludes /configuration /plugin /plugins /build /project4.2 定义领域模型与枚举创建model包并定义以下类。TaskStatus.java (任务状态枚举)package com.example.teamtasksystem.model; public enum TaskStatus { PENDING, // 待分配/待执行 PROCESSING, // 执行中 SUCCESS, // 成功 FAILED // 失败 }Mission.java (主任务实体)package com.example.teamtasksystem.model; import lombok.Data; import java.time.LocalDateTime; import java.util.List; Data public class Mission { private String id; private String name; private String description; private TaskStatus status; // 整体任务状态 private LocalDateTime createTime; private LocalDateTime finishTime; private ListTeam teams; // 包含的队伍 public Mission() { this.createTime LocalDateTime.now(); this.status TaskStatus.PENDING; } }Team.java (任务组实体)package com.example.teamtasksystem.model; import lombok.Data; import java.util.List; import java.util.concurrent.atomic.AtomicInteger; Data public class Team { private String id; private String name; // 如 “黑马队” “白马队” private String missionId; private TaskStatus status; private ListTask tasks; // 进度 (0-100)使用AtomicInteger保证线程安全 private AtomicInteger progress new AtomicInteger(0); public Team() { this.status TaskStatus.PENDING; } // 计算并更新当前进度 public void calculateProgress() { if (tasks null || tasks.isEmpty()) { progress.set(0); return; } long completedCount tasks.stream() .filter(t - t.getStatus() TaskStatus.SUCCESS || t.getStatus() TaskStatus.FAILED) .count(); int newProgress (int) ((completedCount * 100) / tasks.size()); progress.set(newProgress); // 根据进度和任务状态更新队伍状态 boolean allDone tasks.stream().allMatch(t - t.getStatus() TaskStatus.SUCCESS || t.getStatus() TaskStatus.FAILED); boolean anyProcessing tasks.stream().anyMatch(t - t.getStatus() TaskStatus.PROCESSING); if (allDone) { this.status TaskStatus.SUCCESS; // 简化逻辑全部完成即成功 } else if (anyProcessing) { this.status TaskStatus.PROCESSING; } else { this.status TaskStatus.PENDING; } } }Task.java (子任务实体)package com.example.teamtasksystem.model; import lombok.Data; import java.time.LocalDateTime; Data public class Task { private String id; private String description; private Long estimatedDurationSeconds; // 预计耗时 private TaskStatus status; private String resultMessage; // 执行结果信息 private LocalDateTime startTime; private LocalDateTime endTime; private String teamId; // 所属队伍ID public Task() { this.status TaskStatus.PENDING; } }4.3 实现任务执行器与调度服务创建service包和executor包。SimulatedTaskExecutor.java (模拟任务执行器)package com.example.teamtasksystem.executor; import com.example.teamtasksystem.model.Task; import com.example.teamtasksystem.model.TaskStatus; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import java.time.LocalDateTime; import java.util.concurrent.ThreadLocalRandom; Slf4j Component public class SimulatedTaskExecutor { /** * 模拟执行一个任务 * param task 待执行的任务 */ public void execute(Task task) { log.info(队伍[{}]的任务[{}]开始执行: {}, task.getTeamId(), task.getId(), task.getDescription()); task.setStatus(TaskStatus.PROCESSING); task.setStartTime(LocalDateTime.now()); try { // 模拟任务执行耗时范围在预计时间的80%-120%之间波动 long baseDuration task.getEstimatedDurationSeconds(); long actualDuration (long) (baseDuration * (0.8 ThreadLocalRandom.current().nextDouble() * 0.4)); Thread.sleep(actualDuration * 1000); // 转换为毫秒 // 模拟小概率失败 if (ThreadLocalRandom.current().nextDouble() 0.1) { // 10%失败率 throw new RuntimeException(模拟任务执行过程中发生随机错误); } // 执行成功 task.setStatus(TaskStatus.SUCCESS); task.setResultMessage(String.format(任务成功完成实际耗时 %d 秒, actualDuration)); log.info(任务[{}]执行成功耗时{}秒, task.getId(), actualDuration); } catch (InterruptedException e) { Thread.currentThread().interrupt(); task.setStatus(TaskStatus.FAILED); task.setResultMessage(任务被中断: e.getMessage()); log.error(任务[{}]被中断, task.getId(), e); } catch (Exception e) { task.setStatus(TaskStatus.FAILED); task.setResultMessage(执行失败: e.getMessage()); log.error(任务[{}]执行失败, task.getId(), e); } finally { task.setEndTime(LocalDateTime.now()); } } }MissionService.java (主任务服务)package com.example.teamtasksystem.service; import com.example.teamtasksystem.model.*; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import javax.annotation.PostConstruct; import java.util.*; import java.util.concurrent.*; import java.util.stream.Collectors; import java.util.stream.IntStream; Service Slf4j public class MissionService { // 内存存储模拟数据库 private final MapString, Mission missionStore new ConcurrentHashMap(); private final MapString, Team teamStore new ConcurrentHashMap(); Autowired private SimulatedTaskExecutor taskExecutor; // 为每个队伍分配独立的固定大小线程池实现队伍间资源隔离 private final MapString, ExecutorService teamExecutors new ConcurrentHashMap(); PostConstruct public void init() { log.info(MissionService 初始化完成); } /** * 创建并启动一个新的主任务 * param missionName 任务名 * param totalTasks 需要创建的子任务总数 * return 创建的主任务ID */ public String createAndStartMission(String missionName, int totalTasks) { String missionId MISSION- System.currentTimeMillis(); Mission mission new Mission(); mission.setId(missionId); mission.setName(missionName); mission.setDescription(模拟团队竞技任务处理); mission.setStatus(TaskStatus.PROCESSING); // 1. 创建两个队伍 Team blackTeam createTeam(黑马队, missionId); Team whiteTeam createTeam(白马队, missionId); // 2. 拆分任务到两个队伍 (简单奇偶分配) ListTask allTasks generateTasks(totalTasks, missionId); for (int i 0; i allTasks.size(); i) { Task task allTasks.get(i); if (i % 2 0) { task.setTeamId(blackTeam.getId()); blackTeam.getTasks().add(task); } else { task.setTeamId(whiteTeam.getId()); whiteTeam.getTasks().add(task); } } // 3. 保存队伍和任务 teamStore.put(blackTeam.getId(), blackTeam); teamStore.put(whiteTeam.getId(), whiteTeam); mission.setTeams(Arrays.asList(blackTeam, whiteTeam)); missionStore.put(missionId, mission); log.info(主任务[{}]创建成功包含{}个子任务分配给[{}]和[{}], missionId, totalTasks, blackTeam.getName(), whiteTeam.getName()); // 4. 异步启动两个队伍的任务执行 startTeamExecution(blackTeam); startTeamExecution(whiteTeam); return missionId; } private Team createTeam(String teamName, String missionId) { Team team new Team(); team.setId(TEAM- teamName - System.currentTimeMillis()); team.setName(teamName); team.setMissionId(missionId); team.setTasks(new CopyOnWriteArrayList()); // 线程安全的List return team; } private ListTask generateTasks(int count, String missionId) { return IntStream.rangeClosed(1, count) .mapToObj(i - { Task task new Task(); task.setId(TASK- missionId - i); task.setDescription(模拟子任务 # i); // 随机预计耗时 3-10 秒 task.setEstimatedDurationSeconds(ThreadLocalRandom.current().nextLong(3, 11)); return task; }).collect(Collectors.toList()); } /** * 启动一个队伍内所有任务的执行 * 每个队伍使用独立的线程池 */ private void startTeamExecution(Team team) { String teamId team.getId(); // 为队伍创建线程池核心线程数等于任务数最多不超过10个 int poolSize Math.min(team.getTasks().size(), 10); ExecutorService executor Executors.newFixedThreadPool(poolSize, r - new Thread(r, Executor- team.getName() - teamId.substring(teamId.length() - 4))); teamExecutors.put(teamId, executor); log.info(队伍[{}]开始执行线程池大小: {}, team.getName(), poolSize); // 提交所有任务到线程池 ListFuture? futures team.getTasks().stream() .map(task - executor.submit(() - { taskExecutor.execute(task); // 任务执行完后更新队伍进度 team.calculateProgress(); updateMissionStatus(team.getMissionId()); })) .collect(Collectors.toList()); // 添加一个关闭钩子当所有任务完成后关闭线程池实际项目中应有更优雅的管理 CompletableFuture.runAsync(() - { for (Future? future : futures) { try { future.get(); // 等待所有任务完成 } catch (InterruptedException | ExecutionException e) { log.error(等待任务完成时出错, e); } } executor.shutdown(); log.info(队伍[{}]所有任务执行完毕线程池已关闭, team.getName()); }); } /** * 更新主任务状态 */ private void updateMissionStatus(String missionId) { Mission mission missionStore.get(missionId); if (mission null) return; ListTeam teams mission.getTeams(); boolean allTeamsDone teams.stream() .allMatch(t - t.getStatus() TaskStatus.SUCCESS || t.getStatus() TaskStatus.FAILED); if (allTeamsDone) { mission.setStatus(TaskStatus.SUCCESS); // 简化所有队伍完成即任务成功 mission.setFinishTime(LocalDateTime.now()); log.info(主任务[{}]所有队伍执行完毕任务完成, missionId); } } /** * 根据ID获取主任务详情 */ public Mission getMission(String missionId) { Mission mission missionStore.get(missionId); if (mission ! null) { // 实时计算并更新队伍进度 mission.getTeams().forEach(Team::calculateProgress); updateMissionStatus(missionId); } return mission; } /** * 获取所有主任务 */ public ListMission getAllMissions() { return new ArrayList(missionStore.values()); } }4.4 提供REST API控制器创建controller包。MissionController.javapackage com.example.teamtasksystem.controller; import com.example.teamtasksystem.model.Mission; import com.example.teamtasksystem.service.MissionService; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.*; import java.util.List; RestController RequestMapping(/api/missions) public class MissionController { Autowired private MissionService missionService; /** * 创建并启动一个新的团队任务 * param missionName 任务名称 * param taskCount 子任务数量 * return 主任务ID */ PostMapping(/start) public String startNewMission(RequestParam String missionName, RequestParam(defaultValue 20) int taskCount) { if (taskCount 0) { throw new IllegalArgumentException(任务数量必须大于0); } return missionService.createAndStartMission(missionName, taskCount); } /** * 根据ID查询任务详情 */ GetMapping(/{missionId}) public Mission getMission(PathVariable String missionId) { Mission mission missionService.getMission(missionId); if (mission null) { throw new RuntimeException(未找到ID为 missionId 的任务); } return mission; } /** * 获取所有任务列表 */ GetMapping public ListMission getAllMissions() { return missionService.getAllMissions(); } }4.5 运行与验证1. 启动应用找到主启动类TeamTaskSystemApplication(通常由Spring Initializr生成)运行它。# 或者在IDE中直接运行 # 控制台应看到Spring Boot启动日志2. 使用API创建任务使用curl、Postman 或浏览器访问API。创建任务(POST请求):curl -X POST http://localhost:8080/api/missions/start?missionName越披哥一公taskCount10响应示例:MISSION-1712345678901返回的是新创建的主任务ID。同时控制台会打印日志显示任务已创建两个队伍开始执行。查询任务进度(GET请求):curl http://localhost:8080/api/missions/MISSION-1712345678901响应示例 (JSON):{ id: MISSION-1712345678901, name: 越披哥一公, description: 模拟团队竞技任务处理, status: PROCESSING, createTime: 2024-04-06T10:00:00, finishTime: null, teams: [ { id: TEAM-黑马队-1712345678902, name: 黑马队, missionId: MISSION-1712345678901, status: PROCESSING, tasks: [ { id: TASK-MISSION-1712345678901-1, description: 模拟子任务 #1, estimatedDurationSeconds: 7, status: SUCCESS, resultMessage: 任务成功完成实际耗时 6 秒, startTime: 2024-04-06T10:00:01, endTime: 2024-04-06T10:00:07, teamId: TEAM-黑马队-1712345678902 }, // ... 更多任务 ], progress: 40 }, { id: TEAM-白马队-1712345678903, name: 白马队, missionId: MISSION-1712345678901, status: PROCESSING, tasks: [...], progress: 60 } ] }你可以多次调用此接口观察progress字段的变化以及各个Task的status变化直到所有任务完成主任务status变为SUCCESS。查看所有任务:curl http://localhost:8080/api/missions3. 观察控制台日志应用启动后控制台会实时打印各个任务的开始、成功或失败日志类似于... 队伍[黑马队]的任务[TASK-MISSION-...-1]开始执行: 模拟子任务 #1 ... 任务[TASK-MISSION-...-1]执行成功耗时6秒 ... 队伍[白马队]的任务[TASK-MISSION-...-2]开始执行: 模拟子任务 #2 ... 任务[TASK-MISSION-...-2]执行失败 ... ... 主任务[MISSION-...]所有队伍执行完毕任务完成4.6 结果说明通过这个简单的系统我们成功模拟了任务动态拆分与分组将10个子任务按奇偶索引自动分配给了“黑马队”和“白马队”。并行与隔离执行每个队伍使用独立的线程池执行任务队伍间互不阻塞。进度实时监控通过progress字段可以实时查看每个队伍的完成百分比。状态全局管理主任务的状态依赖于所有队伍的状态。容错处理任务执行模拟了10%的随机失败率系统能正常记录失败状态而不影响其他任务。这为构建更复杂的分布式任务调度系统如使用Spring Batch、Quartz或XXL-JOB打下了坚实的基础。5. 常见问题与排查思路在实际开发和运行中你可能会遇到以下问题问题现象可能原因排查思路与解决方案应用启动失败端口冲突8080端口被其他进程占用1. 使用netstat -ano | findstr :8080(Win) 或lsof -i :8080(Mac/Linux) 查找占用进程并终止。2. 在application.properties中修改server.port8081。调用/api/missions/start后无反应日志也没有请求参数错误或任务数量太大1. 检查URL和参数是否正确taskCount必须是正整数。2. 如果taskCount非常大如10000创建任务和线程池需要时间请耐心等待或先使用较小数字测试。查询任务进度返回404任务ID不正确或任务已从内存中清除1. 确认使用的missionId是创建任务时返回的ID。2. 本示例使用内存存储应用重启后数据会丢失。生产环境需接入数据库。任务状态长时间卡在PROCESSING模拟任务线程被阻塞或死锁1. 检查SimulatedTaskExecutor.execute中的Thread.sleep是否正常。2. 检查线程池配置确保有足够线程执行任务。3. 查看日志是否有未捕获的异常导致线程终止。队伍进度 (progress) 计算不准确calculateProgress方法并发更新问题1. 本示例使用了CopyOnWriteArrayList和AtomicInteger基本能保证线程安全。但在极高并发下计算瞬间可能状态已变。2. 更严谨的做法是使用锁或原子引用或将进度计算也放入任务完成的回调中。内存占用持续增长任务对象和线程池未释放1. 本示例在线程池任务完成后调用了shutdown()但missionStore和teamStore会一直保存历史数据。2. 生产环境需要增加任务清理机制或使用具有TTL生存时间的缓存。6. 最佳实践与工程建议将上述Demo升级到生产可用系统需要考虑更多工程化细节持久化存储绝对不要在生产环境使用内存Map存储任务状态。应集成数据库如 MySQL记录任务元数据、Redis存储实时进度、作为缓存。设计合理的表结构对应Mission,Team,Task实体并建立索引如对status,create_time字段。可观测性日志使用SLF4J配合Logback/Log4j2为不同级别INFO, WARN, ERROR和不同组件执行器、调度器配置独立的Appender和日志文件。监控集成Micrometer和Prometheus暴露任务队列长度、线程池活跃线程数、任务执行耗时分布Histogram、成功率等关键指标。链路追踪为每个Mission和Task生成唯一的traceId方便在分布式系统中追踪一个任务的完整生命周期。高可用与弹性线程池配置根据任务类型IO密集型、CPU密集型合理设置线程池参数核心线程数、最大线程数、队列容量、拒绝策略。建议使用ThreadPoolExecutor构造函数而非Executors工厂方法以便更精细控制。故障转移如果部署多个实例需要使用分布式锁如Redis RedLock来保证同一个任务不会被重复调度。任务状态变更应具有幂等性。优雅停机在应用关闭时收到SIGTERM信号应等待正在执行的任务完成并拒绝新任务记录中断点以便重启后恢复。配置化与扩展性策略模式将任务拆分策略如奇偶分、按权重分、按哈希分、任务分配策略从代码中抽象出来通过配置文件或数据库配置进行切换。执行器注册中心可以设计一个ExecutorRegistry支持动态注册不同类型的TaskExecutor如调用HTTP接口、执行Shell脚本、处理消息队列系统根据任务类型自动路由。安全与权限API认证为任务创建、查询等接口添加API Key、JWT Token或OAuth2认证。权限控制不同用户或角色只能操作自己有权限的任务。在查询和更新时必须在Service层加入权限校验。输入校验对所有API参数进行严格校验防止非法输入导致系统异常。任务依赖与工作流当前模型是简单的并行。复杂场景下任务间可能存在依赖关系A任务成功后才能执行B。此时需要引入DAG有向无环图调度引擎如Apache Airflow的核心思想。通过遵循这些最佳实践你可以将一个简单的演示项目逐步演进为一个健壮、可扩展、易于维护的企业级任务调度与处理平台。
返回列表