Executor—->Executors
Array———->Arrays
Collection—>Collections

  1. public class MyThreadPoolDemo {
  2. public static void main(String[] args) {
  3. //ExecutorService threadPool = Executors.newFixedThreadPool(5);//一池固定5个处理线程,执行长期任务
  4. //ExecutorService threadPool = Executors.newSingleThreadExecutor();//1池固定1线程,一个任务一个任务执行
  5. ExecutorService threadPool = Executors.newCachedThreadPool();//1池N线程 执行短期异步的小程序或负载轻的程序
  6. //ExecutorService threadPool = Executors.newScheduledThreadPool(3);//创建一个定长线程池,支持定时及周期性任务执行
  7. try {
  8. for (int i = 0; i <= 10; i++) {
  9. threadPool.execute(()->{
  10. System.out.println(Thread.currentThread().getName() + "\t 办理业务");
  11. });
  12. //让线程暂停一会
  13. //TimeUnit.MILLISECONDS.sleep(1);
  14. }
  15. } catch (Exception e){
  16. e.printStackTrace();
  17. }finally {
  18. threadPool.shutdown();
  19. }
  20. }
  21. }

核心参数源码:

  1. public ThreadPoolExecutor(int corePoolSize,
  2. int maximumPoolSize,
  3. long keepAliveTime,
  4. TimeUnit unit,
  5. BlockingQueue<Runnable> workQueue,
  6. ThreadFactory threadFactory,
  7. RejectedExecutionHandler handler) {
  8. if (corePoolSize < 0 ||
  9. maximumPoolSize <= 0 ||
  10. maximumPoolSize < corePoolSize ||
  11. keepAliveTime < 0)
  12. throw new IllegalArgumentException();
  13. if (workQueue == null || threadFactory == null || handler == null)
  14. throw new NullPointerException();
  15. this.acc = System.getSecurityManager() == null ?
  16. null :
  17. AccessController.getContext();
  18. this.corePoolSize = corePoolSize;
  19. this.maximumPoolSize = maximumPoolSize;
  20. this.workQueue = workQueue;
  21. this.keepAliveTime = unit.toNanos(keepAliveTime);
  22. this.threadFactory = threadFactory;
  23. this.handler = handler;
  24. }

1.七大核心参数解析

1.corePoolSize 线程池中常驻核心线程数
2.maximumPoolSize 线程池中能够允许同时执行的最大线程数,这个值必须要大于等于1
3.keepAliveTime 多余的空闲线程存活时间,当前线程池线程数量超过corePoolSize,线程的空闲时间达到keepAliveTime,多余的空闲线程会销毁,直至线程池中的线程数等于corePoolSize
4.unit keepAliveTime的单位
5.workQueue BlockingQueue阻塞队列 相当于银行候客区,在超出corePoolSize的时候线程任务等待的区域
6.threadFactory 表示生成线程池中工作线程的线程工厂,用于创建线程,一般用默认即可。
7.handler 饱和策略(拒绝策略)

2.可选择的阻塞队列BlockingQueue详解

在重复一下新任务进入时线程池的执行策略:
如果运行的线程少于corePoolSize,则 Executor始终首选添加新的线程,而不进行排队。(如果当前运行的线程小于corePoolSize,则任务根本不会存入queue中,而是直接运行)
如果运行的线程大于等于 corePoolSize,则 Executor始终首选将请求加入队列,而不添加新的线程。
如果无法将请求加入队列,则创建新的线程,除非创建此线程超出 maximumPoolSize,在这种情况下,任务将被拒绝。
主要有3种类型的BlockingQueue:
无界队列
队列大小无限制,常用的为无界的LinkedBlockingQueue,使用该队列做为阻塞队列时要尤其当心,当任务耗时较长时可能会导致大量新任务在队列中堆积最终导致OOM。阅读代码发现,Executors.newFixedThreadPool 采用就是 LinkedBlockingQueue,而楼主踩到的就是这个坑,当QPS很高,发送数据很大,大量的任务被添加到这个无界LinkedBlockingQueue 中,导致cpu和内存飙升服务器挂掉。
有界队列
常用的有两类,一类是遵循FIFO原则的队列如ArrayBlockingQueue,另一类是优先级队列如PriorityBlockingQueue。PriorityBlockingQueue中的优先级由任务的Comparator决定。
使用有界队列时队列大小需和线程池大小互相配合,线程池较小有界队列较大时可减少内存消耗,降低cpu使用率和上下文切换,但是可能会限制系统吞吐量。
在我们的修复方案中,选择的就是这个类型的队列,虽然会有部分任务被丢失,但是我们线上是排序日志搜集任务,所以对部分对丢失是可以容忍的。
同步移交队列
如果不希望任务在队列中等待而是希望将任务直接移交给工作线程,可使用SynchronousQueue作为等待队列。SynchronousQueue不是一个真正的队列,而是一种线程之间移交的机制。要将一个元素放入SynchronousQueue中,必须有另一个线程正在等待接收这个元素。只有在使用无界线程池或者有饱和策略时才建议使用该队列。

3.可选择的饱和策略RejectedExecutionHandler详解

JDK主要提供了4种饱和策略供选择。4种策略都做为静态内部类在ThreadPoolExcutor中进行实现。

3.1 AbortPolicy中止策略

