查看原文
其他

Spring Boot使用@Async实现异步调用:ThreadPoolTaskScheduler线程池的优雅关闭

翟永超 程序猿DD 2019-07-13

上周发了一篇关于Spring Boot中使用 @Async来实现异步任务和线程池控制的文章:《Spring Boot使用@Async实现异步调用:自定义线程池》。由于最近身边也发现了不少异步任务没有正确处理而导致的问题,所以本文就接前面的内容,继续说说线程池的优雅关闭,主要针对 ThreadPoolTaskScheduler线程池。

问题现象

在上篇文章的例子中,我们定义了一个线程池,然后利用 @Async注解写了3个任务,并指定了这些任务执行使用的线程池。在上文的单元测试中,我们没有具体说说shutdown相关的问题,下面我们就来模拟一个问题现场出来。

第一步:如前文一样,我们定义一个 ThreadPoolTaskScheduler线程池:

  1. @SpringBootApplication

  2. public class Application {

  3.    public static void main(String[] args) {

  4.        SpringApplication.run(Application.class, args);

  5.    }

  6.    @EnableAsync

  7.    @Configuration

  8.    class TaskPoolConfig {

  9.        @Bean("taskExecutor")

  10.        public Executor taskExecutor() {

  11.            ThreadPoolTaskScheduler executor = new ThreadPoolTaskScheduler();

  12.            executor.setPoolSize(20);

  13.            executor.setThreadNamePrefix("taskExecutor-");

  14.            return executor;

  15.        }

  16.    }

  17. }

第二步:改造之前的异步任务,让它依赖一个外部资源,比如:Redis

  1. @Slf4j

  2. @Component

  3. public class Task {

  4.    @Autowired

  5.    private StringRedisTemplate stringRedisTemplate;

  6.    @Async("taskExecutor")

  7.    public void doTaskOne() throws Exception {

  8.        log.info("开始做任务一");

  9.        long start = System.currentTimeMillis();

  10.        log.info(stringRedisTemplate.randomKey());

  11.        long end = System.currentTimeMillis();

  12.        log.info("完成任务一,耗时:" + (end - start) + "毫秒");

  13.    }

  14.    @Async("taskExecutor")

  15.    public void doTaskTwo() throws Exception {

  16.        log.info("开始做任务二");

  17.        long start = System.currentTimeMillis();

  18.        log.info(stringRedisTemplate.randomKey());

  19.        long end = System.currentTimeMillis();

  20.        log.info("完成任务二,耗时:" + (end - start) + "毫秒");

  21.    }

  22.    @Async("taskExecutor")

  23.    public void doTaskThree() throws Exception {

  24.        log.info("开始做任务三");

  25.        long start = System.currentTimeMillis();

  26.        log.info(stringRedisTemplate.randomKey());

  27.        long end = System.currentTimeMillis();

  28.        log.info("完成任务三,耗时:" + (end - start) + "毫秒");

  29.    }

  30. }

注意:这里省略了pom.xml中引入依赖和配置redis的步骤

第三步:修改单元测试,模拟高并发情况下ShutDown的情况:

  1. @RunWith(SpringJUnit4ClassRunner.class)

  2. @SpringBootTest

  3. public class ApplicationTests {

  4.    @Autowired

  5.    private Task task;

  6.    @Test

  7.    @SneakyThrows

  8.    public void test() {

  9.        for (int i = 0; i < 10000; i++) {

  10.            task.doTaskOne();

  11.            task.doTaskTwo();

  12.            task.doTaskThree();

  13.            if (i == 9999) {

  14.                System.exit(0);

  15.            }

  16.        }

  17.    }

  18. }

说明:通过for循环往上面定义的线程池中提交任务,由于是异步执行,在执行过程中,利用 System.exit(0)来关闭程序,此时由于有任务在执行,就可以观察这些异步任务的销毁与Spring容器中其他资源的顺序是否安全。

