Home | 简体中文 | 繁体中文 | 杂文 | Github | 知乎专栏 | Facebook | Linkedin | Youtube | 打赏(Donations) | About
知乎专栏

1.35. Spring boot with Async

Callable 实现异步

		
@GetMapping("/email")
public Callable<String> order() {
    System.out.println("主线程开始:" + Thread.currentThread().getName());
    Callable<String> result = () -> {
        System.out.println("副线程开始:" + Thread.currentThread().getName());
        Thread.sleep(1000);
        System.out.println("副线程返回:" + Thread.currentThread().getName());
        return "success";
    };

    System.out.println("主线程返回:" + Thread.currentThread().getName());
    return result;
}		
		
		

WebAsyncTask 实现异步

		
@GetMapping("/webAsyncTask")
public WebAsyncTask<String> webAsyncTask() {
    log.info("外部线程:" + Thread.currentThread().getName());
    WebAsyncTask<String> result = new WebAsyncTask<>(60 * 1000L, new Callable<String>() {
        @Override
        public String call() {
            log.info("内部线程:" + Thread.currentThread().getName());
            return "success";
        }
    });
    result.onTimeout(new Callable<String>() {
        @Override
        public String call() {
            log.info("timeout callback");
            return "timeout callback";
        }
    });
    result.onCompletion(new Runnable() {
        @Override
        public void run() {
            log.info("finish callback");
        }
    });
    return result;
}
		
		

DeferredResult 实现异步返回结果

		
	private DeferredResult<String> deferredResult = new DeferredResult<String>();
		
	@ResponseBody
    @GetMapping("/receive")
    public DeferredResult<String> receive() throws Exception {
        return deferredResult;
    }

    @ResponseBody
    @GetMapping("/send")
    public void send() throws Exception {
        deferredResult.setResult("Helloworld!!!");
    }
		
		
		
	private final List<DeferredResult<String>> deferredResultList = new ArrayList<DeferredResult<String>>();
	
    @ResponseBody
    @GetMapping("/receive")
    public DeferredResult<String> receive() throws Exception {
        DeferredResult<String> deferredResult = new DeferredResult<>();

        //先存起来,等待触发
        deferredResultList.add(deferredResult);
        return deferredResult;
    }

    @ResponseBody
    @GetMapping("/send")
    public void send() throws Exception {
        // 让所有hold住的请求给与响应
        deferredResultList.forEach(d -> d.setResult("say hello to all"));
    }		
		
		

DeferredResult 与 Callback 配合使用,用来获取 Callback 返回值

		
    @GetMapping("/tts")
    @Operation(summary = "音频合成")
    @ResponseBody
    public DeferredResult<ResponseJson> test(@RequestParam("text") String text, @RequestParam("filename") String filename) {
        DeferredResult<ResponseJson> deferredResult = new DeferredResult<ResponseJson>();
        speechSynthesizerService.tts(text, new XfyunCallback() {
            @Override
            public void onCallback(String sid, String text) {
                String audio = aliyunService.uploadMp3FromBase64(text, filename.concat(".mp3"));
                ResponseJson response = new ResponseJson(true, ResponseJson.Code.SUCCESS, "", audio);
                deferredResult.setResult(response);
            }
        });
        return deferredResult;
    }		
		
		

1.35.1. 带有返回值的异步任务

Future

        
@Service
public class DemoAsyncServiceImpl implements DemoAsyncService {

	public static Random random =new Random();
	
	@Async("OneAsync")
	public Future<String> doTaskOne() throws Exception {
	    System.out.println("开始做任务一");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务一,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务一完成");
	}
	
	@Async("TwoAsync")
	public Future<String> doTaskTwo() throws Exception {
	    System.out.println("开始做任务二");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务二,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务二完成");
	}
	
	@Async
	public Future<String> doTaskThree() throws Exception {
	    System.out.println("开始做任务三");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务三,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务三完成");
	}

}
        
			
			
			
			
			

CompletableFuture

			
        CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> {
            throw new RuntimeException();
//            return "程序出现异常";
        }).exceptionally((e) -> {
            System.out.println("程序出现异常");
            return "程序出现异常";
        });
        System.out.println(completableFuture.get());			
			
			
			
