欢迎光临
我们一直在努力

深度解析|DynamicTP 源码解析系列(五):监控告警核心原理(指标采集+告警推送)

前言

在上一篇《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/usercenter/ # 监控日志数据路径,默认 ${user.home}/logs,采集类型非logging不用配置
    monitorInterval: 5 # 监控时间间隔(报警检测、指标采集),默认5s

    # 告警渠道
    platforms: # 通知报警平台配置
    platform: wechat
    platformId: 1 # 平台id,自定义
    urlKey: 3a7001274bda798c53d8b69c # webhook 中的 key
    receivers: test1,test2 # 接受人企微账号

    platform: ding
    platformId: 2 # 平台id,自定义
    urlKey: f80dad441fcd655438f4a08dcd6a # webhook 中的 access_token
    secret: SECb5441fa6f375d5b9d21 # 安全设置在验签模式下才的秘钥,非验签模式没有此值
    receivers: 18888888888 # 钉钉账号手机号

    platform: lark
    platformId: 3
    urlKey: 0d944ae7b24a40 # 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时,先做好监控告警,再谈动态调参。很多小伙伴容易陷入“盲目调参”的误区,殊不知,只有通过监控实时掌握线程池的运行状态(如任务堆积、线程负载),才能精准判断调参方向;同时,务必优先配置“任务拒绝、队列满”两类核心告警,搭配钉钉+企业微信双渠道,确保异常能第一时间触达,避免小问题扩大为生产故障。

    赞(0)
    未经允许不得转载:171主机测评 » 深度解析|DynamicTP 源码解析系列(五):监控告警核心原理(指标采集+告警推送)
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址