REST HTTP Web API场景事件发送实现及高并发扩容方案咨询
功能改造实现方案
需求拆解
- 支持包含6个事件的单场景按自定义每小时发送频次x、自定义持续时长y循环发送
- 最高承载每小时19000个场景的发送量,最长支持连续72小时(3天)稳定运行
- 替换现有单次触发单场景发送的实现逻辑
现有代码核心问题
现有实现存在3个明显缺陷,完全不满足高并发长稳运行要求:
- 每次请求都新建固定5线程的线程池,频繁创建销毁资源,高并发场景下会直接导致OOM
- 仅支持单次提交单场景,无周期调度、运行时长控制能力
- 无任务追踪、限流、重试机制,长时间运行容易出现任务丢失、线程卡死等问题
具体改造步骤
1. 全局线程池架构替换
删除接口内新建线程池的逻辑,改用全局共享的「调度线程池+工作线程池」分离架构,Spring项目可直接注册为Bean全局管理:
// 调度线程池:仅负责周期触发任务,不需要太多线程,10个足够应对最高并发调度需求 private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(10); // 工作线程池:实际执行发送逻辑,按最高每小时19000场景计算,每秒约5.3个场景,20个线程完全足够,搭配有界队列防止溢出 private final ThreadPoolExecutor workerPool = new ThreadPoolExecutor( 10, 20, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(1000), new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时由调用方线程执行,避免任务丢失 ); // 运行中任务存储,用于超时取消、状态查询 private final ConcurrentHashMap<String, ScheduledFuture<?>> runningTasks = new ConcurrentHashMap<>();
2. 新增请求参数与校验
在TestRequest实体类中新增3个字段:
hourlyRate:每小时发送频次(即需求中的x)durationHours:持续运行时长(即需求中的y)taskId:任务唯一标识,自动生成无需传入
同时加参数校验逻辑,避免超出系统承载上限。
3. 周期调度逻辑实现
修改接口核心逻辑,实现按频次调度、超时自动停止的能力:
@PostMapping("/sendtest") public TestResult sendScenario(@RequestBody TestRequest testRequest) throws Exception { // 前置参数校验 if (testRequest.getHourlyRate() > 19000 || testRequest.getDurationHours() > 72) { return TestResult.fail("参数超出限制:最高每小时19000个场景,最长运行3天"); } String taskId = UUID.randomUUID().toString(); testRequest.setTaskId(taskId); // 计算发送间隔:单位毫秒 3600*1000 / 每小时频次 long intervalMs = 3600 * 1000L / testRequest.getHourlyRate(); // 提交周期调度任务 ScheduledFuture<?> scheduledFuture = scheduler.scheduleAtFixedRate( () -> workerPool.submit(() -> { try { // 原方法内部已实现单场景6个事件的发送逻辑,无需修改 this.sendeventfortest(testResult, testRequest.getLoggerURL(), httpEndPoint); } catch (JsonProcessingException e) { logger.error("任务{}发送场景失败", taskId, e); // 可按需增加最多3次重试逻辑,避免偶发网络波动导致任务失败 } }), 0, intervalMs, TimeUnit.MILLISECONDS ); runningTasks.put(taskId, scheduledFuture); // 注册超时自动取消任务,到达运行时长后自动终止调度 scheduler.schedule(() -> { ScheduledFuture<?> future = runningTasks.remove(taskId); if (future != null && !future.isCancelled()) { future.cancel(false); } }, testRequest.getDurationHours(), TimeUnit.HOURS); return TestResult.success(taskId); }
4. 长稳运行优化
- 给
sendeventfortest方法加HTTP请求超时控制,避免单个请求卡住占用线程 - 批量发送6个事件时复用HTTP连接池,减少TCP握手开销,提升发送效率
- 增加线程池监控,队列长度、活跃线程数、任务失败率达到阈值时触发告警
内容的提问来源于stack exchange,提问作者XYZ
相关产品推荐
相关产品推荐