@Service
public class MyService {

    @Async
    public CompletableFuture<String> asyncMethod() {
        try {
            // 异步方法逻辑...
            return CompletableFuture.completedFuture("Success");
        } catch (Exception e) {
            // 处理异常...
            return CompletableFuture.failedFuture(e);
        }
    }
}

// 调用异步方法并处理异常
CompletableFuture<String> future = myService.asyncMethod();
future.exceptionally(e -> {
    // 异常处理逻辑...
    return "Error";
});

			
			
			
    @GetMapping("/completableFutureExceptionally")
    public String completableFutureExceptionally() throws ExecutionException, InterruptedException {

        CompletableFuture.supplyAsync(() -> {
            System.out.println("当前线程名称:" + Thread.currentThread().getName());
            throw new RuntimeException();
        }).exceptionally((e) -> {
            System.out.println(e.getMessage());
            return "程序出现异常";
        });

        return "Done";
    }			
			
			

输出结果

			
当前线程名称:ForkJoinPool.commonPool-worker-1
java.lang.RuntimeException			
			
			

1.35.2. 默认简单线程池 SimpleAsyncTaskExecutor

启用异步执行 @EnableAsync

		
@EnableAsync
@SpringBootApplication
public class ThreadPoolApplication {

    public static void main(String[] args) {
        SpringApplication.run(ThreadPoolApplication.class, args);
    }

}		
		
		

编写异步执行代码

		
@Component
@Slf4j
public class AsyncTask {
    @Async
    public void  asyncRun() throws InterruptedException {
        Thread.sleep(10);
        log.info(Thread.currentThread().getName()+":处理完成");
    }
}		
		
		

配置线程池

默认线程池的配置很简单,配置参数如下:

		
spring.task.execution.pool.core-size:线程池创建时的初始化线程数,默认为8
spring.task.execution.pool.max-size:线程池的最大线程数,默认为int最大值
spring.task.execution.pool.queue-capacity:用来缓冲执行任务的队列,默认为int最大值
spring.task.execution.pool.keep-alive:线程终止前允许保持空闲的时间
spring.task.execution.pool.allow-core-thread-timeout:是否允许核心线程超时
spring.task.execution.shutdown.await-termination:是否等待剩余任务完成后才关闭应用
spring.task.execution.shutdown.await-termination-period:等待剩余任务完成的最大时间
spring.task.execution.thread-name-prefix:线程名的前缀,设置好了之后可以方便我们在日志中查看处理任务所在的线程池		
		
			

具体配置含义如下:

		
spring.task.execution.pool.core-size=8
spring.task.execution.pool.max-size=20
spring.task.execution.pool.queue-capacity=10
spring.task.execution.pool.keep-alive=60s
spring.task.execution.pool.allow-core-thread-timeout=true
spring.task.execution.shutdown.await-termination=true
spring.task.execution.shutdown.await-termination-period=60
spring.task.execution.thread-name-prefix=simple-
		
			
		
spring:
  task:
    execution:
      thread-name-prefix: task- # 线程池的线程名的前缀。默认为 task- ,建议根据自己应用来设置
      pool: # 线程池相关
        core-size: 8 # 核心线程数,线程池创建时候初始化的线程数。默认为 8 。
        max-size: 20 # 最大线程数,线程池最大的线程数,只有在缓冲队列满了之后,才会申请超过核心线程数的线程。默认为 Integer.MAX_VALUE
        keep-alive: 60s # 允许线程的空闲时间,当超过了核心线程之外的线程,在空闲时间到达之后会被销毁。默认为 60 秒
        queue-capacity: 200 # 缓冲队列大小,用来缓冲执行任务的队列的大小。默认为 Integer.MAX_VALUE 。
        allow-core-thread-timeout: true # 是否允许核心线程超时,即开启线程池的动态增长和缩小。默认为 true 。
      shutdown:
        await-termination: true # 应用关闭时,是否等待定时任务执行完成。默认为 false ,建议设置为 true
        await-termination-period: 60 # 等待任务完成的最大时长,单位为秒。默认为 0 ,根据自己应用来设置		
		
			

