Spring Boot实现异步任务案例

wen java案例 3

本文目录导读:

Spring Boot实现异步任务案例

  1. 项目依赖配置
  2. 配置异步任务
  3. 异步任务实现
  4. 控制器层
  5. 异步任务管理器(进阶用法)
  6. 异步回调处理
  7. 测试代码
  8. 注意事项

我来给你一个完整的Spring Boot异步任务案例,包括基础的异步执行、异步回调处理,以及最佳实践。

项目依赖配置

pom.xml

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
</dependencies>

配置异步任务

异步配置类

package com.example.async.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import java.util.concurrent.Executor;
import java.util.concurrent.ThreadPoolExecutor;
@Configuration
public class AsyncConfig {
    @Bean("taskExecutor")
    public Executor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        // 核心线程数
        executor.setCorePoolSize(5);
        // 最大线程数
        executor.setMaxPoolSize(10);
        // 队列容量
        executor.setQueueCapacity(100);
        // 线程名称前缀
        executor.setThreadNamePrefix("async-task-");
        // 拒绝策略:调用者执行
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        // 初始化
        executor.initialize();
        return executor;
    }
    // 可以配置多个线程池处理不同业务
    @Bean("emailExecutor")
    public Executor emailExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(3);
        executor.setMaxPoolSize(6);
        executor.setQueueCapacity(50);
        executor.setThreadNamePrefix("email-task-");
        executor.initialize();
        return executor;
    }
}

或者使用配置文件方式:

# application.yml
spring:
  task:
    execution:
      pool:
        core-size: 5
        max-size: 10
        queue-capacity: 100
        keep-alive: 60s
      thread-name-prefix: async-task-
      shutdown:
        await-termination: true
        await-termination-period: 60s

异步任务实现

基础异步任务类

package com.example.async.service;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture;
@Service
public class AsyncTaskService {
    private static final Logger logger = LoggerFactory.getLogger(AsyncTaskService.class);
    /**
     * 无返回值异步任务
     */
    @Async("taskExecutor")
    public void sendEmail(String to, String content) {
        logger.info("开始发送邮件到: {}", to);
        try {
            // 模拟耗时操作
            Thread.sleep(2000);
            logger.info("邮件发送成功: {}", to);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            logger.error("邮件发送失败: {}", to, e);
        }
    }
    /**
     * 带返回值的异步任务
     */
    @Async("taskExecutor")
    public CompletableFuture<String> generateReport(String reportType) {
        logger.info("开始生成报表: {}", reportType);
        try {
            Thread.sleep(3000);
            String result = "报表-" + reportType + "-" + System.currentTimeMillis();
            logger.info("报表生成完成: {}", result);
            return CompletableFuture.completedFuture(result);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            return CompletableFuture.failedFuture(e);
        }
    }
    /**
     * 处理订单异步任务
     */
    @Async("taskExecutor")
    public void processOrder(Long orderId) {
        logger.info("开始处理订单: {}", orderId);
        // 订单处理逻辑
        try {
            Thread.sleep(1500);
            // 发送通知
            sendNotification(orderId);
            // 更新库存
            updateStock(orderId);
            logger.info("订单处理完成: {}", orderId);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            logger.error("订单处理失败: {}", orderId, e);
        }
    }
    private void sendNotification(Long orderId) {
        logger.info("发送订单通知: {}", orderId);
    }
    private void updateStock(Long orderId) {
        logger.info("更新库存: {}", orderId);
    }
}

异步事件处理

package com.example.async.event;
import org.springframework.context.event.EventListener;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
@Component
public class OrderEventListener {
    private static final Logger logger = LoggerFactory.getLogger(OrderEventListener.class);
    @Async("taskExecutor")
    @EventListener
    public void handleOrderCreatedEvent(OrderCreatedEvent event) {
        logger.info("处理订单创建事件: {}", event.getOrderId());
        // 发送订单创建通知邮件
        try {
            Thread.sleep(2000);
            logger.info("订单通知已发送: {}", event.getOrderId());
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    @Async("taskExecutor")
    @EventListener
    public void handlePaymentSuccessEvent(PaymentSuccessEvent event) {
        logger.info("处理支付成功事件: 订单{} 金额{}", event.getOrderId(), event.getAmount());
        // 更新订单状态等
    }
}
// 事件类
public class OrderCreatedEvent {
    private Long orderId;
    // getter/setter...
}
public class PaymentSuccessEvent {
    private Long orderId;
    private BigDecimal amount;
    // getter/setter...
}

控制器层

package com.example.async.controller;
import com.example.async.service.AsyncTaskService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.multipart.MultipartFile;
import java.util.concurrent.CompletableFuture;
@RestController
@RequestMapping("/api/async")
public class AsyncController {
    @Autowired
    private AsyncTaskService asyncTaskService;
    /**
     * 触发异步任务
     */
    @PostMapping("/email")
    public ApiResponse sendEmail(@RequestParam String to, @RequestParam String content) {
        asyncTaskService.sendEmail(to, content);
        return ApiResponse.success("邮件发送任务已提交");
    }
    /**
     * 获取异步任务结果
     */
    @GetMapping("/report/{type}")
    public CompletableFuture<ApiResponse> generateReport(@PathVariable String type) {
        return asyncTaskService.generateReport(type)
                .thenApply(result -> ApiResponse.success(result));
    }
    /**
     * 批量异步处理
     */
    @PostMapping("/batch-process")
    public ApiResponse batchProcess(@RequestParam("files") MultipartFile[] files) {
        List<CompletableFuture<String>> futures = new ArrayList<>();
        for (MultipartFile file : files) {
            CompletableFuture<String> future = asyncTaskService.generateReport(file.getOriginalFilename());
            futures.add(future);
        }
        // 等待所有任务完成
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
                .thenRun(() -> System.out.println("所有文件处理完成"));
        return ApiResponse.success("批量处理任务已提交,共" + files.length + "个文件");
    }
}

异步任务管理器(进阶用法)

package com.example.async.manager;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Component;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function;
@Component
public class AsyncTaskManager {
    private final Map<String, CompletableFuture<?>> tasks = new ConcurrentHashMap<>();
    public void submitTask(String taskId, Runnable task) {
        CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
            task.run();
        });
        tasks.put(taskId, future);
    }
    public <T> CompletableFuture<T> submitTask(String taskId, Supplier<T> task) {
        CompletableFuture<T> future = CompletableFuture.supplyAsync(task);
        tasks.put(taskId, future);
        return future;
    }
    public CompletableFuture<?> getTaskResult(String taskId) {
        return tasks.get(taskId);
    }
    public boolean isCompleted(String taskId) {
        CompletableFuture<?> future = tasks.get(taskId);
        return future != null && future.isDone();
    }
    public void cancelTask(String taskId) {
        CompletableFuture<?> future = tasks.get(taskId);
        if (future != null) {
            future.cancel(true);
        }
    }
}

