开篇:3个致命痛点,催生DynamicTP的诞生
传统线程池的三大核心痛点,正是 DynamicTP 的发力点:
而 DynamicTP 的核心优势的在于:零代码侵入+动态调参+全方位监控+智能告警,完美适配微服务高并发、流量波动大的场景。
一、DynamicTP核心架构:5大模块+3大核心能力
DynamicTP 采用分层模块化设计,核心围绕“动态配置、监控告警、多组件兼容”三大能力展开,架构清晰且扩展性极强。
1. 核心模块拆解
dynamic-tp/
├── core模块:核心逻辑层(动态调参、监控采集、告警触发)
├── adapter模块:三方组件适配(Tomcat/Dubbo/RocketMQ等线程池管理)
├── starter模块:配置中心集成(Nacos/Apollo/ZK等,开箱即用)
├── logging模块:监控日志输出(JSON格式落地,支持自定义路径)
├── common模块:通用工具类(配置解析、SPI扩展、线程工具)
2. 三大核心能力
| 动态调参 | 配置中心监听+事件驱动+平滑调整 | 运行时修改参数,无需重启服务 |
| 监控采集 | 定时采集20+指标(活跃线程数、拒绝率等) | 消除监控盲区,实时掌握线程池状态 |
| 智能告警 | 多维度阈值告警+多渠道通知 | 故障事前预警,降低损失 |
3. 核心工作流程
二、源码深度拆解:DynamicTP的3个核心秘密
(一)核心类关系:基于JDK线程池的扩展
java.util.concurrent.ThreadPoolExecutor(JDK原生)
↓ 继承
cn.dynamictp.core.threadpool.DtpExecutor(核心扩展类)
↓ 依赖
cn.dynamictp.core.threadpool.support.VariableLinkedBlockingQueue(可变容量队列)
cn.dynamictp.core.event.DtpRefreshEvent(配置变更事件)
cn.dynamictp.core.monitor.DtpMonitor(监控采集类)
(二)秘密1:动态调参的底层实现(配置监听+事件驱动)
动态调参是 DynamicTP 的核心,底层依赖「配置中心监听+Spring事件驱动」实现,全程无锁且不影响任务执行。
1. 配置监听源码(以Nacos为例)
NacosRefresher 类负责监听配置中心变更,核心逻辑如下:
public class NacosRefresher extends AbstractSpringRefresher implements SmartApplicationListener {
//监听事件
public void onApplicationEvent(ApplicationEvent event) {
if (event instanceof NacosConfigEvent) {
this.refresh (this.environment);
}
}
}
org.dromara.dynamictp.core.refresher.AbstractRefresher#refresh(java.lang.Object)
protected void refresh(Object environment) {
BinderHelper.bindDtpProperties(environment, dtpProperties);
doRefresh(dtpProperties);
}
//开始刷新线程池参数
protected void doRefresh(DtpProperties properties) {
DtpRegistry.refresh(properties);
publishEvent(properties);
}
//发布刷新事件
private void publishEvent(DtpProperties dtpProperties) {
RefreshEvent event = new RefreshEvent(this, dtpProperties);
EventBusManager.post(event);
}
2. 事件处理源码(参数刷新核心)
org.dromara.dynamictp.core.DtpRegistry#refresh(org.dromara.dynamictp.core.support.ExecutorWrapper, org.dromara.dynamictp.common.entity.DtpExecutorProps)
private static void refresh(ExecutorWrapper executorWrapper, DtpExecutorProps props) {
//ExecutorWrapper为线程池的包装体,封装了线程池的信息
if (props.coreParamIsInValid()) {
log.error("DynamicTp refresh, invalid parameters exist, properties: {}", props);
return;
}
TpMainFields oldFields = ExecutorConverter.toMainFields(executorWrapper);
//修改线程池参数
doRefresh(executorWrapper, props);
TpMainFields newFields = ExecutorConverter.toMainFields(executorWrapper);
if (oldFields.equals(newFields)) {
log.debug("DynamicTp refresh, main properties of [{}] have not changed.",
executorWrapper.getThreadPoolName());
return;
}
// Get the changed keys
List<FieldInfo> diffFields = EQUATOR.getDiffFields(oldFields, newFields);
List<String> diffKeys = StreamUtil.fetchProperty(diffFields, FieldInfo::getFieldName);
NoticeManager.tryNoticeAsync(executorWrapper, oldFields, diffKeys);
log.info("DynamicTp refresh, tpName: [{}], changed keys: {}, corePoolSize: [{}], maxPoolSize: [{}]," +
" queueType: [{}], queueCapacity: [{}], keepAliveTime: [{}], rejectedType: [{}]," +
" allowsCoreThreadTimeOut: [{}]", executorWrapper.getThreadPoolName(), diffKeys,
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getCorePoolSize(), newFields.getCorePoolSize()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getMaxPoolSize(), newFields.getMaxPoolSize()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getQueueType(), newFields.getQueueType()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getQueueCapacity(), newFields.getQueueCapacity()),
String.format("%ss => %ss", oldFields.getKeepAliveTime(), newFields.getKeepAliveTime()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.getRejectType(), newFields.getRejectType()),
String.format(PROPERTIES_CHANGE_SHOW_STYLE, oldFields.isAllowCoreThreadTimeOut(),
newFields.isAllowCoreThreadTimeOut()));
}
private static void doRefresh(ExecutorWrapper executorWrapper, DtpExecutorProps props) {
ExecutorAdapter<?> executor = executorWrapper.getExecutor();
//修改核心线程数和最大线程数
doRefreshPoolSize(executor, props);
if (!Objects.equals(executor.getKeepAliveTime(props.getUnit()), props.getKeepAliveTime())) {
executor.setKeepAliveTime(props.getKeepAliveTime(), props.getUnit());
}
if (!Objects.equals(executor.allowsCoreThreadTimeOut(), props.isAllowCoreThreadTimeOut())) {
executor.allowCoreThreadTimeOut(props.isAllowCoreThreadTimeOut());
}
// 修改任务队列的长度
updateQueueProps(executor, props);
if (executorWrapper.isDtpExecutor()) {
doRefreshDtp(executorWrapper, props);
return;
}
doRefreshCommon(executorWrapper, props);
}
(三)秘密2:可变容量队列(解决JDK队列不可扩容问题)
JDK 原生 LinkedBlockingQueue 容量固定,DynamicTP 封装 VariableLinkedBlockingQueue 支持动态扩容,适配队列容量的动态调整:
public class VariableLinkedBlockingQueue<E> extends AbstractQueue<E> implements BlockingQueue<E>, Serializable {
//任务队列长度
private volatile int capacity;
//当前任务数
private final AtomicInteger count;
// 动态调整队列容量
public void setCapacity(int capacity) {
int oldCapacity = this.capacity;
this.capacity = capacity;
int size = this.count.get();
if (capacity > size && size >= oldCapacity) {
this.signalNotFull();
}
}
@Override
public boolean offer(E e) {
if (e == null) {
throw new NullPointerException();
} else {
AtomicInteger count = this.count;
if (count.get() >= this.capacity) {
return false;
} else {
int c = –1;
Node<E> node = new Node(e);
ReentrantLock putLock = this.putLock;
putLock.lock();
try {
if (count.get() < this.capacity) {
this.enqueue(node);
c = count.getAndIncrement();
if (c + 1 < this.capacity) {
this.notFull.signal();
}
}
} finally {
putLock.unlock();
}
if (c == 0) {
this.signalNotEmpty();
}
return c >= 0;
}
}
}
}
(四)秘密3:监控告警机制
1. 监控指标采集源码(DtpMonitor)
org.dromara.dynamictp.core.monitor.DtpMonitor#run
private void run() {
Set<String> executorNames = DtpRegistry.getAllExecutorNames();
try {
//检测报警
checkAlarm(executorNames);
//收集监控指标
collectMetrics(executorNames);
} catch (Exception e) {
log.error("DynamicTp monitor, run error", e);
}
}
2. 告警规则源码
org.dromara.dynamictp.common.entity.NotifyItem
public class NotifyItem { //只展示关键属性信息
//统计窗口+触发次数」,减少无效告警
// 告警阈值(如队列容量使用率70%)
private int threshold = 70;
// 统计周期(如120秒)
private int period = 120;
// 周期内触发次数
private int count;
// 静默时间(告警后多久不重复通知)
private int silencePeriod = 120;
}
3. 多渠道告警实现(SPI扩展)
DynamicTP 通过 SPI 接口支持自定义告警渠道
// 告警渠道SPI接口
public interface Notifier {
String platform();
void send(NotifyPlatform var1, String var2);
}
// 钉钉告警实现
public class DtpDingNotifier extends AbstractDtpNotifier {
public DtpDingNotifier(Notifier notifier) {
super(notifier);
}
//报警平台
@Override
public String platform() {
return NotifyPlatformEnum.DING.name().toLowerCase();
}
//通知模板,线程池参数变更
@Override
protected String getNoticeTemplate() {
return DingNotifyConst.DING_CHANGE_NOTICE_TEMPLATE;
}
//报警通知
@Override
protected String getAlarmTemplate() {
return DingNotifyConst.DING_ALARM_TEMPLATE;
}
@Override
protected Pair<String, String> getColors() {
return new ImmutablePair<>(DingNotifyConst.WARNING_COLOR, DingNotifyConst.CONTENT_COLOR);
}
}
三、核心总结
结尾
以上是对 DynamicTP 主要执行流程的拆解,后续还会对其源码进行解析。



