尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

如何自定义线程池拒绝策略?

如何自定义线程池拒绝策略? 在生产环境中默认的AbortPolicy抛出异常可能会导致局部业务中断而DiscardPolicy则会导致任务静默丢失。通过自定义拒绝策略RejectedExecutionHandler我们可以将触发拒绝策略的任务暂存落盘磁盘/日志/MQ/Redis进行事后补偿或实时发送钉钉告警通知从而实现高可用与可观测性。下面提供一个标准的、包含任务落盘与钉钉告警功能的自定义拒绝策略实现。自定义拒绝策略的核心是通过实现RejectedExecutionHandler接口在任务触发拒绝机制时捕获线程池上下文与任务信息进而执行定制的降级、补偿与告警逻辑。一、 自定义拒绝策略的标准实现步骤实现接口编写类实现java.util.concurrent.RejectedExecutionHandler接口。重写核心方法实现rejectedExecution(Runnable r, ThreadPoolExecutor executor)方法。提取运行上下文利用executor参数获取活跃线程数getActiveCount()、当前队列积压数getQueue().size()等实时状态。执行兜底逻辑编写具体的业务降级行为如写入本地磁盘/日志表、推入 MQ/Redis、触发钉钉/企微 HTTP 告警等。注入线程池在创建ThreadPoolExecutor时将自定义策略实例作为第 7 个参数传入构造函数。二、 核心代码实现1. 自定义拒绝策略实现类importjava.io.FileWriter;importjava.io.IOException;importjava.io.PrintWriter;importjava.time.LocalDateTime;importjava.time.format.DateTimeFormatter;importjava.util.concurrent.RejectedExecutionHandler;importjava.util.concurrent.ThreadPoolExecutor;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;/** * 自定义线程池拒绝策略 * 1. 记录日志与指标监控 * 2. 任务落盘持久化到本地文件方便后续补偿 * 3. 实时触发告警如发送钉钉通知 */publicclassFallbackAndNotifyRejectedPolicyimplementsRejectedExecutionHandler{privatestaticfinalLoggerlogLoggerFactory.getLogger(FallbackAndNotifyRejectedPolicy.class);// 线程池名称方便标识来源privatefinalStringthreadPoolName;// 本地任务备份文件路径privatefinalStringbackupFilePath;publicFallbackAndNotifyRejectedPolicy(StringthreadPoolName,StringbackupFilePath){this.threadPoolNamethreadPoolName;this.backupFilePathbackupFilePath;}OverridepublicvoidrejectedExecution(Runnabler,ThreadPoolExecutorexecutor){StringtimestampLocalDateTime.now().format(DateTimeFormatter.ofPattern(yyyy-MM-dd HH:mm:ss.SSS));// 1. 收集线程池运行时的关键指标信息StringpoolInfoString.format([%s] 线程池 [%s] 触发拒绝策略ActiveCount: %d, PoolSize: %d, QueueSize: %d, TaskCount: %long,timestamp,threadPoolName,executor.getActiveCount(),executor.getPoolSize(),executor.getQueue().size(),executor.getTaskCount());// 打印 ERROR 日志log.error({}. 任务内容: {},poolInfo,r.toString());// 2. 任务落盘持久化异步/同步写入本地文件或专用日志表persistTaskToDisk(timestamp,r);// 3. 触发钉钉告警通知sendDingTalkAlert(poolInfo,r);}/** * 将被拒绝的任务详情落盘可用于后续离线补偿脚本读取 */privatevoidpersistTaskToDisk(Stringtimestamp,Runnabler){try(PrintWriterwriternewPrintWriter(newFileWriter(backupFilePath,true))){writer.printf(%s | ThreadPool: %s | Task: %s%n,timestamp,threadPoolName,r.toString());writer.flush();}catch(IOExceptione){log.error(将被拒绝的任务写入磁盘失败,e);}}/** * 发送钉钉机器人告警简单 HTTP POST 请求 */privatevoidsendDingTalkAlert(StringpoolInfo,Runnabler){// 建议异步调用避免阻碍调用线程此处展示核心逻辑newThread(()-{try{// 拼接钉钉 Markdown 告警消息StringmarkdownTextString.format(### ⚠️ 线程池高负载拒绝告警\n- **线程池名称**: %s\n- **告警详情**: %s\n- **被拒绝任务**: %s\n- **建议动作**: 请及时检查流量突增或线程池参数设置,threadPoolName,poolInfo,r.toString());// 实际项目中替换为真实的钉钉 Webhook URLStringwebhookUrlhttps://oapi.dingtalk.com/robot/send?access_tokenYOUR_ACCESS_TOKEN;// 此处调用你的 HTTP Client如 RestTemplate / HttpClient / OkHttp发送请求// HttpClientUtil.postJson(webhookUrl, buildDingTalkJsonPayload(markdownText));log.info(【告警通知已触发】推送至钉钉: {},markdownText);}catch(Exceptione){log.error(发送钉钉告警失败,e);}},DingTalk-Alert-Thread).start();}}三、 在线程池中的配置与使用示例在创建ThreadPoolExecutor时将自定义的策略作为第 7 个参数传入importjava.util.concurrent.ArrayBlockingQueue;importjava.util.concurrent.ThreadPoolExecutor;importjava.util.concurrent.TimeUnit;publicclassThreadPoolDemo{publicstaticvoidmain(String[]args){// 创建带有自定义拒绝策略的线程池ThreadPoolExecutorexecutornewThreadPoolExecutor(2,// corePoolSize2,// maximumPoolSize60L,TimeUnit.SECONDS,// keepAliveTimenewArrayBlockingQueue(2),// 队列容量只有 2newCustomThreadFactory(order-service-pool),// 规范命名线程工厂newFallbackAndNotifyRejectedPolicy(Order-Service-Pool,/data/logs/rejected_tasks.log)// 自定义拒绝策略);// 模拟提交 10 个任务以触发拒绝策略for(inti1;i10;i){finalinttaskIdi;executor.execute(newRunnable(){Overridepublicvoidrun(){try{Thread.sleep(2000);// 模拟任务耗时}catch(InterruptedExceptione){Thread.currentThread().interrupt();}}OverridepublicStringtoString(){returnOrderProcessingTask-#taskId;}});}executor.shutdown();}}三、 生产落地最佳实践防告警轰炸限流/频率控制在高并发触发拒绝策略时短时间内可能会抛出成千上万条拒绝事件。必须在sendDingTalkAlert中加本地缓存限流如令牌桶或滑动窗口例如设置“1 分钟内最多只发送一次告警”避免打爆钉钉 API 或造成告警风暴。异步落盘 / 兜底暂存本地文件适合最底层的降级保底配合定时补偿脚本重新读取并投递。Redis / MQ如果业务非常关键如支付、下单可以尝试将任务参数 Serialization 后推入 Redis 队列或死信队列DLQ。保证 Runnable 的可可读性默认的Runnable对象的toString()只输出类名和内存地址。最好实现自己的NamedRunnable或重写toString()方法打印出任务 ID、业务参数如 orderId方便落盘后排查与人工补单。
返回列表