diff --git a/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java b/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java index ed74589..3951463 100644 --- a/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java +++ b/src/main/java/com/knowledge/base/infrastructure/config/DynamicConfig.java @@ -1,24 +1,30 @@ package com.knowledge.base.infrastructure.config; import cn.hutool.core.util.StrUtil; +import com.knowledge.base.infrastructure.monitor.ThreadPoolMonitorStarter; import lombok.Getter; +import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.stereotype.Component; +import javax.annotation.PostConstruct; import java.util.Arrays; import java.util.Set; import java.util.stream.Collectors; /** - * @author Luke.ye - * @date 2025/5/8 09:59 + * 动态配置类,支持 Nacos 热更新 + * + * @author Luke + * @date 2025/5/8 */ +@Slf4j +@Getter @Component @RefreshScope -@Getter public class DynamicConfig { - //是否记录请求和响应信息 Y-记录 N-不记录 默认记录 + @Value("${micro.saas.doc.parser.recordMsgBody:Y}") private String recordMsgBody; @@ -28,7 +34,6 @@ public class DynamicConfig { @Value("${import.schedule.cron:* 0/30 * * * ?}") private String importScheduleCron; - // token失效时间,默认24小时 @Value("${token.expire.time:86400000}") private long tokenExpireTime; @@ -59,11 +64,41 @@ public class DynamicConfig { @Value("${exclude.file.path.prefix}") private String mdExcludePrefix; + @Value("${thread.pool.monitor.interval.seconds:60}") + private long threadPoolMonitorIntervalSeconds; + + private boolean enableThreadPoolMonitor; + + /** + * 来自 nacos 的开关控制是否启用线程池监控 + */ + @Value("${thread.pool.monitor.enabled:true}") + public void setEnableThreadPoolMonitor(boolean enabled) { + boolean changed = this.enableThreadPoolMonitor != enabled; + this.enableThreadPoolMonitor = enabled; + if (changed) { + if (enabled) { + log.info("[config-refresh] enableThreadPoolMonitor=true,刷新线程池监控"); + ThreadPoolMonitorStarter.getInstance().refreshAll(); + } else { + log.info("[config-refresh] enableThreadPoolMonitor=false,清除线程池监控"); + ThreadPoolMonitorStarter.getInstance().clearAll(); + } + } + } + public Set getLlmActiveSlugs() { - // 支持逗号、分号和换行分隔 return Arrays.stream(llmSharedActiveSlugIds.split("[,;\\n]")) .map(String::trim) .filter(StrUtil::isNotBlank) .collect(Collectors.toSet()); } + + @PostConstruct + public void init() { + log.info("[config-init] DynamicConfig 初始化完成,线程池监控开关 enableThreadPoolMonitor={}", enableThreadPoolMonitor); + if (enableThreadPoolMonitor) { + ThreadPoolMonitorStarter.getInstance().refreshAll(); + } + } } diff --git a/src/main/java/com/knowledge/base/infrastructure/monitor/ThreadPoolMonitorStarter.java b/src/main/java/com/knowledge/base/infrastructure/monitor/ThreadPoolMonitorStarter.java new file mode 100644 index 0000000..45a28c9 --- /dev/null +++ b/src/main/java/com/knowledge/base/infrastructure/monitor/ThreadPoolMonitorStarter.java @@ -0,0 +1,87 @@ +package com.knowledge.base.infrastructure.monitor; + +import cn.hutool.extra.spring.SpringUtil; +import com.knowledge.base.infrastructure.config.DynamicConfig; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import javax.annotation.PreDestroy; +import java.util.Map; +import java.util.concurrent.*; + +/** + * 线程池监控调度器(支持动态配置控制、自动刷新、清理) + */ +@Slf4j +@Component +@RequiredArgsConstructor +public class ThreadPoolMonitorStarter { + + private final DynamicConfig dynamicConfig; + + private final ScheduledExecutorService monitorScheduler = Executors.newSingleThreadScheduledExecutor( + r -> new Thread(r, "thread-monitor-scheduler")); + + private final Map> monitorTasks = new ConcurrentHashMap<>(); + + private final Map registeredPools = new ConcurrentHashMap<>(); + + public void register(String poolName, ExecutorService pool) { + registeredPools.put(poolName, pool); + refresh(poolName, pool); + } + + public void refreshAll() { + clearAll(); + for (Map.Entry entry : registeredPools.entrySet()) { + refresh(entry.getKey(), entry.getValue()); + } + } + + public void refresh(String poolName, ExecutorService pool) { + if (!"Y".equalsIgnoreCase(dynamicConfig.getRecordMsgBody())) { + log.info("[thread-monitor] 未开启线程池监控配置,跳过 {}", poolName); + return; + } + + if (!(pool instanceof ThreadPoolExecutor)) { + log.info("[thread-monitor] {} 非 ThreadPoolExecutor,无法监控", poolName); + return; + } + + ThreadPoolExecutor executor = (ThreadPoolExecutor) pool; + + ScheduledFuture future = monitorScheduler.scheduleAtFixedRate(() -> { + log.info("[thread-monitor] {} - 活跃线程数: {}, 最大线程数: {}, 核心线程数: {}, 排队任务数: {}, 总任务数: {}, 已完成任务数: {}", + poolName, + executor.getActiveCount(), + executor.getMaximumPoolSize(), + executor.getCorePoolSize(), + executor.getQueue().size(), + executor.getTaskCount(), + executor.getCompletedTaskCount()); + }, 0, 1, TimeUnit.MINUTES); + + monitorTasks.put(poolName, future); + log.info("[thread-monitor] 线程池 {} 监控任务已启动", poolName); + } + + public void clearAll() { + monitorTasks.forEach((name, future) -> { + future.cancel(true); + log.info("[thread-monitor] 已取消线程池 {} 的监控任务", name); + }); + monitorTasks.clear(); + } + + @PreDestroy + public void destroy() { + clearAll(); + monitorScheduler.shutdownNow(); + } + + public static ThreadPoolMonitorStarter getInstance() { + return SpringUtil.getBean(ThreadPoolMonitorStarter.class); + } +} diff --git a/src/main/java/com/knowledge/base/infrastructure/util/ThreadPoolUtil.java b/src/main/java/com/knowledge/base/infrastructure/util/ThreadPoolUtil.java index c305446..10c67fb 100644 --- a/src/main/java/com/knowledge/base/infrastructure/util/ThreadPoolUtil.java +++ b/src/main/java/com/knowledge/base/infrastructure/util/ThreadPoolUtil.java @@ -1,14 +1,16 @@ package com.knowledge.base.infrastructure.util; import cn.hutool.core.thread.ThreadFactoryBuilder; +import com.knowledge.base.infrastructure.monitor.ThreadPoolMonitorStarter; +import lombok.extern.slf4j.Slf4j; import java.util.concurrent.*; /** * 通用线程池工具类 - * 支持外部自定义线程池传入,未传入时使用默认线程池 - * @author Luke + * 支持外部自定义线程池执行任务,并可选是否注册监控 */ +@Slf4j public class ThreadPoolUtil { private static final int CORE_POOL_SIZE = Runtime.getRuntime().availableProcessors() + 1; @@ -16,7 +18,6 @@ public class ThreadPoolUtil { private static final int QUEUE_CAPACITY = 500; private static final long KEEP_ALIVE_TIME = 60L; - // 默认线程池 private static final ThreadPoolExecutor DEFAULT_THREAD_POOL = new ThreadPoolExecutor( CORE_POOL_SIZE, MAX_POOL_SIZE, @@ -28,35 +29,52 @@ public class ThreadPoolUtil { ); /** - * 执行任务,使用默认线程池 + * 提交默认线程池任务 */ public static void execute(Runnable task) { - DEFAULT_THREAD_POOL.execute(task); + DEFAULT_THREAD_POOL.execute(wrap(task, "default")); } /** - * 执行任务,允许调用方传入自定义线程池 - * @param task Runnable - * @param executor 若为 null,则用默认线程池 + * 提交自定义线程池任务(默认不监控) */ public static void execute(Runnable task, ExecutorService executor) { - if (executor != null) { - executor.execute(task); - } else { - DEFAULT_THREAD_POOL.execute(task); + execute(task, executor, false, "custom"); + } + + /** + * 提交任务(带监控选项 + 线程池名称) + */ + public static void execute(Runnable task, ExecutorService executor, boolean monitor, String poolName) { + if (executor == null) { + DEFAULT_THREAD_POOL.execute(wrap(task, "default")); + return; + } + executor.execute(wrap(task, poolName)); + if (monitor) { + ThreadPoolMonitorStarter.getInstance().register(poolName, executor); } } - /** - * 优雅关闭(默认线程池) - */ + private static Runnable wrap(Runnable task, String poolName) { + return () -> { + String threadName = Thread.currentThread().getName(); + try { + log.debug("[thread-pool][{}] 执行任务开始", poolName); + task.run(); + log.debug("[thread-pool][{}] 执行任务结束", poolName); + } catch (Exception e) { + log.error("[thread-pool][{}] 执行异常", poolName, e); + } finally { + Thread.currentThread().setName(threadName); // 防止线程池复用导致名称混乱 + } + }; + } + public static void shutdownAndAwait() { shutdownAndAwait(DEFAULT_THREAD_POOL); } - /** - * 优雅关闭(指定线程池) - */ public static void shutdownAndAwait(ExecutorService executor) { if (executor == null) return; executor.shutdown(); @@ -70,9 +88,6 @@ public class ThreadPoolUtil { } } - /** - * 获取默认线程池(如需提交批量任务) - */ public static ExecutorService getDefaultThreadPool() { return DEFAULT_THREAD_POOL; }