@Service/@Component 中异步执行

直接使用 @Async 注解,即可调用默认线程池

        
@Component
public class Task {

	@Async
	public void doTaskOne() throws Exception {
	    // 业务逻辑
	}
	
	@Async
	public void doTaskTwo() throws Exception {
	    // 业务逻辑
	}
	
	@Async
	public void doTaskThree() throws Exception {
	    // 业务逻辑
	}

}			
        
			

this 调用方法 @Async 将失效,编程同步阻塞执行。例如下面的 doTaskThree() 会同步执行 this.doTaskOne(); 和 this.doTaskTwo();

			
@Component
public class Task {

	@Async
	public void doTaskOne() throws Exception {
	    // 业务逻辑
	}
	
	@Async
	public void doTaskTwo() throws Exception {
	    // 业务逻辑
	}
	
	public void doTaskThree() throws Exception {
	    this.doTaskOne();
	    this.doTaskTwo();
	}

}						
			
			

applicationTaskExecutor

SimpleAsyncTaskExecutor Bean 名字是 applicationTaskExecutor

			
	"applicationTaskExecutor": {
          "aliases": [
            "taskExecutor"
          ],
          "scope": "singleton",
          "type": "org.springframework.core.task.SimpleAsyncTaskExecutor",
          "resource": "class path resource [org/springframework/boot/autoconfigure/task/TaskExecutorConfigurations$TaskExecutorConfiguration.class]",
          "dependencies": [
            "org.springframework.boot.autoconfigure.task.TaskExecutorConfigurations$TaskExecutorConfiguration",
            "simpleAsyncTaskExecutorBuilder"
          ]
        }			
			
			
			
package cn.netkiller.service;

import lombok.extern.slf4j.Slf4j;
import org.springframework.cache.annotation.Cacheable;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;

@Service
@Slf4j
public class TestService {

    @Async("applicationTaskExecutor")
    public void thread() {
        log.info("applicationTaskExecutor");
    }
}

			
			
			
2024-02-02T18:20:50.938+08:00  INFO 91307 --- [watch-development] [       simple-2] cn.netkiller.service.TestService           : applicationTaskExecutor
			
			

1.35.3. ThreadPoolTaskExecutor 自定义线程池

Springboot 官方不建议在生产环境使用 SimpleAsyncTaskExecutor,建议使用 自定义线程池。

		
2024-02-02T16:24:38.585+08:00  WARN 86223 --- [watch-development] [nio-8080-exec-2] s.w.s.m.m.a.RequestMappingHandlerAdapter : !!!
Performing asynchronous handling through the default Spring MVC SimpleAsyncTaskExecutor.
This executor is not suitable for production use under load.
Please, configure an AsyncTaskExecutor through the WebMvc config.
-------------------------------
!!!		
		
		

Bean 注入配置线程池

        
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

import java.util.concurrent.Executor;

@SpringBootApplication
@EnableAsync
public class Application {

	public static void main(String[] args) {
	    // close the application context to shut down the custom ExecutorService
	    SpringApplication.run(Application.class, args).close();
	}
	
	@Bean
	public Executor asyncExecutor() {
	    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
	    executor.setCorePoolSize(2);
	    executor.setMaxPoolSize(2);
	    executor.setQueueCapacity(500);
	    executor.setThreadNamePrefix("Netkiller -");
	    executor.initialize();
	    return executor;
	}

	@Bean("thread")
    public Executor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        // 设置核心线程数
        executor.setCorePoolSize(5);
        // 设置最大线程数
        executor.setMaxPoolSize(10);
        // 设置队列容量
        executor.setQueueCapacity(20);
        // 设置线程活跃时间(秒)
        executor.setKeepAliveSeconds(60);
        // 设置线程名称
        executor.setThreadNamePrefix("hello-");
        // 设置拒绝策略
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        // 等待所有任务结束后再关闭线程池
        executor.setWaitForTasksToCompleteOnShutdown(true);

        return executor;
    }

}			
        
		

设置线程池参数

最简单的配置

        
@SpringBootApplication
@EnableAsync
public class Application {
public static void main(String[] args) {
    SpringApplication.run(Application.class, args);
}
}
        
			

        
@Component
public class Task {