第四步:运行上面的单元测试,我们将碰到下面的异常内容。

  1. org.springframework.data.redis.RedisConnectionFailureException: Cannot get Jedis connection; nested exception is redis.clients.jedis.exceptions.JedisConnectionException: Could not get a resource from the pool

  2.    at org.springframework.data.redis.connection.jedis.JedisConnectionFactory.fetchJedisConnector(JedisConnectionFactory.java:204) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  3.    at org.springframework.data.redis.connection.jedis.JedisConnectionFactory.getConnection(JedisConnectionFactory.java:348) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  4.    at org.springframework.data.redis.core.RedisConnectionUtils.doGetConnection(RedisConnectionUtils.java:129) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  5.    at org.springframework.data.redis.core.RedisConnectionUtils.getConnection(RedisConnectionUtils.java:92) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  6.    at org.springframework.data.redis.core.RedisConnectionUtils.getConnection(RedisConnectionUtils.java:79) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  7.    at org.springframework.data.redis.core.RedisTemplate.execute(RedisTemplate.java:194) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  8.    at org.springframework.data.redis.core.RedisTemplate.execute(RedisTemplate.java:169) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  9.    at org.springframework.data.redis.core.RedisTemplate.randomKey(RedisTemplate.java:781) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  10.    at com.didispace.async.Task.doTaskOne(Task.java:26) ~[classes/:na]

  11.    at com.didispace.async.Task$$FastClassBySpringCGLIB$$ca3ff9d6.invoke(<generated>) ~[classes/:na]

  12.    at org.springframework.cglib.proxy.MethodProxy.invoke(MethodProxy.java:204) ~[spring-core-4.3.14.RELEASE.jar:4.3.14.RELEASE]

  13.    at org.springframework.aop.framework.CglibAopProxy$CglibMethodInvocation.invokeJoinpoint(CglibAopProxy.java:738) ~[spring-aop-4.3.14.RELEASE.jar:4.3.14.RELEASE]

  14.    at org.springframework.aop.framework.ReflectiveMethodInvocation.proceed(ReflectiveMethodInvocation.java:157) ~[spring-aop-4.3.14.RELEASE.jar:4.3.14.RELEASE]

  15.    at org.springframework.aop.interceptor.AsyncExecutionInterceptor$1.call(AsyncExecutionInterceptor.java:115) ~[spring-aop-4.3.14.RELEASE.jar:4.3.14.RELEASE]

  16.    at java.util.concurrent.FutureTask.run(FutureTask.java:266) [na:1.8.0_151]

  17.    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180) [na:1.8.0_151]

  18.    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293) [na:1.8.0_151]

  19.    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [na:1.8.0_151]

  20.    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [na:1.8.0_151]

  21.    at java.lang.Thread.run(Thread.java:748) [na:1.8.0_151]

  22. Caused by: redis.clients.jedis.exceptions.JedisConnectionException: Could not get a resource from the pool

  23.    at redis.clients.util.Pool.getResource(Pool.java:53) ~[jedis-2.9.0.jar:na]

  24.    at redis.clients.jedis.JedisPool.getResource(JedisPool.java:226) ~[jedis-2.9.0.jar:na]

  25.    at redis.clients.jedis.JedisPool.getResource(JedisPool.java:16) ~[jedis-2.9.0.jar:na]

  26.    at org.springframework.data.redis.connection.jedis.JedisConnectionFactory.fetchJedisConnector(JedisConnectionFactory.java:194) ~[spring-data-redis-1.8.10.RELEASE.jar:na]

  27.    ... 19 common frames omitted

  28. Caused by: java.lang.InterruptedException: null

  29.    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.reportInterruptAfterWait(AbstractQueuedSynchronizer.java:2014) ~[na:1.8.0_151]

  30.    at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.awaitNanos(AbstractQueuedSynchronizer.java:2088) ~[na:1.8.0_151]

  31.    at org.apache.commons.pool2.impl.LinkedBlockingDeque.pollFirst(LinkedBlockingDeque.java:635) ~[commons-pool2-2.4.3.jar:2.4.3]

  32.    at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:442) ~[commons-pool2-2.4.3.jar:2.4.3]

  33.    at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:361) ~[commons-pool2-2.4.3.jar:2.4.3]

  34.    at redis.clients.util.Pool.getResource(Pool.java:49) ~[jedis-2.9.0.jar:na]

  35.    ... 22 common frames omitted

如何解决

原因分析

从异常信息 JedisConnectionException:Couldnotgeta resourcefromthe pool来看,我们很容易的可以想到,在应用关闭的时候异步任务还在执行,由于Redis连接池先销毁了,导致异步任务中要访问Redis的操作就报了上面的错。所以,我们得出结论,上面的实现方式在应用关闭的时候是不优雅的,那么我们要怎么做呢?

解决方法

要解决上面的问题很简单,Spring的 ThreadPoolTaskScheduler为我们提供了相关的配置,只需要加入如下设置即可:

  1. @Bean("taskExecutor")

  2. public Executor taskExecutor() {

  3.    ThreadPoolTaskScheduler executor = new ThreadPoolTaskScheduler();

  4.    executor.setPoolSize(20);

  5.    executor.setThreadNamePrefix("taskExecutor-");

  6.    executor.setWaitForTasksToCompleteOnShutdown(true);

  7.    executor.setAwaitTerminationSeconds(60);

  8.    return executor;

  9. }

说明: setWaitForTasksToCompleteOnShutdowntrue该方法就是这里的关键,用来设置线程池关闭的时候等待所有任务都完成再继续销毁其他的Bean,这样这些异步任务的销毁就会先于Redis线程池的销毁。同时,这里还设置了 setAwaitTerminationSeconds(60),该方法用来设置线程池中任务的等待时间,如果超过这个时候还没有销毁就强制销毁,以确保应用最后能够被关闭,而不是阻塞住。

完整示例:

读者可以根据喜好选择下面的两个仓库中查看 Chapter4-1-4项目:

  • Github:https://github.com/dyc87112/SpringBoot-Learning/

  • Gitee:https://gitee.com/didispace/SpringBoot-Learning/

如果您对这些感兴趣,欢迎star、follow、收藏、转发给予支持!

阿里云1C2G虚拟机【99/年】羊毛党集合啦!

推荐阅读


长按指纹

一键关注

深入交流、更多福利

扫码加入我的知识星球



点击 “阅读原文” 看看本号其他精彩内容

    您可能也对以下帖子感兴趣

    文章有问题?点此查看未经处理的缓存