前言
在上一篇《DynamicTP 源码解析系列(四):动态调参核心原理》中,我们讲透了 DynamicTP 最核心的动态能力——通过配置中心联动+refresh 方法,实现线程池参数的平滑刷新,快速适配生产流量波动。
但动态调参只是“手段”,不是“终点”——生产中,我们不仅要能“动态调参”,更要能“实时感知线程池状态”:调参后参数是否生效?线程池是否出现任务堆积、拒绝?空闲线程是否过多导致资源浪费?一旦出现异常,如何第一时间收到通知、快速排查?这就是 DynamicTP 监控告警模块的核心价值,也是它成为“生产级框架”的关键支撑。
本文讲透 DynamicTP 监控告警的完整核心原理,聚焦两个核心问题:① 线程池的核心指标(活跃线程数、任务拒绝数等)如何被实时采集?② 当指标触发告警阈值时,如何联动钉钉/企业微信等渠道,实现告警信息的精准推送?
一、先理清:监控告警的核心定位与整体流程
1.1 核心定位(为什么需要监控告警?)
DynamicTP 的监控告警模块,核心定位是「线程池状态的“晴雨表”+ 异常的“预警器”」,解决原生线程池的两大痛点:
-
原生线程池无监控:无法实时获取活跃线程数、任务堆积数、拒绝数等核心指标,出现异常(如任务大量拒绝)后,只能通过业务报错反向排查,排查成本高、耗时久;
-
动态调参无反馈:第四篇讲的动态调参,修改参数后无法实时确认“调参效果”(如核心线程数调大后,活跃线程数是否同步提升),只能靠经验判断,易出现调参无效、参数不合理的问题。
1.2 与动态调参的联动逻辑
监控告警与动态调参不是孤立的,而是“相辅相成”的闭环:
动态调参触发监控更新:第四篇讲的 refresh 方法(动态调参入口),在刷新线程池参数(如 corePoolSize、queueCapacity)后,会同步触发监控采集,实时更新指标,方便开发者确认调参效果;
监控指标指导动态调参:通过监控采集的指标(如任务堆积数持续上升),可以判断当前线程池参数不合理,进而触发动态调参(手动调参),实现“监控→调参”的闭环;
异常告警倒逼问题解决:当监控指标触发告警阈值(如任务拒绝数≥1),告警推送后,开发者可快速排查问题,若为参数不合理导致,可立即通过动态调参修复,避免异常扩大。
1.3 整体流程
DynamicTP 监控告警的完整流程,本质是“指标采集 → 指标判断 → 告警触发 → 告警推送 ”,联动动态调参、配置中心。
-
监控采集的核心是 DtpMonitor:负责从 DtpExecutor 中采集各类指标;
-
告警管理的核心是 AlarmManager:负责加载告警规则、判断指标是否触发阈值、触发告警并推送,是连接“监控采集”和“告警推送”的关键;
-
配置联动:监控告警阈值、告警渠道等,均可在配置中心动态修改(无需重启服务);
-
调参联动:动态调参后,会触发监控采集,确保调参后的指标能实时更新,方便开发者确认调参效果。
二、源码拆解:指标采集核心原理
指标采集是监控告警的“基础”——没有精准、实时的指标采集,后续的告警判断、推送都无从谈起。DynamicTP 的指标采集核心是 DtpMonitor 类,配合 DtpExecutor 的指标暴露能力,实现“实时采集、低损耗、可配置”的采集效果。
2.1核心类:DtpMonitor(监控采集器)
DtpMonitor 的核心作用:按照配置的采集间隔,从所有 DtpExecutor 实例中采集上述核心指标,将采集到的指标存储,同时触发告警判断逻辑。
源码简化(聚焦核心逻辑:初始化、采集触发、指标采集):
/**
* 动态线程池监控采集器,核心负责指标采集、触发告警判断
*/
public class DtpMonitor {
// 监控采集间隔(ms),从配置中心加载(默认 5000ms,可动态修改)
private final long monitorInterval;
// 线程池配置信息
private final DtpProperties dtpProperties;
private ScheduledFuture<?> monitorFuture;
//调度线程池
private static ScheduledExecutorService monitorExecutor;
public DtpMonitor(DtpProperties dtpProperties) {
this.dtpProperties = dtpProperties;
EventBusManager.register(this);
}
//监听事件CustomContextRefreshedEvent
@Subscribe
public synchronized void onContextRefreshedEvent(CustomContextRefreshedEvent event) {
if (this.monitorInterval != this.dtpProperties.getMonitorInterval()) {
if (this.monitorFuture != null) {
this.monitorFuture.cancel(true);
}
// 单线程定时任务池,核心:避免并发采集多个线程池指标,导致锁竞争、性能损耗
if (monitorExecutor == null || monitorExecutor.isShutdown() || monitorExecutor.isTerminated()) {
monitorExecutor = ThreadPoolCreator.newScheduledThreadPool("dtp-monitor", 1);
}
// 延迟 0 s 启动,每隔 monitorInterval s 执行一次采集
this.monitorInterval = this.dtpProperties.getMonitorInterval();
this.monitorFuture = monitorExecutor.scheduleWithFixedDelay(this::run, 0L, (long)this.monitorInterval, TimeUnit.SECONDS);
}
}
/**
* 核心方法:执行监控采集(定时任务执行的逻辑)
*/
private void run() {
// 1. 从注册表中获取所有 DtpExecutor 实例(多线程池场景适配)
Set<String> executorNames = DtpRegistry.getAllExecutorNames();
try {
// 判断是否触发告警
this.checkAlarm(executorNames);
//采集单个线程池的核心指标
this.collectMetrics(executorNames);
} catch (Exception e) {
log.error("DynamicTp monitor, run error", e);
}
}
//判断是否触发告警
private void checkAlarm(Set<String> executorNames) {
executorNames.forEach((name) -> {
ExecutorWrapper wrapper = DtpRegistry.getExecutorWrapper(name);
AlarmManager.checkAndTryAlarmAsync(wrapper, DynamicTpConst.SCHEDULE_NOTIFY_ITEMS);
});
this.publishAlarmCheckEvent();
}
//采集单个线程池的核心指标
private void collectMetrics(Set<String> executorNames) {
//判断是否开启指标监控
if (!dtpProperties.isEnabledCollect()) {
return;
}
//遍历线程池,收集指标
executorNames.forEach(x -> {
ExecutorWrapper wrapper = DtpRegistry.getExecutorWrapper(x);
//遍历收集器,发送指标
doCollect(ExecutorConverter.toMetrics(wrapper));
});
publishCollectEvent();
}
private void doCollect(ThreadPoolStats threadPoolStats) {
try {
CollectorHandler.getInstance().collect(threadPoolStats, this.dtpProperties.getCollectorTypes());
} catch (Exception e) {
log.error("DynamicTp monitor, metrics collect error.", e);
}
}
private void publishCollectEvent() {
CollectEvent event = new CollectEvent(this, this.dtpProperties);
EventBusManager.post(event);
}
private void publishAlarmCheckEvent() {
AlarmCheckEvent event = new AlarmCheckEvent(this, this.dtpProperties);
EventBusManager.post(event);
}
public static void destroy() {
monitorExecutor.shutdownNow();
}
}
//指标收集处理器
public final class CollectorHandler {
private static final Map<String, MetricsCollector> COLLECTORS = Maps.newHashMap();
private CollectorHandler() {
//动态加载指标收集器
List<MetricsCollector> loadedCollectors = ExtensionServiceLoader.get(MetricsCollector.class);
loadedCollectors.forEach(collector -> COLLECTORS.put(collector.type().toLowerCase(), collector));
//自带收集器
MetricsCollector microMeterCollector = new MicroMeterCollector();
LogCollector logCollector = new LogCollector();
InternalLogCollector internalLogCollector = new InternalLogCollector();
JMXCollector jmxCollector = new JMXCollector();
COLLECTORS.put(microMeterCollector.type(), microMeterCollector);
COLLECTORS.put(logCollector.type(), logCollector);
COLLECTORS.put(internalLogCollector.type(), internalLogCollector);
COLLECTORS.put(jmxCollector.type(), jmxCollector);
}
public void collect(ThreadPoolStats poolStats, List<String> types) {
if (poolStats == null || CollectionUtils.isEmpty(types)) {
return;
}
//遍历收集器,发送指标
for (String collectorType : types) {
MetricsCollector collector = COLLECTORS.get(collectorType.toLowerCase());
if (collector != null) {
collector.collect(poolStats);
}
}
}
public static CollectorHandler getInstance() {
return CollectorHandlerHolder.INSTANCE;
}
//静态内部类(单例)
private static class CollectorHandlerHolder {
private static final CollectorHandler INSTANCE = new CollectorHandler();
}
}
实战价值解读:
-
定时任务设计:使用单线程定时任务池(newSingleThreadScheduledExecutor),避免并发采集多个线程池指标导致的锁竞争、性能损耗——生产中,线程池数量可能较多(如10+),单线程采集可确保采集逻辑有序执行,且不占用过多系统资源;
-
指标采集逻辑:所有指标均从 DtpExecutor 和自定义队列中直接获取,无额外 IO 操作;
-
容错设计:采集逻辑中添加全局异常捕获,采集失败仅打印日志,不影响线程池正常运行、不影响告警判断——生产中若监控采集异常,不能导致线程池不可用,这是企业级框架的基本要求;
2.2 关键补充:任务拒绝数的采集原理(自定义计数器)
org.dromara.dynamictp.core.support.ExecutorWrapper#ExecutorWrapper(org.dromara.dynamictp.core.executor.DtpExecutor)
将DtpExecutor封装成ExecutorWrapper,注册到DtpRegistry
public ExecutorWrapper(DtpExecutor executor) {
this.executor = executor;
this.threadPoolName = executor.getThreadPoolName();
this.threadPoolAliasName = executor.getThreadPoolAliasName();
this.notifyItems = executor.getNotifyItems();
this.notifyEnabled = executor.isNotifyEnabled();
this.platformIds = executor.getPlatformIds();
this.awareNames = executor.getAwareNames();
this.rejectEnhanced = executor.isRejectEnhanced();
this.waitForTasksToCompleteOnShutdown = executor.isWaitForTasksToCompleteOnShutdown();
this.awaitTerminationSeconds = executor.getAwaitTerminationSeconds();
//线程池统计信息
this.threadPoolStatProvider = ThreadPoolStatProvider.of(this);
}
public class ThreadPoolStatProvider {
private final ExecutorWrapper executorWrapper;
//统计任务拒绝数
private final LongAdder rejectCount = new LongAdder();
private final PerformanceProvider performanceProvider = new PerformanceProvider();
private ThreadPoolStatProvider(ExecutorWrapper executorWrapper) {
this.executorWrapper = executorWrapper;
}
public static ThreadPoolStatProvider of(ExecutorWrapper executorWrapper) {
val provider = new ThreadPoolStatProvider(executorWrapper);
if (executorWrapper.isDtpExecutor()) {
DtpExecutor dtpExecutor = (DtpExecutor) executorWrapper.getExecutor();
provider.setRunTimeout(dtpExecutor.getRunTimeout());
provider.setQueueTimeout(dtpExecutor.getQueueTimeout());
provider.setTryInterrupt(dtpExecutor.isTryInterrupt());
}
return provider;
}
public ExecutorWrapper getExecutorWrapper() {
return executorWrapper;
}
//获取任务拒绝数
public long getRejectedTaskCount() {
return rejectCount.sum();
}
//增加任务拒绝数
public void incRejectCount(int count) {
rejectCount.add(count);
}
public PerformanceProvider getPerformanceProvider() {
return this.performanceProvider;
}
}
org.dromara.dynamictp.core.reject.RejectedInvocationHandler#invoke
public class RejectedInvocationHandler implements InvocationHandler {
private final Object target;
public RejectedInvocationHandler(Object target) {
this.target = target;
}
@Override
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
//拒绝前
beforeReject((Runnable) args[0], (Executor) args[1]);
try {
//执行拒绝策略
return method.invoke(target, args);
} catch (InvocationTargetException ex) {
throw ex.getCause();
} finally {
//拒绝后
afterReject((Runnable) args[0], (Executor) args[1]);
}
}
/**
* Do sth before reject.
*
* @param runnable the runnable
* @param executor ThreadPoolExecutor instance
*/
private void beforeReject(Runnable runnable, Executor executor) {
AwareManager.beforeReject(runnable, executor);
}
/**
* Do sth after reject.
*
* @param runnable the runnable
* @param executor ThreadPoolExecutor instance
*/
private void afterReject(Runnable runnable, Executor executor) {
AwareManager.afterReject(runnable, executor);
}
}
public class TaskRejectAware extends TaskStatAware {
@Override
public int getOrder() {
return AwareTypeEnum.TASK_REJECT_AWARE.getOrder();
}
@Override
public String getName() {
return AwareTypeEnum.TASK_REJECT_AWARE.getName();
}
@Override
public void beforeReject(Runnable runnable, Executor executor) {
ThreadPoolStatProvider statProvider = statProviders.get(executor);
if (Objects.isNull(statProvider)) {
return;
}
//拒绝任务数 +1
statProvider.incRejectCount(1);
//异步触发报警
AlarmManager.tryAlarmAsync(statProvider.getExecutorWrapper(), REJECT, runnable);
ExecutorAdapter<?> executorAdapter = statProvider.getExecutorWrapper().getExecutor();
String logMsg = CharSequenceUtil.format("DynamicTp execute, thread pool is exhausted, tpName: {}, traceId: {}, " +
"poolSize: {} (active: {}, core: {}, max: {}, largest: {}), " +
"task: {} (completed: {}), queueCapacity: {}, (currSize: {}, remaining: {}) ," +
"executorStatus: (isShutdown: {}, isTerminated: {}, isTerminating: {})",
statProvider.getExecutorWrapper().getThreadPoolName(), MDC.get(TRACE_ID), executorAdapter.getPoolSize(),
executorAdapter.getActiveCount(), executorAdapter.getCorePoolSize(), executorAdapter.getMaximumPoolSize(),
executorAdapter.getLargestPoolSize(), executorAdapter.getTaskCount(), executorAdapter.getCompletedTaskCount(),
statProvider.getExecutorWrapper().getExecutor().getQueueCapacity(), executorAdapter.getQueue().size(),
executorAdapter.getQueue().remainingCapacity(),
executorAdapter.isShutdown(), executorAdapter.isTerminated(), executorAdapter.isTerminating());
log.warn(logMsg);
}
}
三、源码拆解:告警推送核心原理(AlarmManager 主导)
指标采集完成后,下一步就是“告警判断与推送”——这是监控告警模块的“核心价值体现”。DynamicTP 的告警推送核心是 AlarmManager 类,配合告警规则配置、告警渠道实现类,完成“指标判断→告警封装→渠道推送”的全流程。
我们从「告警规则配置→告警判断逻辑→告警渠道推送」三个维度,拆解源码,结合生产中最常用的钉钉推送,补充实战配置。
3.1 前置:告警规则配置示例
告警规则(阈值、渠道、频率)均支持在配置中心动态配置,无需重启服务——先看生产中常见的告警配置,后续源码拆解均围绕该配置展开:
DynamicTP 监控告警配置(Nacos 中配置,对应 DtpProperties 类)
dynamictp:
enabled: true # 是否启用 dynamictp,默认true
enabledCollect: true # 是否开启监控指标采集,默认true
collectorTypes: micrometer,logging # 监控数据采集器类型(logging | micrometer | internal_logging | JMX),默认micrometer
logPath: /home/logs/dynamictp/user–center/ # 监控日志数据路径,默认 ${user.home}/logs,采集类型非logging不用配置
monitorInterval: 5 # 监控时间间隔(报警检测、指标采集),默认5s
# 告警渠道
platforms: # 通知报警平台配置
– platform: wechat
platformId: 1 # 平台id,自定义
urlKey: 3a700–127–4bd–a798–c53d8b69c # webhook 中的 key
receivers: test1,test2 # 接受人企微账号
– platform: ding
platformId: 2 # 平台id,自定义
urlKey: f80dad441fcd655438f4a08dcd6a # webhook 中的 access_token
secret: SECb5441fa6f375d5b9d21 # 安全设置在验签模式下才的秘钥,非验签模式没有此值
receivers: 18888888888 # 钉钉账号手机号
– platform: lark
platformId: 3
urlKey: 0d944ae7–b24a–40 # webhook 中的 token
secret: 3a750012874bdac5c3d8b69c # 安全设置在签名校验模式下才的秘钥,非验签模式没有此值
receivers: test1,test2 # 接受人username / openid
– platform: email
platformId: 4
receivers: 123456@qq.com,789789@qq.com # 收件人邮箱,多个用逗号隔开
# 全局配置
globalExecutorProps: # 线程池配置 > 全局配置 > 字段默认值
rejectedHandlerType: CallerRunsPolicy
queueType: VariableLinkedBlockingQueue
waitForTasksToCompleteOnShutdown: true
awaitTerminationSeconds: 3
taskWrapperNames: ["swTrace", "ttl", "mdc"]
queueTimeout: 300
runTimeout: 300
notifyItems: # 报警项,不配置自动会按默认值配置(变更通知、容量报警、活性报警、拒绝报警、任务超时报警)
– type: change # 线程池核心参数变更通知
silencePeriod: 120 # 通知静默时间(单位:s),默认值1,0表示不静默
– type: capacity # 队列容量使用率,报警项类型,查看源码 NotifyTypeEnum枚举类
threshold: 80 # 报警阈值,意思是队列使用率达到70%告警;默认值=70
count: 2 # 在一个统计周期内,如果触发阈值的数量达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默,默认值=120
– type: liveness # 线程池活性
threshold: 80 # 报警阈值,意思是活性达到70%告警;默认值=70
count: 3 # 在一个统计周期内,如果触发阈值的数量达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120
– type: reject # 触发任务拒绝告警
count: 1 # 在一个统计周期内,如果触发拒绝策略次数达到 count,则触发报警;默认值=1
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120
– type: run_timeout # 任务执行超时告警
count: 20 # 在一个统计周期内,如果执行超时次数达到 count,则触发报警;默认值=10
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 30 # 报警静默时间(单位:s),0表示不静默;默认值=120
– type: queue_timeout # 任务排队超时告警
count: 5 # 在一个统计周期内,如果排队超时次数达到 count,则触发报警;默认值=10
period: 30 # 报警统计周期(单位:s),默认值=120
silencePeriod: 0 # 报警静默时间(单位:s),0表示不静默;默认值=120
# 线程池配置
executors: # 动态线程池配置,都有默认值,采用默认值的可以不配置该项,减少配置量
– threadPoolName: dtpExecutor1 # 线程池名称,必填
threadPoolAliasName: 测试线程池 # 线程池别名,可选
executorType: common # 线程池类型 common、eager、ordered、scheduled、priority,默认 common
corePoolSize: 6 # 核心线程数,默认1
maximumPoolSize: 8 # 最大线程数,默认cpu核数
queueCapacity: 2000 # 队列容量,默认1024
queueType: VariableLinkedBlockingQueue # 任务队列,查看源码QueueTypeEnum枚举类,默认VariableLinkedBlockingQueue
rejectedHandlerType: CallerRunsPolicy # 拒绝策略,查看RejectedTypeEnum枚举类,默认AbortPolicy
keepAliveTime: 60 # 空闲线程等待超时时间,默认60
threadNamePrefix: test # 线程名前缀,默认dtp
allowCoreThreadTimeOut: false # 是否允许核心线程池超时,默认false
waitForTasksToCompleteOnShutdown: true # 参考spring线程池设计,优雅关闭线程池,默认true
awaitTerminationSeconds: 5 # 优雅关闭线程池时,阻塞等待线程池中任务执行时间,默认3,单位(s)
preStartAllCoreThreads: false # 是否预热所有核心线程,默认false
runTimeout: 200 # 任务执行超时阈值,单位(ms),默认0(不统计)
queueTimeout: 100 # 任务在队列等待超时阈值,单位(ms),默认0(不统计)
tryInterrupt: false # 执行超时后是否中断线程,默认false
taskWrapperNames: ["ttl", "mdc"] # 任务包装器名称,继承TaskWrapper接口
notifyEnabled: true # 是否开启报警,默认true
platformIds: [1,2] # 报警平台id,不配置默认拿上层platforms配置的所有平台
配置解读:
-
告警频率限制:monitorInterval ,监控时间间隔(报警检测、指标采集),默认5s;
-
告警类型:支持 REJECT、RUN_TIMEOUT、QUEUE_TIMEOUT等多种类型,覆盖生产中常见异常场景;
-
多渠道适配:可同时配置多个告警渠道(钉钉、企业微信),确保告警信息能及时触达(生产中建议至少配置两个渠道,避免单一渠道故障导致告警丢失)。
3.2 核心类 1:AlarmManager(告警管理器)
AlarmManager 的核心作用:加载告警配置、接收 DtpMonitor 采集的指标、判断指标是否触发告警阈值、触发告警并推送给对应的渠道,是告警推送的“调度中心”。
源码简化(聚焦核心逻辑:告警判断、告警触发,剔除冗余的配置校验代码)
public class AlarmManager {
//单线程的线程池
private static final ExecutorService ALARM_EXECUTOR = ThreadPoolBuilder.newBuilder()
.threadFactory("dtp-alarm")
.corePoolSize(1)
.maximumPoolSize(1)
.workQueue(LINKED_BLOCKING_QUEUE.getName(), 2000)
.rejectedExecutionHandler(RejectedTypeEnum.DISCARD_OLDEST_POLICY.getName())
.rejectEnhanced(false)
.taskWrappers(TaskWrappers.getInstance().getByNames(Sets.newHashSet("mdc")))
.buildDynamic();
private static final InvokerChain<BaseNotifyCtx> ALARM_INVOKER_CHAIN;
static {
//报警执行链
ALARM_INVOKER_CHAIN = NotifyFilterBuilder.getAlarmInvokerChain();
}
private AlarmManager() { }
public static void initAlarm(String poolName, List<NotifyItem> notifyItems) {
notifyItems.forEach(x -> initAlarm(poolName, x));
}
public static void initAlarm(String poolName, NotifyItem notifyItem) {
//静默检测
AlarmLimiter.initAlarmLimiter(poolName, notifyItem);
//次数统计,比如拒绝任务数
AlarmCounter.initAlarmCounter(poolName, notifyItem);
}
public static void initAlarmLimiter(String poolName, NotifyItem notifyItem) {
AlarmLimiter.initAlarmLimiter(poolName, notifyItem);
}
public static void initAlarmCounter(String poolName, NotifyItem notifyItem) {
AlarmCounter.initAlarmCounter(poolName, notifyItem);
}
public static void tryAlarmAsync(ExecutorWrapper executorWrapper, NotifyItemEnum notifyType, Runnable runnable) {
preAlarm(runnable);
try {
ALARM_EXECUTOR.execute(() -> doTryAlarm(executorWrapper, notifyType));
} finally {
postAlarm(runnable);
}
}
//检测、异步报警 (调用地址:org.dromara.dynamictp.core.monitor.DtpMonitor#run)
public static void checkAndTryAlarmAsync(ExecutorWrapper executorWrapper, List<NotifyItemEnum> notifyTypes) {
//循环遍历线程池下的报警类型
ALARM_EXECUTOR.execute(() -> notifyTypes.forEach(x -> doCheckAndTryAlarm(executorWrapper, x)));
}
public static void doCheckAndTryAlarm(ExecutorWrapper executorWrapper, NotifyItemEnum notifyType) {
NotifyHelper.getNotifyItem(executorWrapper, notifyType).ifPresent(notifyItem -> {
//检测是否达到报警阈值
if (hasReachedThreshold(executorWrapper, notifyType, notifyItem)) {
//达到阈值,发送报警通知
ALARM_INVOKER_CHAIN.proceed(new AlarmCtx(executorWrapper, notifyItem));
}
});
}
public static void tryAlarmAsync(ExecutorWrapper executorWrapper, List<NotifyItemEnum> notifyTypes) {
ALARM_EXECUTOR.execute(() -> notifyTypes.forEach(x -> doTryAlarm(executorWrapper, x)));
}
public static void doTryAlarm(ExecutorWrapper executorWrapper, NotifyItemEnum notifyType) {
NotifyHelper.getNotifyItem(executorWrapper, notifyType).ifPresent(notifyItem -> {
val alarmCtx = new AlarmCtx(executorWrapper, notifyItem);
ALARM_INVOKER_CHAIN.proceed(alarmCtx);
});
}
private static void preAlarm(Runnable runnable) {
if (runnable instanceof DtpRunnable) {
MDC.put(TRACE_ID, ((DtpRunnable) runnable).getTraceId());
}
}
private static void postAlarm(Runnable runnable) {
if (runnable instanceof DtpRunnable) {
MDC.remove(TRACE_ID);
}
}
/**
* 检测是否达阈值(任务堆积百分比、活跃线程百分比)
*/
private static boolean hasReachedThreshold(ExecutorWrapper executor, NotifyItemEnum notifyType, NotifyItem notifyItem) {
switch (notifyType) {
case CAPACITY:
return checkCapacity(executor, notifyItem);
case LIVENESS:
return checkLiveness(executor, notifyItem);
case REJECT:
case RUN_TIMEOUT:
case QUEUE_TIMEOUT:
return true;
default:
log.error("Unsupported alarm type [{}]", notifyType);
return false;
}
}
//检测活性告警(触发此类报警的原因:任务处理速度跟不上请求速率,导致系统响应变慢。线程池配置不合理,最大线程数过小。)
private static boolean checkLiveness(ExecutorWrapper executorWrapper, NotifyItem notifyItem) {
val executor = executorWrapper.getExecutor();
int maximumPoolSize = executor.getMaximumPoolSize();
double div = NumberUtil.div(executor.getActiveCount(), maximumPoolSize, 2) * 100;
//默认阈值70,线程池活跃线程数量超过70%,报警
if (div >= notifyItem.getThreshold()) {
log.warn("DynamicTp monitor, current liveness [{}] >= threshold [{}], threadPoolName: {}",
div, notifyItem.getThreshold(), executorWrapper.getThreadPoolName());
return true;
}
return false;
}
//检测队列容量使用率
private static boolean checkCapacity(ExecutorWrapper executorWrapper, NotifyItem notifyItem) {
val executor = executorWrapper.getExecutor();
if (executor.getQueueSize() <= 0) {
return false;
}
double div = NumberUtil.div(executor.getQueueSize(), executor.getQueueCapacity(), 2) * 100;
//默认阈值为70,任务堆积超过70%就会报警
if (div >= notifyItem.getThreshold()) {
log.warn("DynamicTp monitor, current queue utilization [{}] >= threshold [{}], threadPoolName: {}",
div, notifyItem.getThreshold(), executorWrapper.getThreadPoolName());
return true;
}
return false;
}
//关闭线程池
public static void destroy() {
ALARM_EXECUTOR.shutdownNow();
}
}
org.dromara.dynamictp.core.notifier.manager.NotifyFilterBuilder#getAlarmInvokerChain
public static InvokerChain<BaseNotifyCtx> getAlarmInvokerChain() {
val filters = ContextManagerHelper.getBeansOfType(NotifyFilter.class);
Collection<NotifyFilter> alarmFilters = Lists.newArrayList(filters.values());
alarmFilters.add(new BaseAlarmFilter());
alarmFilters.add(new SilentCheckFilter());
alarmFilters = alarmFilters.stream()
.filter(x -> x.supports(NotifyTypeEnum.ALARM))
.sorted(Comparator.comparing(Filter::getOrder))
.collect(Collectors.toList());
//构建执行链
return InvokerChainFactory.buildInvokerChain(new AlarmInvoker(), alarmFilters.toArray(new NotifyFilter[0]));
}
org.dromara.dynamictp.core.notifier.chain.invoker.AlarmInvoker#invoke
public void invoke(BaseNotifyCtx context) {
val executorWrapper = context.getExecutorWrapper();
val notifyItem = context.getNotifyItem();
try {
DtpNotifyCtxHolder.set(context);
//发送报警通知
NotifierHandler.getInstance().sendAlarm(NotifyItemEnum.of(notifyItem.getType()));
AlarmCounter.reset(executorWrapper.getThreadPoolName(), notifyItem.getType());
} finally {
DtpNotifyCtxHolder.remove();
}
}
org.dromara.dynamictp.core.handler.NotifierHandler#sendAlarm
public void sendAlarm(NotifyItemEnum notifyItemEnum) {
//报警类型
NotifyItem notifyItem = DtpNotifyCtxHolder.get().getNotifyItem();
//遍历配置的告警渠道,发送报警通知
for (String platformId : notifyItem.getPlatformIds()) {
NotifyHelper.getPlatform(platformId).ifPresent(p -> {
DtpNotifier notifier = NOTIFIERS.get(p.getPlatform().toLowerCase());
if (notifier != null) {
notifier.sendAlarmMsg(p, notifyItemEnum);
}
});
}
}
核心逻辑拆解:
告警判断入口:checkAlarm 方法,由 DtpMonitor 的 run方法调用,周期性的调用,判断是否需要告警;
阈值判断逻辑:hasReachedThreshold方法,根据不同的告警类型,判断对应的指标是否超过阈值;
多渠道推送:sendAlarm 方法,遍历配置的告警渠道,调用对应渠道的 send 方法推送告警——单个渠道推送失败不影响其他渠道,容错性强(生产中,钉钉、企业微信双渠道推送,可避免单一渠道故障导致告警丢失)。
3.3 核心类 2:DtpWechatNotifier (以企业微信为例,最常用)
DynamicTP 采用“接口+实现”的方式,适配多告警渠道(钉钉、企业微信等)——核心是 Notifier接口,不同渠道实现该接口,重写 send 方法,实现个性化推送。
我们重点拆解生产中最常用的企业微信渠道实现,其他渠道(钉钉)逻辑类似,可直接复用思路。
public interface Notifier {
//报警平台名称
String platform();
//推送报警信息
void send(NotifyPlatform platform, String content);
}
public abstract class AbstractNotifier implements Notifier {
@Override
public final void send(NotifyPlatform platform, String content) {
try {
send0(platform, content);
} catch (Exception e) {
log.error("DynamicTp notify, {} send failed.", platform(), e);
}
}
protected abstract void send0(NotifyPlatform platform, String content);
}
public abstract class AbstractHttpNotifier extends AbstractNotifier {
@Override
protected void send0(NotifyPlatform platform, String content) {
val url = buildUrl(platform);
val msgBody = buildMsgBody(platform, content);
HttpRequest request = HttpRequest.post(url)
.setConnectionTimeout(platform.getTimeout())
.setReadTimeout(platform.getTimeout())
.body(msgBody);
if (platform.getProxyType() != Proxy.Type.DIRECT) {
request.setProxy(new Proxy(platform.getProxyType(), new InetSocketAddress(platform.getProxyHost(), platform.getProxyPort())));
}
HttpResponse response = request.execute();
if (Objects.nonNull(response)) {
log.info("DynamicTp notify, {} send success, response: {}, request: {}",
platform(), response.body(), msgBody);
}
}
protected abstract String buildMsgBody(NotifyPlatform platform, String content);
protected abstract String buildUrl(NotifyPlatform platform);
}
/**
* 企业微信告警渠道实现
*/
public class WechatNotifier extends AbstractHttpNotifier {
@Override
public String platform() {
return NotifyPlatformEnum.WECHAT.name().toLowerCase();
}
//推送的消息:配置的变更、告警
@Override
protected String buildMsgBody(NotifyPlatform platform, String content) {
MarkdownReq markdownReq = new MarkdownReq();
markdownReq.setMsgtype("markdown");
MarkdownReq.Markdown markdown = new MarkdownReq.Markdown();
markdown.setContent(content);
markdownReq.setMarkdown(markdown);
return JsonUtil.toJson(markdownReq);
}
//企业微信请求地址
@Override
protected String buildUrl(NotifyPlatform platform) {
if (StringUtils.isBlank(platform.getUrlKey())) {
return platform.getWebhook();
}
UrlBuilder builder = UrlBuilder.of(Optional.ofNullable(platform.getWebhook()).orElse(WechatNotifyConst.WECHAT_WEB_HOOK));
if (StringUtils.isBlank(builder.getQuery().get(WechatNotifyConst.KEY_PARAM))) {
builder.addQuery(WechatNotifyConst.KEY_PARAM, platform.getUrlKey());
}
return builder.build();
}
}
org.dromara.dynamictp.core.notifier.DtpWechatNotifier
public class DtpWechatNotifier extends AbstractDtpNotifier {
//绑定WechatNotifier
public DtpWechatNotifier(Notifier notifier) {
super(notifier);
}
@Override
public String platform() {
return NotifyPlatformEnum.WECHAT.name().toLowerCase();
}
//配置变更 通知模板
@Override
protected String getNoticeTemplate() {
return WechatNotifyConst.WECHAT_CHANGE_NOTICE_TEMPLATE;
}
//告警模板
@Override
protected String getAlarmTemplate() {
return WechatNotifyConst.WECHAT_ALARM_TEMPLATE;
}
@Override
protected Pair<String, String> getColors() {
return new ImmutablePair<>(WechatNotifyConst.WARNING_COLOR, WechatNotifyConst.COMMENT_COLOR);
}
@Override
protected String formatReceivers(String receives) {
return Arrays.stream(StringUtils.split(receives, ','))
.map(receiver -> "<@" + receiver + ">")
.collect(Collectors.joining(","));
}
}
结尾
到这里,DynamicTP监控告警模块的核心源码、设计思路与实战落地细节,就全部拆解完成了。结合前四篇内容,我们已经完整掌握了DynamicTP最核心的两大能力——「动态调参」与「监控告警」,这也是DynamicTP能成为生产级动态线程池框架的核心原因,更是我们在面试中体现“源码能力+生产思维”的关键亮点。
作为一名10年Java开发,这里给大家一个生产落地的核心建议:落地DynamicTP时,先做好监控告警,再谈动态调参。很多小伙伴容易陷入“盲目调参”的误区,殊不知,只有通过监控实时掌握线程池的运行状态(如任务堆积、线程负载),才能精准判断调参方向;同时,务必优先配置“任务拒绝、队列满”两类核心告警,搭配钉钉+企业微信双渠道,确保异常能第一时间触达,避免小问题扩大为生产故障。