	@Async
	public void doTaskOne() throws Exception {
	    // 业务逻辑
	}
	
	@Async("asyncExecutor")
	public void doTaskTwo() throws Exception {
	    // 业务逻辑
	}
	
	@Async("thread")
	public void doTaskThree() throws Exception {
	    // 业务逻辑
	}

}			
        
			

注意:@Async 不会用到刚刚定义的线程池,@Async("asyncExecutor"),@Async("thread") 会正确调用

队列

线程池能接受多少队列?

下面配置是 executor.setQueueCapacity(10); 也就是 10个,但是实测结果跟你想的不同

			
package cn.netkiller.config;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableAsync;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;

import java.util.concurrent.ThreadPoolExecutor;

@Configuration
@EnableAsync
public class ThreadPoolTaskExecutorConfiguration {
    @Bean("asyncExecutor")
    public ThreadPoolTaskExecutor executor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setThreadGroupName("job");
        executor.setThreadNamePrefix("async-job-");
        executor.setCorePoolSize(5);
        executor.setMaxPoolSize(10);
        executor.setQueueCapacity(10);
        executor.setKeepAliveSeconds(60);
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
        executor.setAwaitTerminationSeconds(60);
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.initialize();
        return executor;
    }
}
			
			

实测结果是,首次执行可以容纳 20 个线程,20个线程执行完毕之后,再添加任务,就只接受 10 个,超过的部分会跑出异常

			
 Executor [java.util.concurrent.ThreadPoolExecutor@7e729046[Running, pool size = 10, active threads = 10, queued tasks = 10, completed tasks = 0]] did not accept task: org.springframework.aop.interceptor.AsyncExecutionInterceptor$$Lambda$1775/0x0000000801b6afb0@20eaccc2			
			
			

这是因为线程池可以容纳 10 个任务,队列还能排队 10 个任务。

			
			
			
			

定义多个线程池

        
package cn.netkiller.wallet.config;

import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.EnableAsync;

@Configuration
@EnableAsync
public class ExecutorConfiguration {
	/** Set the ThreadPoolExecutor's core pool size. */
	private int corePoolSize = 10;
	/** Set the ThreadPoolExecutor's maximum pool size. */
	private int maxPoolSize = 200;
	/** Set the capacity for the ThreadPoolExecutor's BlockingQueue. */
	private int queueCapacity = 10;
	
	@Bean
	public Executor OneAsync() {
	    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
	    executor.setCorePoolSize(corePoolSize);
	    executor.setMaxPoolSize(maxPoolSize);
	    executor.setQueueCapacity(queueCapacity);
	    executor.setThreadNamePrefix("MySimpleExecutor-");
	    executor.initialize();
	    return executor;
	}
	
	@Bean
	public Executor TwoAsync() {
	    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
	    executor.setCorePoolSize(corePoolSize);
	    executor.setMaxPoolSize(maxPoolSize);
	    executor.setQueueCapacity(queueCapacity);
	    executor.setThreadNamePrefix("MyExecutor-");
	
	    // rejection-policy:当pool已经达到max size的时候,如何处理新任务
	    // CALLER_RUNS:不在新线程中执行任务,而是有调用者所在的线程来执行
	    executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
	    executor.initialize();
	    return executor;
	}

}
        
        
			
        
@Service
public class DemoAsyncServiceImpl implements DemoAsyncService {

	public static Random random =new Random();
	
	@Async("OneAsync")
	public Future<String> doTaskOne() throws Exception {
	    System.out.println("开始做任务一");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务一,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务一完成");
	}
	
	@Async("TwoAsync")
	public Future<String> doTaskTwo() throws Exception {
	    System.out.println("开始做任务二");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务二,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务二完成");
	}
	
	@Async
	public Future<String> doTaskThree() throws Exception {
	    System.out.println("开始做任务三");
	    long start = System.currentTimeMillis();
	    Thread.sleep(random.nextInt(10000));
	    long end = System.currentTimeMillis();
	    System.out.println("完成任务三,耗时:" + (end - start) + "毫秒");
	    return new AsyncResult<>("任务三完成");
	}

}
        
			

