使用线程池

线程池的创建

创建线程池直接new ThreadPoolExecutor这个类就好了

public ThreadPoolExecutor(int corePoolSize,
                              int maximumPoolSize,
                              long keepAliveTime,
                              TimeUnit unit,
                              BlockingQueue<Runnable> workQueue,
                              ThreadFactory threadFactory,
                              RejectedExecutionHandler handler)

构造方法有几个重载,这个事参数最全的,如下分别介绍这些参数

  • corePoolSize: 核心线程的数量,当提交一个任务到线程池时,当线程数量小于corePoolSize时,进来新任务都会创建新的线程。如果调用了线程池的prestartAllCoreThreads()方法,线程池会提前创建并启动所有核心线程,而不是等到任务进来
  • maximumPoolSize:最大线程数量,如果队列满了,已有线程数小于最大线程数,则线程池会再创建新的线程执行任务。值得注意的是,如果使用了无界的任务队列这个参数就没什么效果。
  • keepAliveTime: 线程池的工作线程空闲后,保持存活的时间。所以,如果任务很多,并且每个任务执行的时间比较短,可以调大时间,提高线程的利用率。
  • unit: 线程保活时间单位
  • workQueue: 用于保存等待执行的任务的阻塞队列。可以选择以下几个阻塞队列。ArrayBlockingQueue,LinkedBlockingQueue,SynchronousQueue,PriorityBlockingQueue
  • threadFactory:用于设置创建线程的工厂,可以通过线程工厂给每个创建出来的线程设置更有意义的名字。要实现ThreadFactory接口
  • handler:拒绝策略,当队列和线程池都满了,说明线程池处于饱和状态,那么必须采取一种策略处理提交的新任务。这个策略默认情况下是AbortPolicy,表示无法处理新任务时抛出异常。jdk有几种默认策略,也可以自定义策略,默认策略如下
    • AbortPolicy:直接抛出异常。
    • CallerRunsPolicy:只用调用者所在线程来运行任务。
    • DiscardOldestPolicy:丢弃队列里最近的一个任务,并执行当前任务。
    • DiscardPolicy:不处理,丢弃掉。

这一段就是从书上借鉴(抄书不叫抄,叫借鉴)过来的,来个demo可能会更加清楚一些

ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(2, 2, 10, TimeUnit.SECONDS, new ArrayBlockingQueue<>(3), new ThreadFactory() {
    @Override
    public Thread newThread(Runnable r) {
        Thread thread = new Thread(r);
        thread.setName("欢迎关注微信公众号“大雄和你一起学编程”" + atomicInteger.getAndIncrement());
        return thread;
    }
}, new RejectedExecutionHandler() {
    @Override
    public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
        System.out.println("任务被拒绝了哦,欢迎关注微信公众号“大雄和你一起学编程”");
    }
});

这个例子里边,核心线程数2,最大线程数2, 保活时间10秒,工作队列长度为3的ArrayBlokingQueue, 线程工厂里给线程起了一个好听的名字 “欢迎关注微信公众号“大雄和你一起学编程”数字”,拒绝测试时打印一段非常有意思的一段话 "任务被拒绝了哦,欢迎关注微信公众号“大雄和你一起学编程”"

向线程池提交任务

方式1:excute

execute()方法用于提交不需要返回值的任务,所以无法判断任务是否被线程池执行成功。

接着上边的例子

