配置类

@Configuration
@EnableAsync
public class AsyncConfig extends AsyncConfigurerSupport {

    @Override
    public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler(){
        return new CusAsyncExceptionHandler();
    }

    @Bean("executor")
    public Executor executor(){
        System.out.println("executor 线程启动---");
        ThreadPoolTaskExecutor scheduler = new ThreadPoolTaskExecutor();
        scheduler.setCorePoolSize(10);    //基本线程数量
        scheduler.setQueueCapacity(100);  //队列最大长度
        scheduler.setMaxPoolSize(30);    //最大线程数量
        scheduler.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        scheduler.setThreadNamePrefix("ASYNC-executor-");
        scheduler.setKeepAliveSeconds(60); //允许空闲时间
        scheduler.initialize();
        return scheduler;
    }

需要异步调用的方法

@Component
public class AsyncMethod {

    //自定义线程池 ("executor")
    @Async("executor")
    public void asyncMethod(){
        System.out.println("asyncMethod--调用了");
        //抛出异常
        int a=1/0;
    }
    //自定义线程池 ("executor")
    @Async("executor")
    public Future<Integer> asyncSquare(int x) {
        System.out.println("calling asyncSquare," + Thread.currentThread().getName() + "," + new Date());
        try {
            Thread.sleep(2000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
        System.out.println("asyncSquare Finished," + Thread.currentThread().getName() + "," + new Date());
        return new AsyncResult<Integer>(x);
    }

    //自定义线程池 ("executor")
    @Async("executor")
    public void exeAsync(){
        System.out.println("--------------");
        int a=1/0;
    }
}

异常处理

@Logger
public class CusAsyncExceptionHandler extends SimpleAsyncUncaughtExceptionHandler {
    private static final Log log = LogFactory.getLog(CusAsyncExceptionHandler.class);
    @Override
    public void handleUncaughtException(Throwable throwable, Method method, Object... params) {
        StringBuilder info = new StringBuilder();
        String msg = throwable.getMessage() != null ? throwable.getMessage() : throwable.getClass().getSimpleName();
        info.append("出现异常:").append(msg).append(" methodName: "+method).append("\n");
        for (StackTraceElement stackTrace : throwable.getStackTrace()) {
            info.append(stackTrace.toString()).append("\n");
        }
        log.error(info.toString());
 
    }
}

 使用案例

import jdk.nashorn.internal.runtime.logging.Logger;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
 
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
 
@RestController
@Logger
public class AsyncController {
    private static final Log log = LogFactory.getLog(AsyncController.class);
    @Autowired
    private AsyncMethod asyncMethod;

    @GetMapping(value = "/async/method2")
    public Object method2() throws ExecutionException, InterruptedException {
        Future<Integer> result =  asyncMethod.asyncSquare(9);
        while (true){
            if(result.isCancelled()){
                break;
            }
            if(result.isDone()){
                break;
            }
        }

        return "asyncTestMethodAdd===="+result.get();
    }


    @GetMapping(value = "/async/method1")
    public String method1(){
        asyncMethod.exeAsync();
        return "str----";
    }



}

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