实现 AsyncConfigurer 接口方式创建自定义连接池

这种方式可以直接使用 @Async

			
package cn.netkiller.config;

import org.springframework.beans.factory.annotation.Value;
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 ThreadPoolTaskExecutorConfiguration implements AsyncConfigurer {

    @Value("${spring.task.execution.pool.core-size}")
    private int corePoolSize;
    @Value("${spring.task.execution.pool.max-size}")
    private int maxPoolSize;
    @Value("${spring.task.execution.pool.queue-capacity}")
    private int queueCapacity;
    //    @Value("${spring.task.execution.thread-name-prefix}")
    private final String threadNamePrefix = "async-";
    @Value("${spring.task.execution.pool.keep-alive}")
    private int keepAliveSeconds;

    @Override
    public Executor getAsyncExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        //最大线程数
        executor.setMaxPoolSize(maxPoolSize);
        //核心线程数
        executor.setCorePoolSize(corePoolSize);
        //任务队列的大小
        executor.setQueueCapacity(queueCapacity);
        //线程前缀名
        executor.setThreadNamePrefix(threadNamePrefix);
        //线程存活时间
        executor.setKeepAliveSeconds(keepAliveSeconds);

        /**
         * 拒绝处理策略
         * CallerRunsPolicy():交由调用方线程运行,比如 main 线程。
         * AbortPolicy():直接抛出异常。
         * DiscardPolicy():直接丢弃。
         * DiscardOldestPolicy():丢弃队列中最老的任务。
         */
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        //线程初始化
        executor.initialize();
        return executor;
    }
}
    
			
			

继承 AsyncConfigurerSupport 创建自定义连接池

			
package cn.netkiller.config;

import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.scheduling.annotation.AsyncConfigurerSupport;
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 ThreadPoolTaskExecutorConfiguration extends AsyncConfigurerSupport {

    @Value("${spring.task.execution.pool.core-size}")
    private int corePoolSize;
    @Value("${spring.task.execution.pool.max-size}")
    private int maxPoolSize;
    @Value("${spring.task.execution.pool.queue-capacity}")
    private int queueCapacity;
    //    @Value("${spring.task.execution.thread-name-prefix}")
    private final String threadNamePrefix = "async-";
    @Value("${spring.task.execution.pool.keep-alive}")
    private int keepAliveSeconds;

    @Override
    public Executor getAsyncExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        //最大线程数
        executor.setMaxPoolSize(maxPoolSize);
        //核心线程数
        executor.setCorePoolSize(corePoolSize);
        //任务队列的大小
        executor.setQueueCapacity(queueCapacity);
        //线程前缀名
        executor.setThreadNamePrefix(threadNamePrefix);
        //线程存活时间
        executor.setKeepAliveSeconds(keepAliveSeconds);

        /**
         * 拒绝处理策略
         * CallerRunsPolicy():交由调用方线程运行,比如 main 线程。
         * AbortPolicy():直接抛出异常。
         * DiscardPolicy():直接丢弃。
         * DiscardOldestPolicy():丢弃队列中最老的任务。
         */
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        //线程初始化
        executor.initialize();
        return executor;
    }
}			
			
			

生产环境完整代码 @Bean 注入方式

这种方式多用在 Spring 2.x 中

			
package cn.netkiller.config;

import org.springframework.beans.factory.annotation.Value;
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 ThreadPoolTaskExecutorConfiguration {

    @Value("${spring.task.execution.pool.core-size}")
    private int corePoolSize;
    @Value("${spring.task.execution.pool.max-size}")
    private int maxPoolSize;
    @Value("${spring.task.execution.pool.queue-capacity}")
    private int queueCapacity;
    //    @Value("${spring.task.execution.thread-name-prefix}")
    private final String threadNamePrefix = "";
    @Value("${spring.task.execution.pool.keep-alive}")
    private int keepAliveSeconds;
    
    @Bean
    public Executor threadPoolTaskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        //最大线程数
        executor.setMaxPoolSize(maxPoolSize);
        //核心线程数
        executor.setCorePoolSize(corePoolSize);
        //任务队列的大小
        executor.setQueueCapacity(queueCapacity);
        //线程前缀名
        executor.setThreadNamePrefix(threadNamePrefix);
        //线程存活时间
        executor.setKeepAliveSeconds(keepAliveSeconds);

        /**
         * 拒绝处理策略
         * CallerRunsPolicy():交由调用方线程运行,比如 main 线程。
         * AbortPolicy():直接抛出异常。
         * DiscardPolicy():直接丢弃。
         * DiscardOldestPolicy():丢弃队列中最老的任务。
         */
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        //线程初始化
        executor.initialize();
        return executor;
    }
}			
			
			