异步回调处理

package com.example.async.callback;
import org.springframework.stereotype.Component;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@Component
public class AsyncCallbackHandler {
    /**
     * 演示异步回调处理
     */
    public CompletableFuture<User> getUserInfoAsync(Long userId) {
        return CompletableFuture.supplyAsync(() -> {
            // 模拟异步获取用户信息
            try {
                Thread.sleep(2000);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            User user = new User();
            user.setId(userId);
            user.setName("User" + userId);
            return user;
        }).thenApply(user -> {
            // 异步结果处理
            user.setStatus("PROCESSED");
            return user;
        }).exceptionally(ex -> {
            System.err.println("获取用户信息失败: " + ex.getMessage());
            return null;
        });
    }
    /**
     * 多个异步任务的组合
     */
    public CompletableFuture<OrderInfo> getOrderDetails(Long orderId) {
        CompletableFuture<User> userFuture = getUserInfoAsync(1L);
        CompletableFuture<Order> orderFuture = CompletableFuture.supplyAsync(() -> {
            // 模拟异步获取订单
            return new Order(orderId);
        });
        // 组合多个异步任务的结果
        return userFuture.thenCombine(orderFuture, (user, order) -> {
            OrderInfo info = new OrderInfo();
            info.setUser(user);
            info.setOrder(order);
            return info;
        });
    }
}

测试代码

package com.example.async;
import com.example.async.service.AsyncTaskService;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@SpringBootTest
public class AsyncTaskTest {
    @Autowired
    private AsyncTaskService asyncTaskService;
    @Test
    public void testAsyncTask() throws Exception {
        // 测试无返回值异步任务
        asyncTaskService.sendEmail("test@example.com", "测试邮件");
        System.out.println("主线程继续执行...");
        // 等待异步任务完成
        Thread.sleep(5000);
        // 测试带返回值异步任务
        CompletableFuture<String> future = asyncTaskService.generateReport("SALES");
        String result = future.get(10, TimeUnit.SECONDS);
        System.out.println("报表结果: " + result);
    }
    @Test
    public void testConcurrentAsyncTasks() throws Exception {
        List<CompletableFuture<String>> futures = new ArrayList<>();
        // 同时触发多个异步任务
        for (int i = 0; i < 5; i++) {
            futures.add(asyncTaskService.generateReport("报表" + i));
        }
        // 等待所有任务完成
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
                .join();
        // 获取所有结果
        for (CompletableFuture<String> future : futures) {
            System.out.println("结果: " + future.get());
        }
    }
}

注意事项

常见问题1:同类调用失效

@Service
public class AsyncService {
    public void methodA() {
        // 这样调用不会触发异步,因为代理不会处理同类调用
        this.methodB();
        // 正确方式:注入自身或使用ApplicationContext获取代理
    }
    @Async
    public void methodB() {
        // 异步逻辑
    }
    // 正确方式示例
    @Autowired
    private ApplicationContext context;
    public void methodC() {
        AsyncService proxy = context.getBean(AsyncService.class);
        proxy.methodB(); // 这样才会触发异步
    }
}

自定义异常处理

package com.example.async.config;
import org.springframework.aop.interceptor.AsyncUncaughtExceptionHandler;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.AsyncConfigurer;
import java.lang.reflect.Method;
import java.util.Arrays;
@Configuration
public class AsyncExceptionConfig implements AsyncConfigurer {
    @Override
    public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
        return new AsyncUncaughtExceptionHandler() {
            @Override
            public void handleUncaughtException(Throwable ex, Method method, Object... params) {
                System.err.println("异步任务异常:");
                System.err.println("方法: " + method.getName());
                System.err.println("参数: " + Arrays.toString(params));
                System.err.println("异常信息: " + ex.getMessage());
                // 可以在这里发送告警通知等
            }
        };
    }
}

这个案例涵盖了Spring Boot异步任务的核心用法,包括:

  • 基础异步执行:@Async注解的基本使用
  • 带返回值的异步任务:使用CompletableFuture
  • 异步事件处理:结合Spring事件机制
  • 异步回调处理:链式调用和组合
  • 最佳实践:线程池配置、异常处理、性能优化

根据你的具体业务需求,可以选择合适的方式进行实现。

抱歉,评论功能暂时关闭!