ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Dubbo线程池策略详解与性能优化

Dubbo线程池策略详解与性能优化 1. HTML基础概念解析HTMLHyperText Markup Language作为构建网页的基础语言其核心功能在于定义文档的结构和内容。与CSS负责样式、JavaScript负责行为不同HTML专注于内容的语义化组织。这种分工明确的体系使得Web开发能够保持清晰的层次结构。超文本hypertext的概念是HTML的关键特征之一。它通过链接将不同文档相互关联这种非线性组织结构正是万维网的核心特征。在实际开发中合理使用 标签创建链接不仅能提升用户体验更是SEO优化的重要环节。2. HTML元素与标签详解2.1 基本文档结构一个标准的HTML文档包含以下必要元素!DOCTYPE html html langzh-CN head meta charsetUTF-8 title页面标题/title /head body !-- 页面内容 -- /body /html重要提示lang属性的正确设置会影响浏览器的拼写检查和行为同时也有助于屏幕阅读器等辅助技术的正确解析。2.2 常用元素分类结构元素header、footer、article、section、nav文本元素h1-h6、p、span、strong、em媒体元素img、video、audio表单元素form、input、textarea、select列表元素ul、ol、li3. HTML5新特性实践3.1 语义化标签的应用现代HTML开发强烈推荐使用语义化标签article header h1文章标题/h1 time datetime2023-07-202023年7月20日/time /header section p正文内容.../p figure img srcexample.jpg alt示例图片 figcaption图片说明/figcaption /figure /section footer p版权信息/p /footer /article3.2 多媒体支持HTML5原生支持多媒体播放video controls width640 source srcmovie.mp4 typevideo/mp4 track kindsubtitles srcsubtitles.vtt srclangzh label中文 您的浏览器不支持视频播放 /video4. 表单与用户交互4.1 现代表单元素form label foremail电子邮箱/label input typeemail idemail nameemail required label forrange满意度/label input typerange idrange namerange min0# 1. 概述 本文分享 **Dubbo 的线程池策略**。在 [《精尽 Dubbo 源码分析 —— 线程池》](http://svip.iocoder.cn/Dubbo/thread-policy/?self) 一文中我们已经看到Dubbo 提供了**五种**线程池的实现 - fixed 固定大小线程池启动时建立线程不关闭一直持有。(**缺省**) - cached 缓存线程池空闲一分钟自动删除需要时重建。 - limited 可伸缩线程池但池中的线程数只会增长不会收缩。只增长不收缩的目的是为了避免收缩时突然来了大流量引起的性能问题。 - eager 优先创建Worker线程池。在任务数大于corePoolSize但是小于maximumPoolSize时优先创建Worker来处理任务。当任务数大于maximumPoolSize时将任务放入阻塞队列中。阻塞队列充满时抛出RejectedExecutionException。(相比于cached:cached在任务数量超过maximumPoolSize时直接抛出异常而不是将任务放入阻塞队列) - direct 直接线程池在需要时创建新线程不会有任何限制以避免出现新的请求时找不到可用的线程。 本文涉及的类如下图所示 ![类图](http://www.iocoder.cn/images/Dubbo/2018_03_01/01.png) # 2. ThreadPool com.alibaba.dubbo.common.threadpool.ThreadPool 线程池接口。代码如下 java SPI(fixed) public interface ThreadPool { /** * 线程池 * * param url 线程参数 * return 线程池 */ Adaptive({Constants.THREADPOOL_KEY}) Executor getExecutor(URL url); }SPI(fixed)注解Dubbo SPI拓展点默认为fixed。Adaptive({Constants.THREADPOOL_KEY})注解基于 Dubbo SPI Adaptive 机制加载对应的线程池实现使用URL.threadpool属性。#getExecutor(url)方法获得对应的线程池的执行器。3. FixedThreadPoolcom.alibaba.dubbo.common.threadpool.support.fixed.FixedThreadPool实现 ThreadPool 接口固定大小线程池启动时建立线程不关闭一直持有。代码如下public class FixedThreadPool implements ThreadPool { Override public Executor getExecutor(URL url) { // 线程名 String name url.getParameter(Constants.THREAD_NAME_KEY, Constants.DEFAULT_THREAD_NAME); // 线程数 int threads url.getParameter(Constants.THREADS_KEY, Constants.DEFAULT_THREADS); // 队列数 int queues url.getParameter(Constants.QUEUES_KEY, Constants.DEFAULT_QUEUES); // 创建执行器 return new ThreadPoolExecutor(threads, threads, 0, TimeUnit.MILLISECONDS, queues 0 ? new SynchronousQueueRunnable() : (queues 0 ? new LinkedBlockingQueueRunnable() : new LinkedBlockingQueueRunnable(queues)), new NamedThreadFactory(name, true), new AbortPolicyWithReport(name, url)); } }默认情况下采用FixedThreadPool。具体参数解析见代码注释。创建的执行器是java.util.concurrent.ThreadPoolExecutor。4. CachedThreadPoolcom.alibaba.dubbo.common.threadpool.support.cached.CachedThreadPool实现 ThreadPool 接口缓存线程池空闲一定时长自动删除需要时重建。代码如下public class CachedThreadPool implements ThreadPool { Override public Executor getExecutor(URL url) { // 线程名 String name url.getParameter(Constants.THREAD_NAME_KEY, Constants.DEFAULT_THREAD_NAME); // 核心线程数 int cores url.getParameter(Constants.CORE_THREADS_KEY, Constants.DEFAULT_CORE_THREADS); // 最大线程数 int threads url.getParameter(Constants.THREADS_KEY, Integer.MAX_VALUE); // 队列数 int queues url.getParameter(Constants.QUEUES_KEY, Constants.DEFAULT_QUEUES); // 空闲线程过多久则回收 int alive url.getParameter(Constants.ALIVE_KEY, Constants.DEFAULT_ALIVE); // 创建执行器 return new ThreadPoolExecutor(cores, threads, alive, TimeUnit.MILLISECONDS, queues 0 ? new SynchronousQueueRunnable() : (queues 0 ? new LinkedBlockingQueueRunnable() : new LinkedBlockingQueueRunnable(queues)), new NamedThreadFactory(name, true), new AbortPolicyWithReport(name, url)); } }具体参数解析见代码注释。创建的执行器是java.util.concurrent.ThreadPoolExecutor。5. LimitedThreadPoolcom.alibaba.dubbo.common.threadpool.support.limited.LimitedThreadPool实现 ThreadPool 接口可伸缩线程池但池中的线程数只会增长不会收缩。只增长不收缩的目的是为了避免收缩时突然来了大流量引起的性能问题。代码如下public class LimitedThreadPool implements ThreadPool { Override public Executor getExecutor(URL url) { // 线程名 String name url.getParameter(Constants.THREAD_NAME_KEY, Constants.DEFAULT_THREAD_NAME); // 核心线程数 int cores url.getParameter(Constants.CORE_THREADS_KEY, Constants.DEFAULT_CORE_THREADS); // 最大线程数 int threads url.getParameter(Constants.THREADS_KEY, Constants.DEFAULT_THREADS); // 队列数 int queues url.getParameter(Constants.QUEUES_KEY, Constants.DEFAULT_QUEUES); // 创建执行器 return new ThreadPoolExecutor(cores, threads, Long.MAX_VALUE, TimeUnit.MILLISECONDS, queues 0 ? new SynchronousQueueRunnable() : (queues 0 ? new LinkedBlockingQueueRunnable() : new LinkedBlockingQueueRunnable(queues)), new NamedThreadFactory(name, true), new AbortPolicyWithReport(name, url)); } }具体参数解析见代码注释。创建的执行器是java.util.concurrent.ThreadPoolExecutor。6. EagerThreadPoolcom.alibaba.dubbo.common.threadpool.support.eager.EagerThreadPool实现 ThreadPool 接口优先创建Worker线程池。在任务数大于corePoolSize但是小于maximumPoolSize时优先创建Worker来处理任务。当任务数大于maximumPoolSize时将任务放入阻塞队列中。阻塞队列充满时抛出RejectedExecutionException。(相比于cached:cached在任务数量超过maximumPoolSize时直接抛出异常而不是将任务放入阻塞队列)6.1 EagerThreadPoolpublic class EagerThreadPool implements ThreadPool { Override public Executor getExecutor(URL url) { // 线程名 String name url.getParameter(Constants.THREAD_NAME_KEY, Constants.DEFAULT_THREAD_NAME); // 核心线程数 int cores url.getParameter(Constants.CORE_THREADS_KEY, Constants.DEFAULT_CORE_THREADS); // 最大线程数 int threads url.getParameter(Constants.THREADS_KEY, Integer.MAX_VALUE); // 队列数 int queues url.getParameter(Constants.QUEUES_KEY, Constants.DEFAULT_QUEUES); // 空闲线程过多久则回收 int alive url.getParameter(Constants.ALIVE_KEY, Constants.DEFAULT_ALIVE); // 初始化队列和线程池 // 创建执行器 TaskQueueRunnable taskQueue new TaskQueueRunnable(queues 0 ? 1 : queues); EagerThreadPoolExecutor executor new EagerThreadPoolExecutor(cores, threads, alive, TimeUnit.MILLISECONDS, taskQueue, new NamedThreadFactory(name, true), new AbortPolicyWithReport(name, url)); taskQueue.setExecutor(executor); return executor; } }具体参数解析见代码注释。创建的执行器是EagerThreadPoolExecutor。6.2 TaskQueuecom.alibaba.dubbo.common.threadpool.support.eager.TaskQueue实现java.util.concurrent.LinkedBlockingQueue类任务队列。通过重写#offer(...)方法实现已提交的任务数量小于等于线程数时优先创建线程而不是放入队列。代码如下Override public boolean offer(Runnable runnable) { if (executor null) { throw new RejectedExecutionException(executor not initialized.); } // 当前线程数 int currentPoolThreadSize executor.getPoolSize(); // 有空闲线程 if (executor.getSubmittedTaskCount() currentPoolThreadSize) { return super.offer(runnable); } // 当前线程数小于最大线程数创建新线程 if (currentPoolThreadSize executor.getMaximumPoolSize()) { return false; } // 提交到队列 return super.offer(runnable); }当executor.getSubmittedTaskCount() currentPoolThreadSize时表示有空闲线程直接提交到队列。当currentPoolThreadSize executor.getMaximumPoolSize()时创建新线程返回false避免提交到队列。其它提交到队列。6.3 EagerThreadPoolExecutorcom.alibaba.dubbo.common.threadpool.support.eager.EagerThreadPoolExecutor实现java.util.concurrent.ThreadPoolExecutor类快速消费任务线程池执行器。代码如下public class EagerThreadPoolExecutor extends ThreadPoolExecutor { /** * 已提交的任务数量 */ private final AtomicInteger submittedTaskCount new AtomicInteger(0); public EagerThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, TaskQueueRunnable workQueue, ThreadFactory threadFactory, RejectedExecutionHandler handler) { super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler); } public int getSubmittedTaskCount() { return submittedTaskCount.get(); } Override protected void afterExecute(Runnable r, Throwable t) { submittedTaskCount.decrementAndGet(); } Override public void execute(Runnable command) { if (command null) { throw new NullPointerException(); } // 提交任务数 1 submittedTaskCount.incrementAndGet(); try { // 执行任务 super.execute(command); } catch (RejectedExecutionException rx) { // 发生拒绝异常尝试重新加入队列 final TaskQueue queue (TaskQueue) super.getQueue(); try { if (!queue.retryOffer(command, 0, TimeUnit.MILLISECONDS)) { submittedTaskCount.decrementAndGet(); throw new RejectedExecutionException(Queue capacity is full., rx); } } catch (InterruptedException x) { submittedTaskCount.decrementAndGet(); throw new RejectedExecutionException(x); } } catch (Throwable t) { // 提交任务数 -1 submittedTaskCount.decrementAndGet(); throw t; } } }submittedTaskCount属性已提交的任务数量。通过AtomicInteger实现线程安全。#execute(command)方法提交任务数 1 而后调用父类方法执行任务。若发生异常当发生 RejectedExecutionException 异常时调用TaskQueue#retryOffer(command, 0, TimeUnit.MILLISECONDS)方法尝试重新加入队列。若失败抛出 RejectedExecutionException 异常。当发生其他异常时提交任务数 -1 抛出该异常。#afterExecute(r, t)方法执行完成提交任务数 -1 。7. DirectExecutorcom.alibaba.dubbo.common.threadpool.support.direct.DirectExecutor实现 ThreadPool 接口直接线程池在需要时创建新线程不会有任何限制以避免出现新的请求时找不到可用的线程。代码如下public class DirectExecutor implements ThreadPool { Override public Executor getExecutor(URL url) { return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS, new SynchronousQueueRunnable(), new NamedThreadFactory(DirectExecutor, true), new AbortPolicyWithReport(DirectExecutor, url)); } }创建的执行器是java.util.concurrent.ThreadPoolExecutor。8. AbortPolicyWithReportcom.alibaba.dubbo.common.threadpool.support.AbortPolicyWithReport实现java.util.concurrent.RejectedExecutionHandler接口拒绝策略实现类。打印 JStack 分析线程状态。代码如下1: public class AbortPolicyWithReport extends AbstractRejectedExecutionHandler { 2: 3: /** 4: * 线程名 5: */ 6: protected final String threadName; 7: 8: /** 9: * URL 对象 10: */ 11: protected final URL url; 12: 13: /** 14: * 最后打印时间 15: */ 16: private static volatile long lastPrintTime 0; 17: 18: /** 19: * 信号量大小为 1 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: Override 29: public void rejectedExecution(Runnable r, ThreadPoolExecutor e) { 30: // 打印告警日志 31: String msg String.format(Thread pool is EXHAUSTED! 32: Thread Name: %s, Pool Size: %d (active: %d, core: %d, max: %d, largest: %d), Task: %d (completed: %d), 33: Executor status:(isShutdown:%s, isTerminated:%s, isTerminating:%s), in %s://%s:%d!, 34: threadName, e.getPoolSize(), e.getActiveCount(), e.getCorePoolSize(), e.getMaximumPoolSize(), e.getLargestPoolSize(), 35: e.getTaskCount(), e.getCompletedTaskCount(), e.isShutdown(), e.isTerminated(), e.isTerminating(), 36: url.getProtocol(), url.getIp(), url.getPort()); 37: logger.warn(msg); 38: // 打印 JStack 分析线程状态 39: dumpJStack(); 40: // 抛出 RejectedExecutionException 异常 41: throw new RejectedExecutionException(msg); 42: } 43: 44: /** 45: * 打印 JStack 46: */ 47: private void dumpJStack() { 48: long now System.currentTimeMillis(); 49: 50: // 每 10 分钟打印一次 51: // dump every 10 minutes 52: if (now - lastPrintTime 10 * 60 * 1000) { 53: return; 54: } 55: 56: // 获得信号量 57: if (!guard.tryAcquire()) { 58: return; 59: } 60: 61: // 创建线程池后台执行打印 JStack 62: ExecutorService pool Executors.newSingleThreadExecutor(); 63: pool.execute(new Runnable() { 64: Override 65: public void run() { 66: // 获得系统 67: String dumpPath url.getParameter(Constants.DUMP_DIRECTORY, System.getProperty(user.home)); 68: 69: SimpleDateFormat sdf; 70: // 获得系统 71: String os System.getProperty(os.name).toLowerCase(); 72: // window system dont support : in file name 73: if (os.contains(win)) { 74: sdf new SimpleDateFormat(yyyy-MM-dd_HH-mm-ss); 75: } else { 76: sdf new SimpleDateFormat(yyyy-MM-dd_HH:mm:ss); 77: } 78: 79: String dateStr sdf.format(new Date()); 80: // 获得输出流 81: FileOutputStream jstackStream null; 82: try { 83: jstackStream new FileOutputStream(new File(dumpPath, Dubbo_JStack.log . dateStr)); 84: // 打印 JStack 85: JVMUtil.jstack(jstackStream); 86: } catch (Throwable t) { 87: logger.error(dump jstack error, t); 88: } finally { 89: // 释放信号量 90: guard.release(); 91: // 关闭输出流 92: if (jstackStream ! null) { 93: try { 94: jstackStream.flush(); 95: jstackStream.close(); 96: } catch (IOException e) { 97: } 98: } 99: } 100: 101: // 记录最后打印时间 102: lastPrintTime System.currentTimeMillis(); 103: } 104: }); 105: // 关闭线程池 106: pool.shutdown(); 107: } 108: 109: }#rejectedExecution(Runnable r, ThreadPoolExecutor e)实现方法第 30 至 37 行打印告警日志。例如[14/04/18 11:29:05:029 WARN] AbortPolicyWithReport: [DUBBO] Thread pool is EXHAUSTED! Thread Name: DubboServerHandler-172.16.132.166:20880, Pool Size: 200 (active: 200, core: 200, max: 200, largest: 200), Task: 295 (completed: 95), Executor status:(isShutdown:false, isTerminated:false, isTerminating:false), in dubbo://172.16.132.166:20880!一般情况下线程池满的原因是服务响应慢阻塞执行线程逐渐撑满线程池。此时可以通过dumpJStack()方法查看线程都在执行什么代码从而排查问题。第 39 行调用#dumpJStack()方法打印JStack分析线程状态。第 41 行抛出 RejectedExecutionException 异常。#dumpJStack()方法第 51 至 54 行每 10 分钟打印一次。第 56 至 59 行获得信号量。保证同一时间有且仅有一个线程执行打印。第 62 行创建线程池后台执行打印 JStack 。第 63 至 104 行调用JVMUtil#jstack(OutputStream)方法打印 JStack 。
返回列表