注意:使用@Bean方式必须配合 @Async("threadPoolTaskExecutor")

通过 @Bean 覆盖掉 SimpleAsyncTaskExecutor

自定义连接池之后,系统内会存在两个连接吃 SimpleAsyncTaskExecutor 和 ThreadPoolTaskExecutor

			
   @Bean
    public Executor applicationTaskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(2);
        executor.setMaxPoolSize(2);
        executor.setQueueCapacity(500);
        executor.setThreadNamePrefix("Netkiller-");
        executor.initialize();
        return executor;
    }			
			
			

1.35.4. 自定义线程池 ThreadPoolExecutor

自定义线程池

ThreadPoolExecutor

			
    @Bean("queueThreadPool")
    public ThreadPoolExecutor queueThreadPool() {
        ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(
                5,
                10,
                60,
                TimeUnit.SECONDS,
                new LinkedBlockingDeque<>(10),
                Executors.defaultThreadFactory(),
                new ThreadPoolExecutor.AbortPolicy());
        return threadPoolExecutor;
    }
			
			

注入自定义线程池bean

		
	 // 注入自定义线程池bean
    @Autowired
    private ThreadPoolExecutor threadPoolExecutor;
    
    threadPoolExecutor.execute(new Runnable() {
        @Override
        public void run() {
            System.out.println("=======");

        }
    });
    
		
			

设置线程名称

		
public static Random random = new Random();

    @Async("jobExecutor")
    public void doAsyncTask() throws InterruptedException {
        Thread.currentThread().setName("测试线程-" + random.nextInt(1000));
        System.out.println("开始做任务一");
        long start = System.currentTimeMillis();
        Thread.sleep(random.nextInt(100000));
        long end = System.currentTimeMillis();
        System.out.println("完成任务一,耗时:" + (end - start) + "毫秒");
    }
		
		

线程池监控

监控指标

		
neo@MacBook-Pro-M2 ~> curl -s http://www.netkiller.cn:8080/actuator/metrics | jq | grep executor
    "executor.active",
    "executor.completed",
    "executor.pool.core",
    "executor.pool.max",
    "executor.pool.size",
    "executor.queue.remaining",
    "executor.queued",		
		
		

获取指标

		
neo@MacBook-Pro-M2 ~> curl -s http://www.netkiller.cn:8080/actuator/metrics/executor.active | jq
{
  "name": "executor.active",
  "description": "The approximate number of threads that are actively executing tasks",
  "baseUnit": "threads",
  "measurements": [
    {
      "statistic": "VALUE",
      "value": 0
    }
  ],
  "availableTags": [
    {
      "tag": "name",
      "values": [
        "asyncExecutor"
      ]
    }
  ]
}		
		
		

		
    @Autowired
    ThreadPoolTaskExecutor threadPoolTaskExecutor;
    		
    @GetMapping("/pool")
    public String pool() {
        int activeCount = threadPoolTaskExecutor.getActiveCount();
        long completedTaskCount = threadPoolTaskExecutor.getThreadPoolExecutor().getCompletedTaskCount();
        long taskCount = threadPoolTaskExecutor.getThreadPoolExecutor().getTaskCount();
        int queue = threadPoolTaskExecutor.getThreadPoolExecutor().getQueue().size();
        String monitor = String.format("Task: %d, Queue: %d, Active: %d, Completed: %d\n", taskCount, queue, activeCount, completedTaskCount);
        log.info(monitor);
        return monitor;
    }	
		
		

注意事项

再 Service 中 通过 this 访问另一个 @Async 方法,异步会失效,必须使用 @@Autowired 入住。