Dubbo 线程池源码解析

本文首发于个人微信公众号《andyqian》,期待你的关注!

前言

之前文章《Java线程池ThreadPoolExecutor》《ThreadPoolExecutor 原理解析》中,分别讲述了ThreadPoolExecutor 的概念以及原理,今天就一起来看看其在 Dubbo 框架中的应用。

ThreadFactory 与 AbortPolicy

  Dubbo 为我们提供了几种不同类型的线程池实现,其底层均使用的是 JDK 中的 ThreadPoolExecutor 线程池。ThreadPoolExecutor 我们都已经非常熟悉,其构造函数中有几个非常重要的参数。其中就包括:拒绝策略( ThreadPoolExecutor.AbortPolicy ) 以及 ThreadFactory,在 Dubbo 中自定义了 ThreadPoolExecutor.AbortPolicy 以及 ThreadFactory。在学习线程池之前,我们先来看看这两者的实现,更有益于后面的理解。

在Dubbo中NamedInternalThreadFactory 为自定义的线程 ThreadFactory 的子类。其类图如下:

其中 NamedInternalThreadFactory 类,其实现如下所示:

1public class NamedInternalThreadFactory extends NamedThreadFactory { 2 3 public NamedInternalThreadFactory() { 4 super(); 5 } 6 7 public NamedInternalThreadFactory(String prefix) { 8 super(prefix, false); 9 } 10 11 public NamedInternalThreadFactory(String prefix, boolean daemon) { 12 super(prefix, daemon); 13 } 14 15 @Override 16 public Thread newThread(Runnable runnable) { 17 String name = mPrefix + mThreadNum.getAndIncrement(); 18 InternalThread ret = new InternalThread(mGroup, runnable, name, 0); 19 ret.setDaemon(mDaemon); 20 return ret; 21 } 22}

其中 NamedThreadFactory 类的实现如下:

1public class NamedThreadFactory implements ThreadFactory { 2 3 protected static final AtomicInteger POOL_SEQ = new AtomicInteger(1); 4 5 protected final AtomicInteger mThreadNum = new AtomicInteger(1); 6 7 protected final String mPrefix; 8 9 protected final boolean mDaemon; 10 11 protected final ThreadGroup mGroup; 12 13 public NamedThreadFactory() { 14 this("pool-" + POOL_SEQ.getAndIncrement(), false); 15 } 16 17 public NamedThreadFactory(String prefix) { 18 this(prefix, false); 19 } 20 21 public NamedThreadFactory(String prefix, boolean daemon) { 22 mPrefix = prefix + "-thread-"; 23 mDaemon = daemon; 24 SecurityManager s = System.getSecurityManager(); 25 mGroup = (s == null) ? Thread.currentThread().getThreadGroup() : s.getThreadGroup(); 26 } 27 28 @Override 29 public Thread newThread(Runnable runnable) { 30 String name = mPrefix + mThreadNum.getAndIncrement(); 31 Thread ret = new Thread(mGroup, runnable, name, 0); 32 ret.setDaemon(mDaemon); 33 return ret; 34 } 35 36 public ThreadGroup getThreadGroup() { 37 return mGroup; 38 }

到这里,上述代码描述的是Dubbo对线程池中线程的命名规则,其作用是为了方便追踪信息。

接下来,我们来看下拒绝策略 AbortPolicyWithReport 类的实现,其类图如下所示:

源码如下:

1public class AbortPolicyWithReport extends ThreadPoolExecutor.AbortPolicy { 2 3 protected static final Logger logger = LoggerFactory.getLogger(AbortPolicyWithReport.class); 4 5 private final String threadName; 6 7 private final URL url; 8 9 private static volatile long lastPrintTime = 0; 10 11 private static final long TEN_MINUTES_MILLS = 10 * 60 * 1000; 12 13 private static final String OS_WIN_PREFIX = "win"; 14 15 private static final String OS_NAME_KEY = "os.name"; 16 17 private static final String WIN_DATETIME_FORMAT = "yyyy-MM-dd_HH-mm-ss"; 18 19 private static final String DEFAULT_DATETIME_FORMAT = "yyyy-MM-dd_HH:mm:ss"; 20 21 private static Semaphore guard = new Semaphore(1); 22 23 public AbortPolicyWithReport(String threadName, URL url) { 24 this.threadName = threadName; 25 this.url = url; 26 } 27 28 // 覆盖 父类 ThreadPoolExecutor.AbortPolicy 的 rejectedExecution 方法。 29 @Override 30 public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { 31 // 构造 warn 参数,其中包括:线程状态,线程池数量,活跃数量,核心线程池数量,最大线程池数量 等信息。 32 String msg = String.format("Thread pool is EXHAUSTED!" + 33 " Thread Name: %s, Pool Size: %d (active: %d, core: %d, max: %d, largest: %d), Task: %d (completed: " 34 + "%d)," + 35 " Executor status:(isShutdown:%s, isTerminated:%s, isTerminating:%s), in %s://%s:%d!", 36 threadName, e.getPoolSize(), e.getActiveCount(), e.getCorePoolSize(), e.getMaximumPoolSize(), 37 e.getLargestPoolSize(), 38 e.getTaskCount(), e.getCompletedTaskCount(), e.isShutdown(), e.isTerminated(), e.isTerminating(), 39 url.getProtocol(), url.getIp(), url.getPort()); 40 logger.warn(msg); 41 // dump 堆栈信息 42 dumpJStack(); 43 throw new RejectedExecutionException(msg); 44 } 45 46 // 当执行 rejectedExecution 方法时,会执行该方法。将会 dump 堆栈信息 至 DUMP_DIRECTORY 目录,默认为:user.name 目录下。 47 private void dumpJStack() { 48 long now = System.currentTimeMillis(); 49 50 //dump every 10 minutes 51 if (now - lastPrintTime < TEN_MINUTES_MILLS) { 52 return; 53 } 54 55 if (!guard.tryAcquire()) { 56 return; 57 } 58 59 ExecutorService pool = Executors.newSingleThreadExecutor(); 60 pool.execute(() -> { 61 String dumpPath = url.getParameter(DUMP_DIRECTORY, System.getProperty("user.home")); 62 63 SimpleDateFormat sdf; 64 65 String os = System.getProperty(OS_NAME_KEY).toLowerCase(); 66 67 // window system don't support ":" in file name 68 if (os.contains(OS_WIN_PREFIX)) { 69 sdf = new SimpleDateFormat(WIN_DATETIME_FORMAT); 70 } else { 71 sdf = new SimpleDateFormat(DEFAULT_DATETIME_FORMAT); 72 } 73 74 String dateStr = sdf.format(new Date()); 75 //try-with-resources 76 try (FileOutputStream jStackStream = new FileOutputStream( 77 new File(dumpPath, "Dubbo_JStack.log" + "." + dateStr))) { 78 // 工具类,此处实现省略,有兴趣的可以查看。 79 JVMUtil.jstack(jStackStream); 80 } catch (Throwable t) { 81 logger.error("dump jStack error", t); 82 } finally { 83 guard.release(); 84 } 85 lastPrintTime = System.currentTimeMillis(); 86 }); 87 //must shutdown thread pool ,if not will lead to OOM 88 pool.shutdown(); 89 } 90}

上面的代码不难,都是打日志,dump 堆栈信息,其目的就是:用于在线程池被打满时,也就是记录执行AbortPolicy时现场信息,主要是便于后期的分析与问题排查。

线程池的实现

  上面讲述了Dubbo线程池中自定义的 ThreadFactory 类 以及 AbortPolicyWithReport 类。接下来,我们继续讲解 Dubbo 提供的不同线程池实现,其类图如下所示:

1. LimitedThreadPool 线程池

源码如下:

1public class LimitedThreadPool implements ThreadPool { 2 3 @Override 4 public Executor getExecutor(URL url) { 5 String name = url.getParameter(THREAD_NAME_KEY, DEFAULT_THREAD_NAME); 6 int cores = url.getParameter(CORE_THREADS_KEY, DEFAULT_CORE_THREADS); 7 int threads = url.getParameter(THREADS_KEY, DEFAULT_THREADS); 8 int queues = url.getParameter(QUEUES_KEY, DEFAULT_QUEUES); 9 return new ThreadPoolExecutor(cores, threads, Long.MAX_VALUE, TimeUnit.MILLISECONDS, 10 queues == 0 ? new SynchronousQueue<Runnable>() : 11 (queues < 0 ? new LinkedBlockingQueue<Runnable>() 12 : new LinkedBlockingQueue<Runnable>(queues)), 13 new NamedInternalThreadFactory(name, true), new AbortPolicyWithReport(name, url)); 14 }

其中:

  1. THREAD_NAME_KEY 值为:threadname ,表示为:线程名,其默认值为:Dubbo。

  2. CORE_THREADS_KEY 值为:corethreads,表示:核心线程池数量,其默认值为:0。

  3. THREADS_KEY 值为:threads 表示:最大线程数,默认值为:200。

  4. QUEUES_KEY 值为:queues 表示:阻塞队列大小,默认值为:0。

备注:

  1. 该线程池中的 cores,threads 参数由外部制定,其中 keepAliveTime 值为:Long.MAX_VALUE,TimeUnit 为 TimeUnit.MILLISECONDS (毫秒)。(意味着线程池中的所有线程永不过期,理论上大于Long.MAX_VALUE 即会过期,因为其足够大,这里可以看为是永不过期 )。

  2. 此处使用了三目运算符:

    当 queues = 0 时,BlockingQueue为SynchronousQueue。

    当 queues < 0 时,则构造一个新的LinkedBlockingQueue。

    当 queues > 0 时,构造一个指定元素的LinkedBlockingQueue。

    1queues == 0 ? new SynchronousQueue<Runnable>() : 2 (queues < 0 ? new LinkedBlockingQueue<Runnable>() 3 : new LinkedBlockingQueue<Runnable>(queues)

该线程池的特点是:可以创建若干个线程,其默认值为 200,线程池中的线程生命周期非常长,甚至可以看做是永不过期。

2.  CachedThreadPool 线程池

源码:

1 public class CachedThreadPool implements ThreadPool { 2 3 @Override 4 public Executor getExecutor(URL url) { 5 String name = url.getParameter(THREAD_NAME_KEY, DEFAULT_THREAD_NAME); 6 int cores = url.getParameter(CORE_THREADS_KEY, DEFAULT_CORE_THREADS); 7 int threads = url.getParameter(THREADS_KEY, Integer.MAX_VALUE); 8 int queues = url.getParameter(QUEUES_KEY, DEFAULT_QUEUES); 9 int alive = url.getParameter(ALIVE_KEY, DEFAULT_ALIVE); 10 return new ThreadPoolExecutor(cores, threads, alive, TimeUnit.MILLISECONDS, 11 queues == 0 ? new SynchronousQueue<Runnable>() : 12 (queues < 0 ? new LinkedBlockingQueue<Runnable>() 13 : new LinkedBlockingQueue<Runnable>(queues)), 14 new NamedInternalThreadFactory(name, true), new AbortPolicyWithReport(name, url)); 15 } 16}

其中:

  1. THREAD_NAME_KEY 值为:threadname , 表示为:线程名,其默认值为:Dubbo。

  2. CORE_THREADS_KEY 值为:corethreads,表示为:核心线程池数量,其默认值为:0。

  3. THREADS_KEY 值为:threads,表示:最大线程数,默认值为:Integer.MAX_VALUE。

  4. QUEUES_KEY 值为:queues,表示:阻塞队列大小,默认值为:0。

  5. ALIVE_KEY 值为:alive, 表示: keepAliveTime 表示线程池中线程的存活时间,其默认值为:60 * 1000 (毫秒) 也就是一分钟。

该线程池的特点是:可创建无限多线程(在操作系统的限制下,会远远低于Integer.MAX_VALUE值,这里视为无限大),其线程的最大存活时间默认为 1 分钟。意味着可以创建无限多线程,但是线程的生命周期默认较短!

3. FixedThreadPool 线程池

1public class FixedThreadPool implements ThreadPool { 2 3 @Override 4 public Executor getExecutor(URL url) { 5 String name = url.getParameter(THREAD_NAME_KEY, DEFAULT_THREAD_NAME); 6 int threads = url.getParameter(THREADS_KEY, DEFAULT_THREADS); 7 int queues = url.getParameter(QUEUES_KEY, DEFAULT_QUEUES); 8 return new ThreadPoolExecutor(threads, threads, 0, TimeUnit.MILLISECONDS, 9 queues == 0 ? new SynchronousQueue<Runnable>() : 10 (queues < 0 ? new LinkedBlockingQueue<Runnable>() 11 : new LinkedBlockingQueue<Runnable>(queues)), 12 new NamedInternalThreadFactory(name, true), new AbortPolicyWithReport(name, url)); 13 }

其中:

  1. THREADS_KEY 值为:threads,表示:最大线程数,默认值为:200。

  2. QUEUES_KEY 值为:queues,表示:阻塞队列大小,默认值为:0。

  3. corePoolSize,maximumPoolSize 的线程数量均为:threads,(也就意味着核心线程数等于最大线程数)。

  4. keepAliveTime 的默认值为0,当线程数大于corePoolSize 时,多余的空闲线程会立即终止。

该线程池的特点是:该线程池中corePoolSize 数量 与 maxinumPoolSize 数量一致,当提交的任务大于核心线程池时,则会将其放入到LinkedBlockingQueue队列中等待执行,也是Dubbo中默认使用的线程池。

4. EagerThreadPool 线程池

源码:

1public class EagerThreadPool implements ThreadPool { 2 3 @Override 4 public Executor getExecutor(URL url) { 5 String name = url.getParameter(THREAD_NAME_KEY, DEFAULT_THREAD_NAME); 6 int cores = url.getParameter(CORE_THREADS_KEY, DEFAULT_CORE_THREADS); 7 int threads = url.getParameter(THREADS_KEY, Integer.MAX_VALUE); 8 int queues = url.getParameter(QUEUES_KEY, DEFAULT_QUEUES); 9 int alive = url.getParameter(ALIVE_KEY, DEFAULT_ALIVE); 10 11 // init queue and executor 12 TaskQueue<Runnable> taskQueue = new TaskQueue<Runnable>(queues <= 0 ? 1 : queues); 13 EagerThreadPoolExecutor executor = new EagerThreadPoolExecutor(cores, 14 threads, 15 alive, 16 TimeUnit.MILLISECONDS, 17 taskQueue, 18 new NamedInternalThreadFactory(name, true), 19 new AbortPolicyWithReport(name, url)); 20 taskQueue.setExecutor(executor); 21 return executor; 22 } 23}

备注:

  该线程池与上面的线程池实现方式有些不一样,上面是直接使用了ThreadPoolExecutor 类的构造函数。在该线程池实现中,首先构造了一个自定义的 EagerThreadPoolExecutor 线程池,其底层实现也是基于 ThreadPoolExecutor 类的,其代码如下所示:

1public class EagerThreadPoolExecutor extends ThreadPoolExecutor { 2 3 /** 4 * task count 5 */ 6 private final AtomicInteger submittedTaskCount = new AtomicInteger(0); 7 8 public EagerThreadPoolExecutor(int corePoolSize, 9 int maximumPoolSize, 10 long keepAliveTime, 11 TimeUnit unit, TaskQueue<Runnable> workQueue, 12 ThreadFactory threadFactory, 13 RejectedExecutionHandler handler) { 14 super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler); 15 } 16 17 /** 18 * @return current tasks which are executed 19 */ 20 public int getSubmittedTaskCount() { 21 return submittedTaskCount.get(); 22 } 23 24 @Override 25 protected void afterExecute(Runnable r, Throwable t) { 26 // 执行任务数依次递减 27 submittedTaskCount.decrementAndGet(); 28 } 29 30 @Override 31 public void execute(Runnable command) { 32 if (command == null) { 33 throw new NullPointerException(); 34 } 35 // do not increment in method beforeExecute! 36 // 提交任务书 依次递加。 37 submittedTaskCount.incrementAndGet(); 38 try { 39 // 调用父类方法执行线程任务 40 super.execute(command); 41 } catch (RejectedExecutionException rx) { 42 // 将任务重新添加到队列中 43 final TaskQueue queue = (TaskQueue) super.getQueue(); 44 try { 45 //如果添加失败,则减少任务数,并抛出异常。 46 if (!queue.retryOffer(command, 0, TimeUnit.MILLISECONDS)) { 47 submittedTaskCount.decrementAndGet(); 48 throw new RejectedExecutionException("Queue capacity is full.", rx); 49 } 50 } catch (InterruptedException x) { 51 submittedTaskCount.decrementAndGet(); 52 throw new RejectedExecutionException(x); 53 } 54 } catch (Throwable t) { 55 // decrease any way 56 submittedTaskCount.decrementAndGet(); 57 throw t; 58 } 59 }

在这里我们发现,在 EagerThreadPoolExecutor 类中,重载了父类ThreadPoolExecutor 类的几个方法,分别如下:afterExecuteexecute方法。分别加入 submittedTaskCount 属性进行任务的统计,当父类的execute方法抛出 RejectedExecutionExcetion 异常时,则会将任务重新放入队列中执行,其TaskQueue代码如下:

1public class TaskQueue<R extends Runnable> extends LinkedBlockingQueue<Runnable> { 2 3 private static final long serialVersionUID = -2635853580887179627L; 4 5 private EagerThreadPoolExecutor executor; 6 7 public TaskQueue(int capacity) { 8 super(capacity); 9 } 10 11 public void setExecutor(EagerThreadPoolExecutor exec) { 12 executor = exec; 13 } 14 15 @Override 16 public boolean offer(Runnable runnable) { 17 if (executor == null) { 18 throw new RejectedExecutionException("The task queue does not have executor!"); 19 } 20 21 int currentPoolThreadSize = executor.getPoolSize(); 22 // have free worker. put task into queue to let the worker deal with task. 23 // 当提交的任务数,小于 当前线程时,则之间调用父类的offer 方法。 24 if (executor.getSubmittedTaskCount() < currentPoolThreadSize) { 25 return super.offer(runnable); 26 } 27 28 // return false to let executor create new worker. 29 // 当当前线程数大小小于,最大线程数时,则直接返回false,创建worker。 30 if (currentPoolThreadSize < executor.getMaximumPoolSize()) { 31 return false; 32 } 33 34 // currentPoolThreadSize >= max 35 return super.offer(runnable); 36 } 37 38 /** 39 * retry offer task 40 * 41 * @param o task 42 * @return offer success or not 43 * @throws RejectedExecutionException if executor is terminated. 44 */ 45 public boolean retryOffer(Runnable o, long timeout, TimeUnit unit) throws InterruptedException { 46 if (executor.isShutdown()) { 47 throw new RejectedExecutionException("Executor is shutdown!"); 48 } 49 return super.offer(o, timeout, unit); 50 } 51}

该线程池的特点是:可以重新将拒绝掉的task,重新添加的work queue中执行。相当于有一个重试机制!

结语

  通过上面的分析,我相信大家对Dubbo中线程池应该有所了解。如果还有不清楚的地方,可以通过debug的方式进行跟踪分析。其实在很多的开源框架中,都有自定义的线程池,但其底层最终使用的还是 ThreadPoolExecutor 线程池,这个知识点建议大家一定要掌握,无论是实际工作还是面试,都是一个常用的知识点。


相关阅读:

你所不知道的 BigDecimal

ThreadPoolExecutor 原理解析

Java线程池ThreadPoolExecutor

使用 Mybatis 真心不要偷懒!

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

Dubbo 线程池源码解析 - HelloWorld