该策略是默认饱和策略。

  1. 1. public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
  2. 2. throw new RejectedExecutionException("Task " + r.toString() +
  3. 3. " rejected from " +
  4. 4. e.toString());
  5. 5. }

使用该策略时在饱和时会抛出RejectedExecutionException(继承自RuntimeException),调用者可捕获该异常自行处理。

3.2 DiscardPolicy抛弃策略

  1. 1. public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
  2. 2. }

如代码所示,不做任何处理直接抛弃任务

3.3 DiscardOldestPolicy抛弃旧任务策略

  1. 1. public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
  2. 2. if (!e.isShutdown()) {
  3. 3. e.getQueue().poll();
  4. 4. e.execute(r);
  5. 5. }
  6. 6. }

如代码,先将阻塞队列中的头元素出队抛弃,再尝试提交任务。如果此时阻塞队列使用PriorityBlockingQueue优先级队列,将会导致优先级最高的任务被抛弃,因此不建议将该种策略配合优先级队列使用。

3.4 CallerRunsPolicy调用者运行

  1. 1. public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
  2. 2. if (!e.isShutdown()) {
  3. 3. r.run();
  4. 4. }
  5. 5. }

既不抛弃任务也不抛出异常,直接运行任务的run方法,换言之将任务回退给调用者来直接运行。使用该策略时线程池饱和后将由调用线程池的主线程自己来执行任务,因此在执行任务的这段时间里主线程无法再提交新任务,从而使线程池中工作线程有时间将正在处理的任务处理完成。

4、有界队列和饱和策略的应用

  1. public class Test2 {
  2. public static void main(String[] args) {
  3. ThreadPoolExecutor tpe = new ThreadPoolExecutor(2, 5, 300, TimeUnit.MILLISECONDS, new ArrayBlockingQueue(3));
  4. // ThreadPoolExecutor tpe = new ThreadPoolExecutor(2, 5, 300, TimeUnit.MILLISECONDS, new ArrayBlockingQueue(3), new RejectedExecutionHandler() {
  5. // @Override
  6. // public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
  7. // if (!e.isShutdown()) {
  8. // r.run();
  9. // }
  10. // }
  11. // });
  12. //当i<=9时拒绝策略会抛出异常
  13. //RejectedExecutionHandler 可以将超出队列的任务让主线程去执行
  14. for(int i = 1; i <= 8; i++){
  15. //核心线程数
  16. System.out.println("线程池线程数=" + tpe.getPoolSize());
  17. //队列等待任务数
  18. System.out.println("队列等待任务数=" + tpe.getQueue().size());
  19. //完成任务数
  20. //System.out.println("完成任务数=" + tpe.getCompletedTaskCount());
  21. tpe.execute(new MyRunnable2(i));
  22. System.out.println("------------------------------");
  23. }
  24. tpe.shutdown();
  25. }
  26. }
  27. class MyRunnable2 implements Runnable{
  28. int num;
  29. public MyRunnable2(int num){
  30. this.num = num;
  31. }
  32. @Override
  33. public void run() {
  34. System.out.println(Thread.currentThread().getName() + "窗口正在办理业务==" + num);
  35. try {
  36. Thread.sleep(2000);
  37. } catch (InterruptedException e) {
  38. e.printStackTrace();
  39. }
  40. System.out.println(Thread.currentThread().getName() + "窗口完成办理业务==" + num);
  41. }
  42. }

5、创建线程池的四种方式

一程一线池———newSingleThreadExecutor()
一程3线池———-newFixedThreadPool(3);
一程N线池———-newCachedThreadPool();
一程多线池———Executors.newScheduledThreadPool(3);

  1. import java.util.concurrent.ExecutorService;
  2. import java.util.concurrent.Executors;
  3. import java.util.concurrent.ScheduledExecutorService;
  4. import java.util.concurrent.TimeUnit;
  5. /*
  6. * 线程池的工具类
  7. * 常用的四种线程池的快速创建方式
  8. */
  9. public class MyThreadPool2 {
  10. public static void main(String[] args) {
  11. //一程一线池,ThreadPoolExecutor(1,1,0**);
  12. //处理周期性任务,一个一个任务的执行
  13. ExecutorService es1=Executors.newSingleThreadExecutor();
  14. //一池3线程(指定线程数)==》处理数量固定,周期较长的任务
  15. ExecutorService es2=Executors.newFixedThreadPool(3);
  16. //一池n线程==》处理大量负载轻的任务
  17. //例如发短信,量大(注册,改密码)
  18. ExecutorService es3=Executors.newCachedThreadPool();
  19. //一池多线程==》处理指定周期的任务,自行定时执行任务
  20. ScheduledExecutorService es4=Executors.newScheduledThreadPool(3);
  21. //分别执行
  22. //5个人调用线程es3.execute
  23. // for (int i = 1; i <=5; i++) {
  24. // es3.execute(new ThreadBiz(i));
  25. // }
  26. // es3.shutdown();
  27. for (int i = 1; i <=3; i++) {
  28. //2秒后执行任务,每6秒执行一次
  29. es4.scheduleWithFixedDelay(new ThreadBiz(i), 2, 6, TimeUnit.SECONDS);
  30. }
  31. // es3.shutdown();es4线程池不需要关闭
  32. }