for (int i = 0; i < 6; i++) {
    // 我们提交6个任务,按照上边的配置,应该会拒绝一个
    int finalI = i;
    threadPoolExecutor.execute(() -> {
        try {
            Thread.sleep(1000);
            System.out.println("任务" + finalI + "执行完了,欢迎关注大雄和你一起学编程");
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    });
}
// 运行的结果如下

// 任务被拒绝了哦,欢迎关注微信公众号“大雄和你一起学编程”
// 任务1执行完了,欢迎关注大雄和你一起学编程
// 任务0执行完了,欢迎关注大雄和你一起学编程
// 任务2执行完了,欢迎关注大雄和你一起学编程
// 任务3执行完了,欢迎关注大雄和你一起学编程
// 任务4执行完了,欢迎关注大雄和你一起学编程

方式2:submit

  • submit()方法用于提交需要返回值的任务

  • 线程池会返回一个future类型的对象,通过这个future对象可以判断任务是否执行成功,并且可以通过future的get()方法来获取返回值,get()方法会阻塞当前线程直到任务完成

  • 使用get(longtimeout,TimeUnit unit)方法则会阻塞当前线程一段时间后立即返回,这时候有可能任务没有执行完

既然是由返回值,那么任务就要用Callable了

for (int i = 0; i < 6; i++) {
    int finalI = i;
    Future<Integer> future = threadPoolExecutor.submit(new Callable<Integer>() {
        @Override
        public Integer call() throws Exception {
            int sum = 0;
            for (int j = 0; j < finalI * 10; j++) {
                sum += j;
                if (j == 39) {
                    sum = j / 0;
                }
            }
            return sum;
        }
    });
    try {
        Integer integer = null;
        integer = future.get();
        System.out.println("任务执行结果:" + integer);
    } catch (InterruptedException e) {
        e.printStackTrace();
    } catch (ExecutionException e) {
        e.printStackTrace();
    }
}
// 输出如下:
// 任务执行结果:0
// 任务执行结果:45
// 任务执行结果:190
// 任务执行结果:435
// java.util.concurrent.ExecutionException: java.lang.ArithmeticException: / by zero
//     at java.util.concurrent.FutureTask.report(FutureTask.java:122)
//     at java.util.concurrent.FutureTask.get(FutureTask.java:192)
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll.main(Demo_06_01_1_ThreadPoll.java:54)
// Caused by: java.lang.ArithmeticException: / by zero
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll$3.call(Demo_06_01_1_ThreadPoll.java:46)
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll$3.call(Demo_06_01_1_ThreadPoll.java:39)
//     at java.util.concurrent.FutureTask.run(FutureTask.java:266)
//     at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
//     at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
//     at java.lang.Thread.run(Thread.java:748)
// java.util.concurrent.ExecutionException: java.lang.ArithmeticException: / by zero
//     at java.util.concurrent.FutureTask.report(FutureTask.java:122)
//     at java.util.concurrent.FutureTask.get(FutureTask.java:192)
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll.main(Demo_06_01_1_ThreadPoll.java:54)
// Caused by: java.lang.ArithmeticException: / by zero
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll$3.call(Demo_06_01_1_ThreadPoll.java:46)
//     at com.xiong.concurrent.threadpoll.Demo_06_01_1_ThreadPoll$3.call(Demo_06_01_1_ThreadPoll.java:39)
//     at java.util.concurrent.FutureTask.run(FutureTask.java:266)
//     at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
//     at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
//     at java.lang.Thread.run(Thread.java:748)

关闭线程池

调用关闭线程池方法后

遍历线程池中的工作线程,然后逐个调用线程的interrupt方法来中断线程,所以无法响应中断的任务可能永远无法终止

shutdown

shutdown只是将线程池的状态设置成SHUTDOWN状态,然后中断所有没有正在执行任务的线程

shutdownNow

shutdownNow首先将线程池的状态设置成STOP,然后尝试停止所有的正在执行或暂停任务的线程,并返回等待执行任务的列表


调用任意一个isShutdown方法就会返回true

当所有的任务都已关闭后,才表示线程池关闭成功,这时调用isTerminaed方法会返回true

创建线程池的参数怎么给

前面的例子中是拍脑袋给的参数,说要2个核心线程啥的,但是应该如何正确决策呢

cpu核心数通过Runtime.getRuntime().availableProcessors()获得

按任务类型考虑

cpu密集型:核心线程数=核心数+1

io密集型:核心线程数=2*核心数

混合型的任务:如果可以拆分,将其拆分成一个CPU密集型任务和一个IO密集型任务,只要这两个任务执行的时间相差不是太大,那么分解后执行的吞吐量将高于串行执行的吞吐量

按优先级考虑

优先级不同的任务可以使用优先级队列PriorityBlockingQueue来处理。它可以让优先级高的任务先执行。

一些原则

  • 执行时间不同的任务可以交给不同规模线程池处理

  • 依赖数据库连接池的任务:线程数可以适当给多一些

  • 建议使用有界队列

results matching ""

    No results matching ""