Java异步处理实例

采用的实现流程:

用户发起异步批量申请请求 ---> Service创建任务对象存入任务池中,并存入redis中 ---> Service方法中通过TaskId异步调用执行器 ---> 执行器中通过TaskId去redis中读取任务信息 ---> 执行器执行后更新redis任务状态

image.png

这里需要注意一点,也是我在写这块内容的时候发现的,也是异步处理的关键点,这个返回给用户TaskId并不是在执行完任务之后返回给用户的,这个TaskId是在创建存入Redis后直接返回给用户的,至于任务异步处理都是在后台自己执行的。 这时候就要问了,那我怎么知道怎么这个任务是否完成了,我们需要再写一个接口去查看redis中任务完成情况。 采用redis作为缓存,异步执行缓存内任务实现异步操作

具体项目实施:

模块异步配置文件:

java
复制代码
@Configuration @EnableAsync public class AsyncConfig { /** * 批量审核异步线程池 * <p> * Bean 名称 {@code AuditBatchExecutor},供批量提交、批量审批的 * {@code @Async("AuditBatchExecutor")} 方法使用。 * </p> * * @return 线程池执行器 */ @Bean("AuditBatchExecutor") public Executor AuditBatchExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); // 1. 核心线程数:常驻工作线程 executor.setCorePoolSize(2); // 2. 最大线程数:高峰时可扩容 executor.setMaxPoolSize(5); // 3. 队列容量:等待执行的任务上限 executor.setQueueCapacity(100); // 4. 线程名前缀:便于日志区分异步线程 executor.setThreadNamePrefix("cbs-audit-batch-"); // 5. 队列满时由调用线程执行,避免任务直接丢弃 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }

存入Redis操作: 构建一个任务对象

java
复制代码
AuditBatchTask task = new AuditBatchTask(); // 创建对象AuditBatchTask需要的参数 // ... // 调用异步批量的函数方法,创建任务并且存入Redis中 batchTaskStore.createProcessingTask(task); // 执行异步处理,通过TaskId,这个是一个执行器创建的对象 batchAsyncExecutor.execute(taskId, request); // 这里可以直接返回给前端,任务Id,处理状态、任务处理条数,这些都根据自己实际情况来 return ...

方法createProcessingTask()实现:

java
复制代码
@Override public void createProcessingTask(CbsPaymentAuditBatchTask task) { // Constants.BATCH_TASK_REDIS_TTL_HOURS为设定的常量过期时间 redisCache.set(buildSubmitKey(task.getTaskId()), task, Constants.BATCH_TASK_REDIS_TTL_HOURS, TimeUnit.HOURS); }

异步执行器中batchAsyncExecutor.execute(taskId, request)方法实现:

java
复制代码
//lazy 不加这个会出现Resource循环 @Autowired @Lazy // 这个就根据你的具体 private PaymentAuditService PaymentAuditService; @Autowired private AuditBatchTaskStore batchTaskStore; // 异步处理注解 根据我之前学习注解的经验,这个就是一个说明执行器方法在class中的名字 @Async("cbsAuditBatchExecutor") public void execute(String taskId, AuditBatchSubmitRequest request) { // 去redis中获取task任务,判断任务是否存在 AuditBatchTask task = batchTaskStore.getTask(taskId); if (task == null) { return; } // 2. 调用 try { // 具体执行,这个根据具体情况执行 AuditBatchTaskStatusVO statusVO = AuditService.processBatchSubmit(request); // 创建更新后的任务对象,AuditConstants.BATCH_TASK_STATUS_COMPLETED我设定的常量 task.setTaskStatus(AuditConstants.BATCH_TASK_STATUS_COMPLETED); task.setProcessedCount(statusVO.getProcessedCount()); task.setSuccessCount(statusVO.getSuccessCount()); task.setFailCount(statusVO.getFailCount()); task.setResults(statusVO.getResults()); task.setFinishTime(LocalDateTime.now()); batchTaskStore.updateTask(task); } catch (Exception e) { task.setTaskStatus(CbsAuditConstants.BATCH_TASK_STATUS_FAILED); task.setErrorMessage(e.getMessage()); task.setFinishTime(LocalDateTime.now()); batchTaskStore.updateTask(task); } }
0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP