ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Spring Boot事件驱动异步任务实战:线程池与状态管理

Spring Boot事件驱动异步任务实战:线程池与状态管理 大家好我是老猫。最近在业务里遇到一个比较典型的性能问题接口本身逻辑不复杂但同步执行时耗时不稳定高峰期偶尔会把请求线程拖住导致整个服务出现连锁超时。当时没有立刻引入消息中间件而是先用 Spring 的事件机制加异步线程池做了一版轻量改造效果非常明显。趁这次机会我把这套方案整理成一个从零搭建的小项目项目代号就叫The Lamp and the Genie。名字听起来有点童话映射到技术上其实很好理解Lamp灯是任务的提交入口Genie精灵是后台干活的任务处理器。你擦亮灯精灵就会响应你的愿望对应到系统里就是你提交一个业务请求后台异步任务负责执行后续的逻辑最后再把结果写回状态。整套项目用 Spring Boot 3.x 实现包含任务实体、事件发布、异步监听、线程池配置、状态查询和失败重试代码都是可以直接复制运行的。本文适合有一定 Spring Boot 基础、正在做接口性能优化或异步任务设计的开发者。如果你之前只听说过Async但不太清楚事件发布和事务之间的关系这篇文章应该能帮你补上最后一环。1. 背景与核心概念1.1 业务诉求为什么需要异步处理先看一个很常见的同步场景。用户在前端点击“生成月度报表”后端接口需要做以下事情校验用户权限和参数。查询大量业务数据。聚合、计算、生成 Excel 文件。上传到对象存储或发送文件下载链接。这个流程如果全部在接口线程里同步执行接口响应时间可能会从几十毫秒飙升到几秒甚至几十秒。一旦并发量上来Tomcat 默认工作线程会被长时间占用其他请求全部排队最终表现为系统整体响应变慢。异步处理的核心思路就是把耗时的非核心操作从请求线程中分离出去。接口只负责校验、保存任务状态、返回“已受理”真正耗费时间的逻辑放到后台线程池异步执行。这样做的好处非常直接接口响应时间大幅下降。请求线程可以快速释放。系统吞吐量显著提升。核心业务逻辑和非核心任务逻辑解耦。1.2 Lamp 与 Genie事件驱动模型设计“The Lamp and the Genie”项目借鉴了“神灯与精灵”的意象把它设计成一个轻量的事件驱动任务处理系统。项目中我们定义两个核心角色角色名称职责Lamp灯任务提交入口负责接收请求、保存任务状态、发布事件Genie精灵后台异步工作者负责消费事件、执行任务、更新任务结果整个流程是这样的客户端提交一个“愿望”业务请求系统先创建一条任务记录状态为PENDING。系统发布一个任务创建事件。Genie 监听器在事件发布后开始异步执行任务。任务执行过程中状态变为PROCESSING。执行成功状态变为SUCCESS失败则变为FAILED并将错误信息记录下来。这种方式把“请求”和“执行”彻底解耦。请求线程只负责交给 Lamp至于精灵几点干完、中间是否失败重试请求线程完全不关心。1.3 本地事件、异步调用与消息队列的区别很多初学者会把“异步执行”和“消息队列”混在一起实际上它们是不同层次的东西。异步调用是方法层面的执行方式比如Async注解修饰的方法会丢进线程池异步执行本质还是同一个 JVM 进程内的调用。本地事件是对象层面的解耦发布者不需要知道谁在监听监听者也不需要知道发布者是谁。最简单的方式是使用 Spring 的ApplicationEventPublisher发布事件配合EventListener监听。消息队列是跨进程、跨系统的通信方式比如 RocketMQ、RabbitMQ、Kafka。消息会持久化到 Broker生产者和消费者完全独立适合分布式场景。三者关系可以这样理解方案范围持久化可靠性适用场景Async异步调用JVM 进程内无低单纯不阻塞主线程Spring 本地事件JVM 进程内无中模块解耦、同进程任务消息队列跨进程有高分布式、削峰、可靠投递The Lamp and the Genie 项目采用了本地事件 异步监听的组合适合单体应用或较小规模的异步任务。如果未来业务增长到需要分布式可靠投递整个设计可以平移到消息队列模式只需要替换事件发布和监听部分即可。1.4 本文项目目标这个项目不是单纯演示Async的 Hello World而是要完成一套带状态的、可查询、可重试的异步任务闭环。看完本文并动手实践后你将掌握如何使用EnableAsync开启异步能力。如何自定义线程池并应用到异步任务。如何使用TransactionalEventListener保证事务提交后再执行业务。如何设计任务状态机。如何实现任务查询、失败重试和幂等控制。2. 环境准备与项目结构2.1 技术栈与版本说明本文示例的完整技术栈如下JDK 17Maven 3.8Spring Boot 3.xSpring Data JPAH2 内存数据库Spring ValidationSpring Web版本需要根据你的项目实际情况调整。本文示例以常见的 Spring Boot 3.x 为例重点演示配置思路如果你使用 Spring Boot 2.x也可以参考同样的设计部分依赖坐标会略有差异。2.2 创建项目骨架你可以直接通过 Spring Initializr 生成基础项目选择以下依赖Spring WebSpring Data JPAH2 DatabaseValidation也可以手动创建 Maven 工程把下面第 4 章里的pom.xml拷贝进去。项目名称建议使用lamp-genie-demo包名使用com.example.lampgenie。2.3 总体模块划分项目代码量不大但为了后面讲解清晰我先把模块职责拆开模块职责entity任务实体WishRecordrepository数据访问层WishRepositorydto请求参数WishRequestevent事件定义WishCreatedEventservice业务层WishService负责提交任务listener异步监听器GenieListener负责消费任务config异步线程池配置AsyncConfigcontroller对外接口层WishController这种拆分方式符合 Spring Boot 项目的基本规范后续从本地事件切换成消息队列时只需要替换listener和service层的实现。3. 核心设计任务模型与状态流转3.1 任务状态机异步任务系统最重要的一环是状态管理。状态设计不合理后面做补偿、重试、统计都会非常痛苦。The Lamp and the Genie 项目中任务状态机设计如下状态含义可流向PENDING任务已受理等待执行PROCESSINGPROCESSING任务正在执行中SUCCESS、FAILEDSUCCESS任务执行成功无FAILED任务执行失败PENDING重试项目启动时受理接口创建的任务初始状态一定是PENDING。Genie 异步开始执行时立刻把状态更新为PROCESSING防止任务被重复执行。任务执行完成后根据结果更新为SUCCESS或FAILED。这里有一个容易踩的坑如果执行期间服务重启PROCESSING状态的任务会卡死。生产环境中通常需要引入超时重置机制把超过一定时间仍处于PROCESSING的任务重新置回PENDING或标记为异常。3.2 事件流转图整个项目的流转过程用文字可以描述为客户端 POST 请求 - WishController 接收参数 - WishService 保存 WishRecordPENDING - 发布 WishCreatedEvent - 事务提交后触发 GenieListener - 异步线程池执行任务 - 更新状态 PROCESSING - 执行业务逻辑 - 更新状态 SUCCESS / FAILED如果你已经接触过 RocketMQ 或 RabbitMQ会发现这个流程与消息队列的“生产 - 消费”模型非常相似。事件发布者就是生产者异步监听器就是消费者Spring 的事件机制只是把消息队列的“跨进程传递”简化成了“JVM 内传递”。3.3 表结构设计任务表名为wish_record字段设计如下字段类型说明idBIGINT主键自增request_idVARCHAR(64)业务幂等键唯一contentVARCHAR(512)任务内容statusVARCHAR(32)任务状态error_messageVARCHAR(1024)错误信息retry_countINT重试次数created_atDATETIME创建时间updated_atDATETIME更新时间request_id由客户端生成服务端用它做幂等判断。如果同一个request_id重复提交直接返回已有记录避免重复创建任务。4. 实战代码实现下面开始写完整代码。你可以按照文件路径逐个创建。4.1 pom.xml 依赖文件路径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 version3.2.5/version relativePath/ /parent groupIdcom.example/groupId artifactIdlamp-genie-demo/artifactId version1.0.0/version namelamp-genie-demo/name descriptionThe Lamp and the Genie demo project/description properties java.version17/java.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-jpa/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-validation/artifactId /dependency dependency groupIdcom.h2database/groupId artifactIdh2/artifactId scoperuntime/scope /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependencies build plugins plugin groupIdorg.springframework.boot/groupId artifactIdspring-boot-maven-plugin/artifactId /plugin /plugins /build /project这里使用 H2 内存数据库方便本地直接运行不需要额外安装 MySQL。如果你要切换到 MySQL只需要替换依赖并在application.yml中修改数据源配置。4.2 application.yml 配置文件路径src/main/resources/application.ymlserver: port: 8080 spring: application: name: lamp-genie-demo datasource: url: jdbc:h2:mem:wishdb;DB_CLOSE_DELAY-1 driver-class-name: org.h2.Driver username: sa password: jpa: hibernate: ddl-auto: update show-sql: true properties: hibernate: format_sql: true h2: console: enabled: true path: /h2-console logging: level: com.example.lampgenie: DEBUG需要注意H2 是内存模式服务重启后数据会丢失。本项目定位是演示所以内存库足够。生产环境建议使用独立数据库并且关闭ddl-auto: update改用 Flyway 或 Liquibase 管理表结构。4.3 枚举与实体文件路径src/main/java/com/example/lampgenie/entity/WishStatus.javapackage com.example.lampgenie.entity; public enum WishStatus { PENDING, PROCESSING, SUCCESS, FAILED }文件路径src/main/java/com/example/lampgenie/entity/WishRecord.javapackage com.example.lampgenie.entity; import jakarta.persistence.*; import java.time.LocalDateTime; Entity Table(name wish_record, uniqueConstraints UniqueConstraint(columnNames requestId)) public class WishRecord { Id GeneratedValue(strategy GenerationType.IDENTITY) private Long id; Column(nullable false, length 64) private String requestId; Column(nullable false, length 512) private String content; Enumerated(EnumType.STRING) Column(nullable false, length 32) private WishStatus status; Column(length 1024) private String errorMessage; Column(nullable false) private Integer retryCount 0; Column(nullable false) private LocalDateTime createdAt; Column(nullable false) private LocalDateTime updatedAt; PrePersist public void prePersist() { LocalDateTime now LocalDateTime.now(); createdAt now; updatedAt now; if (retryCount null) { retryCount 0; } } PreUpdate public void preUpdate() { updatedAt LocalDateTime.now(); } public Long getId() { return id; } public void setId(Long id) { this.id id; } public String getRequestId() { return requestId; } public void setRequestId(String requestId) { this.requestId requestId; } public String getContent() { return content; } public void setContent(String content) { this.content content; } public WishStatus getStatus() { return status; } public void setStatus(WishStatus status) { this.status status; } public String getErrorMessage() { return errorMessage; } public void setErrorMessage(String errorMessage) { this.errorMessage errorMessage; } public Integer getRetryCount() { return retryCount; } public void setRetryCount(Integer retryCount) { this.retryCount retryCount; } public LocalDateTime getCreatedAt() { return createdAt; } public void setCreatedAt(LocalDateTime createdAt) { this.createdAt createdAt; } public LocalDateTime getUpdatedAt() { return updatedAt; } public void setUpdatedAt(LocalDateTime updatedAt) { this.updatedAt updatedAt; } }实体类中使用了 JPA 注解PrePersist和PreUpdate负责在插入、更新前自动填充时间字段这样业务代码里就不用手动维护时间了。4.4 Repository 仓库层文件路径src/main/java/com/example/lampgenie/repository/WishRepository.javapackage com.example.lampgenie.repository; import com.example.lampgenie.entity.WishRecord; import com.example.lampgenie.entity.WishStatus; import org.springframework.data.jpa.repository.JpaRepository; import java.util.List; import java.util.Optional; public interface WishRepository extends JpaRepositoryWishRecord, Long { OptionalWishRecord findByRequestId(String requestId); ListWishRecord findByStatus(WishStatus status); }findByRequestId是幂等控制的关键方法提交任务前先根据requestId查询是否已经存在。4.5 事件定义文件路径src/main/java/com/example/lampgenie/event/WishCreatedEvent.javapackage com.example.lampgenie.event; public record WishCreatedEvent(Long wishId) { }这里使用 Java 17 的record定义事件对象简洁且不可变。事件里只需要携带任务主键wishId监听器拿到主键后再查询完整数据避免事件对象携带过多数据导致序列化开销。4.6 LampService 提交任务文件路径src/main/java/com/example/lampgenie/service/WishService.javapackage com.example.lampgenie.service; import com.example.lampgenie.dto.WishRequest; import com.example.lampgenie.entity.WishRecord; import com.example.lampgenie.entity.WishStatus; import com.example.lampgenie.event.WishCreatedEvent; import com.example.lampgenie.repository.WishRepository; import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; Service public class WishService { private final WishRepository wishRepository; private final ApplicationEventPublisher eventPublisher; public WishService(WishRepository wishRepository, ApplicationEventPublisher eventPublisher) { this.wishRepository wishRepository; this.eventPublisher eventPublisher; } Transactional public WishRecord submit(WishRequest request) { return wishRepository.findByRequestId(request.requestId()) .map(existing - { // 幂等处理重复提交直接返回已有记录 return existing; }) .orElseGet(() - { WishRecord record new WishRecord(); record.setRequestId(request.requestId()); record.setContent(request.content()); record.setStatus(WishStatus.PENDING); record.setRetryCount(0); wishRepository.save(record); // 发布事件Genie 会在事务提交后异步消费 eventPublisher.publishEvent(new WishCreatedEvent(record.getId())); return record; }); } Transactional public void republish(String requestId) { WishRecord record wishRepository.findByRequestId(requestId) .orElseThrow(() - new IllegalArgumentException(wish not found: requestId)); record.setStatus(WishStatus.PENDING); record.setErrorMessage(null); record.setRetryCount(record.getRetryCount() 1); wishRepository.save(record); eventPublisher.publishEvent(new WishCreatedEvent(record.getId())); } }重点看一下submit方法里的幂等逻辑同一个requestId第二次提交时不会重复创建任务而是直接返回已有记录。这条逻辑在接口超时重试、前端重复点击时非常有用。事件发布的位置也值得注意publishEvent是在Transactional方法里调用的但事件真正触发是由监听器那里的TransactionalEventListener控制保证只有事务提交成功后Genie 才开始执行。这样 Genie 查询任务时一定能查到刚才保存的数据。4.7 GenieListener 异步执行文件路径src/main/java/com/example/lampgenie/listener/GenieListener.javapackage com.example.lampgenie.listener; import com.example.lampgenie.entity.WishRecord; import com.example.lampgenie.entity.WishStatus; import com.example.lampgenie.event.WishCreatedEvent; import com.example.lampgenie.repository.WishRepository; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Component; import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalEventListener; Component public class GenieListener { private static final Logger log LoggerFactory.getLogger(GenieListener.class); private final WishRepository wishRepository; public GenieListener(WishRepository wishRepository) { this.wishRepository wishRepository; } Async(genieExecutor) TransactionalEventListener(phase TransactionPhase.AFTER_COMMIT) public void onWishCreated(WishCreatedEvent event) { WishRecord record wishRepository.findById(event.wishId()).orElse(null); if (record null) { log.warn(任务不存在wishId{}, event.wishId()); return; } // 1. 更新为处理中 record.setStatus(WishStatus.PROCESSING); wishRepository.save(record); log.info(Genie 开始处理任务requestId{}, wishId{}, record.getRequestId(), record.getId()); try { // 2. 模拟业务执行耗时 Thread.sleep(2000L); // 3. 模拟任务失败场景任务内容包含 fail 关键字 if (record.getContent().contains(fail)) { throw new IllegalStateException(模拟业务执行失败); } // 4. 更新为成功 record.setStatus(WishStatus.SUCCESS); record.setErrorMessage(null); wishRepository.save(record); log.info(Genie 任务执行成功requestId{}, record.getRequestId()); } catch (Exception e) { // 5. 更新为失败并记录错误信息 record.setStatus(WishStatus.FAILED); record.setErrorMessage(e.getMessage()); wishRepository.save(record); log.error(Genie 任务执行失败requestId{}, error{}, record.getRequestId(), e.getMessage(), e); } } }这个监听器是整个项目的核心。Async(genieExecutor)让它在线程池中异步执行TransactionalEventListener(phase TransactionPhase.AFTER_COMMIT)保证在事务提交后才运行。模拟失败逻辑使用任务内容包含fail关键字来触发这样方便后面测试重试接口。4.8 异步线程池配置文件路径src/main/java/com/example/lampgenie/config/AsyncConfig.javapackage com.example.lampgenie.config; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.scheduling.annotation.AsyncConfigurer; import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; import java.util.concurrent.Executor; import java.util.concurrent.ThreadPoolExecutor; Configuration EnableAsync public class AsyncConfig implements AsyncConfigurer { Bean(genieExecutor) public Executor genieExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); // 核心线程数 executor.setCorePoolSize(2); // 最大线程数 executor.setMaxPoolSize(5); // 队列容量 executor.setQueueCapacity(200); // 线程名称前缀 executor.setThreadNamePrefix(genie-); // 拒绝策略由调用者线程执行 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }线程池参数需要根据实际机器配置和任务量调整。这里给出一组保守参数核心线程 2最大线程 5队列容量 200。拒绝策略采用 CallerRunsPolicy意思是当线程池和队列都满了新任务会由提交任务的线程自己执行这样可以避免任务直接丢弃代价是调用线程会被阻塞。后续章节会详细分析参数设置。4.9 控制器接口层文件路径src/main/java/com/example/lampgenie/controller/WishController.javapackage com.example.lampgenie.controller; import com.example.lampgenie.dto.WishRequest; import com.example.lampgenie.entity.WishRecord; import com.example.lampgenie.repository.WishRepository; import com.example.lampgenie.service.WishService; import jakarta.validation.Valid; import org.springframework.http.HttpStatus; import org.springframework.web.bind.annotation.*; import org.springframework.web.server.ResponseStatusException; RestController RequestMapping(/api/wishes) public class WishController { private final WishService wishService; private final WishRepository wishRepository; public WishController(WishService wishService, WishRepository wishRepository) { this.wishService wishService; this.wishRepository wishRepository; } /** * 提交愿望任务接口快速返回 */ PostMapping public WishRecord submit(Valid RequestBody WishRequest request) { return wishService.submit(request); } /** * 根据 requestId 查询任务状态 */ GetMapping(/{requestId}) public WishRecord status(PathVariable String requestId) { return wishRepository.findByRequestId(requestId) .orElseThrow(() - new ResponseStatusException( HttpStatus.NOT_FOUND, wish not found)); } /** * 失败任务重试 */ PostMapping(/{requestId}/retry) public WishRecord retry(PathVariable String requestId) { return wishRepository.findByRequestId(requestId) .map(record - { if (record.getStatus() ! com.example.lampgenie.entity.WishStatus.FAILED) { throw new IllegalStateException(只有失败状态的任务才能重试); } wishService.republish(requestId); return wishRepository.findByRequestId(requestId).orElseThrow(); }) .orElseThrow(() - new ResponseStatusException( HttpStatus.NOT_FOUND, wish not found)); } }文件路径src/main/java/com/example/lampgenie/dto/WishRequest.javapackage com.example.lampgenie.dto; import jakarta.validation.constraints.NotBlank; public record WishRequest( NotBlank(message requestId不能为空) String requestId, NotBlank(message content不能为空) String content ) { }控制器层保持轻薄只做参数接收和状态查询。注意retry接口里先检查任务是否处于FAILED状态只有失败任务才允许重新投递。4.10 启动类文件路径src/main/java/com/example/lampgenie/LampGenieApplication.javapackage com.example.lampgenie; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; SpringBootApplication public class LampGenieApplication { public static void main(String[] args) { SpringApplication.run(LampGenieApplication.class, args); } }到这里整个项目的核心代码已经完成。你可以直接用 Maven 启动。5. 运行与验证5.1 启动项目在项目根目录执行mvn spring-boot:run启动成功后控制台会输出类似日志Tomcat started on port 8080 Started LampGenieApplication in 3.2 seconds同时可以看到 Spring Data JPA 自动创建wish_record表的日志。5.2 提交“愿望”任务打开另一个终端执行以下请求curl -X POST http://localhost:8080/api/wishes \ -H Content-Type: application/json \ -d {requestId: req-001, content: 请帮我生成一份月度业务报表}预期返回{ id: 1, requestId: req-001, content: 请帮我生成一份月度业务报表, status: PENDING, errorMessage: null, retryCount: 0, createdAt: 2025-01-01T12:00:00, updatedAt: 2025-01-01T12:00:00 }接口瞬间返回响应时间不会因为休眠 2 秒而变慢。这就是异步任务的核心效果。5.3 查询任务状态执行curl http://localhost:8080/api/wishes/req-001如果任务已经执行完成查询结果中status应该为SUCCESS。5.4 观察异步执行日志服务端日志会输出类似内容Genie 开始处理任务requestIdreq-001, wishId1 Genie 任务执行成功requestIdreq-001注意线程名称前缀是genie-这说明任务确实是在异步线程池中执行的而不是 Tomcat 请求线程。5.5 重试失败任务我们再来模拟一个失败场景。请求内容里包含fail关键字curl -X POST http://localhost:8080/api/wishes \ -H Content-Type: application/json \ -d {requestId: req-fail, content: 这个任务会 fail}等待 3 秒左右查询curl http://localhost:8080/api/wishes/req-fail此时status应为FAILED并且errorMessage记录了模拟异常信息。然后执行重试接口curl -X POST http://localhost:8080/api/wishes/req-fail/retry注意这里的重试接口返回的是PENDING状态但异步线程会再次执行任务。由于任务内容仍然包含fail关键字最终会再次失败。你可以把任务内容改成正常文本后重新提交即可看到从失败到成功的完整生命周期。6. 常见问题与排查思路问题现象常见原因解决思路Async不生效方法同步执行在同一个类内部调用异步方法代理未生效把异步方法放到独立的 Bean 中或者通过依赖注入调用自身代理事件监听器查询不到任务数据事件在事务提交前触发数据还没提交使用TransactionalEventListener(phase AFTER_COMMIT)任务丢失监听器内部异常没有记录没有捕获异常在异步方法内使用 try-catch 包裹失败时更新任务状态异步线程池被占满大量任务堆积线程池参数设置不合理根据任务耗时、并发量合理设置 corePoolSize、maxPoolSize、queueCapacity重复提交导致重复执行业务没有做幂等控制增加唯一索引和findByRequestId校验服务重启后任务卡在PROCESSING内存状态没有恢复机制引入启动扫描重置超时未完成的任务下面针对几个高频问题展开说明。6.1 Async 不生效的问题Async是基于 Spring AOP 代理实现的。如果你在WishService内部直接调用同类中的Async方法调用的是this对象的方法不会经过代理异步自然不生效。比如这样写是有问题的Service public class BadExampleService { public void doSomething() { // 同类调用不会走代理 this.asyncMethod(); } Async public void asyncMethod() { // 仍然是同步执行 } }正确的做法是把异步监听器放在独立的 Bean 中也就是本文中GenieListener的写法。调用方通过 Spring 容器注入代理对象代理才会拦截方法并交给线程池执行。6.2 事件监听器查询不到刚保存的记录如果使用普通的EventListener监听器会在publishEvent调用时立刻执行。但此时Transactional方法还没提交事务数据库中还查不到新记录。解决办法是使用TransactionalEventListener(phase TransactionPhase.AFTER_COMMIT)它会把监听器的执行时机推迟到事务提交之后。这既解决了数据可见性问题也保证了业务操作和异步任务的一致性。6.3 任务丢失的问题异步方法的异常默认不会传播到主线程。如果监听器内抛出异常且没有捕获任务很可能就“悄无声息”地失败了任务状态永远不会更新。所以异步处理器中一定要有兜底逻辑。本文的做法是在onWishCreated方法内使用 try-catch把任务状态更新为FAILED并记录errorMessage。这样即使业务逻辑异常状态机仍然完整。7. 最佳实践与工程建议7.1 异步线程池参数设计线程池参数不能拍脑袋。考虑以下公式如果是 CPU 密集型任务线程数一般设置为CPU 核心数 1。如果是 IO 密集型任务线程数可以设置得更大比如CPU 核心数 * 2。队列容量取决于系统允许的最大等待量。拒绝策略必须是明确的不能默认丢弃。ThreadPoolTaskExecutor的核心参数参数说明corePoolSize核心线程数即使空闲也不会回收maxPoolSize最大线程数超过核心线程数且队列满时才继续创建queueCapacity阻塞队列容量keepAliveSeconds非核心线程空闲存活时间rejectedExecutionHandler拒绝策略要特别说明ThreadPoolTaskExecutor的一个特性任务先进入队列队列满了才会创建额外线程而不是核心线程满了立即创建最大线程。所以如果queueCapacity设置过大maxPoolSize可能永远不会生效。生产环境下建议给线程池增加监控记录活跃线程数、队列长度、拒绝次数。这些指标能帮助你判断参数是否需要调优。7.2 事务与事件发布的边界事件发布和事务提交是异步设计中最容易出问题的地方。记住三个原则在事务内发布事件保证事件与业务操作在同一个事务里。在事务提交后消费事件使用TransactionalEventListener(phase AFTER_COMMIT)。监听器内部要有独立的事务控制不要依赖发布者的事务。如果对可靠性要求极高本地事件机制是不够的。因为 JVM 崩溃时未消费的事件会丢失。这种情况下应该引入消息队列利用 Broker 的持久化能力保证事件不丢。The Lamp and the Genie 的项目结构可以平滑迁移到 MQ 方案因为监听器和事件对象都是现成的。7.3 幂等与防重设计面向用户的接口一定要考虑重复请求。重复点击、前端重试、网关重发任何一个环节都可能导致同一笔业务被重复处理。本项目使用requestId 唯一索引实现幂等。具体做法是数据库wish_record表对request_id添加唯一约束。提交时先查询查到就直接返回已有记录。未查询到就插入新记录。在高并发场景下先查后插会有并发窗口。更稳妥的方式是直接插入捕获唯一索引冲突后重新查询。生产环境建议把这一步做好避免超卖、重复建单等严重问题。7.4 重试与失败补偿异步任务一定会遇到失败。失败后的处理策略取决于业务类型业务类型推荐策略通知类短信、邮件允许失败记录日志即可数据同步类自动重试 2-3 次间隔递增账户/订单类进入人工审核队列禁止无限自动重试重试时要控制次数防止网络抖动恢复后大量积压任务同时冲击下游系统。建议使用指数退避算法比如第一次等待 1 秒第二次 5 秒第三次 30 秒。对于重试 N 次仍然失败的任务可以标记为FINAL_FAILED或推送告警由人工介入处理。7.5 监控与日志异步任务的监控和同步调用不同问题不会立刻体现在接口耗时上而是积累在后台慢慢爆发。建议至少监控以下指标任务积压数量PENDING状态数量。执行中任务数量PROCESSING状态数量。任务成功率、失败率。线程池队列深度。线程池拒绝次数。日志方面每个任务从提交到完成建议串联同一个requestId这样排查问题时可以通过日志快速检索全链路。本文的日志已经包含了requestId在实际项目中可以把requestId放入 MDC让所有关联日志自动带上这个标识。7.6 安全与参数校验异步任务虽然不直接面对用户但入参校验仍然不能省略。项目中WishRequest使用了NotBlank校验防止空参进入系统。另外要注意任务内容如果可能包含用户输入后续打印日志时要避免输出完整敏感信息。生产环境中建议对用户信息做脱敏处理比如手机号、身份证号、银行卡号等字段只保留前后几位。8. 总结与扩展方向这套基于The Lamp and the Genie思想实现的异步任务系统麻雀虽小但五脏俱全。核心亮点有这么几个一是通过事件机制把“任务提交”和“任务执行”彻底解耦接口响应性能大幅提升二是使用事务提交后监听避免了数据不可见的坑三是设计了 PENDING、PROCESSING、SUCCESS、FAILED 四种状态让任务可追踪、可重试四是把幂等、线程池、失败重试这些生产级话题都落地到了代码里。如果只是单机应用这套方案完全够用。说实话不少小型项目的异步任务做得还不如这个规范比如任务没有状态、失败没有重试、线程池参数随意写这些问题在真实业务里迟早会暴露。如果你想继续深入可以从以下几个方向扩展把本地事件替换成 RocketMQ 或 RabbitMQ实现跨服务异步任务。增加分布式任务调度比如集成 XXL-JOB 或 ElasticJob解决服务重启后的任务恢复问题。引入定时任务扫描把长时间卡在PROCESSING的任务自动重置。扩展告警能力任务连续失败时通过钉钉、企业微信等渠道通知开发人员。为线程池增加 Prometheus 指标暴露在 Grafana 里配置监控大盘。我个人的建议是先搞清楚本地事件模型再上消息队列。很多团队一开始就引入重型中间件结果架构复杂度和运维成本都上去了业务收益却没那么明显。像本文这种轻量异步方案反而是日常开发中性价比最高的选择。最后分享一个经验异步不是银弹。异步任务让接口变快了但排查问题的难度会上升状态一致性、幂等、重试、监控每一项都需要额外投入。设计时一定要问自己这个操作真的需要异步吗如果用户必须等待结果异步反而会让体验变差。合适的技术用在合适的位置才是工程能力的真正体现。如果你在本地跑通了这套代码建议再改一改线程池参数观察不同配置下的执行效果也可以把 H2 换成 MySQL模拟一下服务重启后的任务状态。动手试过一遍比看十遍文章都管用。希望这篇文章对你有帮助有问题欢迎在评论区一起讨论。
返回列表