AI 对话耗时优化、自定义 SSE 连接池、异步线程池设计,支撑高并发百万级吞吐调优
AI 对话耗时优化、自定义 SSE 连接池、异步线程池设计,支撑高并发百万级吞吐调优
文章标签:#SpringBoot3 #SSE #AI 大模型 #线程池 #高并发 #性能调优 #后端实战
适合人群:Java 后端、AI 应用开发、大模型服务网关、微服务性能调优、架构师
阅读收益:
- 定位 AI 流式对话场景常见性能瓶颈,理清 SSE 高并发下的痛点
- 掌握自定义 SSE 连接池实现,管理长连接生命周期,避免连接泄露
- 学会 AI 业务专属异步线程池设计,区分 IO 密集 / CPU 密集任务
- 掌握超时、排队、熔断、背压策略,实现百万级吞吐生产落地
- 拿到可直接复制改造的生产级 Java 代码,规避线上踩坑
一、背景痛点
基于 SpringBoot SSE 做大模型流式问答,很多同学上线后遇到这些线上问题:
- 用户并发上来,SSE 长连接不断累积,连接没有正常关闭,发生连接泄露,文件句柄耗尽,服务报错。
- 使用默认
SimpleAsyncTaskExecutor,无边界线程,并发暴涨直接 OOM,线程疯狂创建销毁,CPU 飙升。 - 大模型接口响应慢,大量请求堆积,没有排队限流策略,雪崩效应。
- SSE 推送事件和调用大模型接口共用同一个线程,IO 阻塞导致流式输出卡顿,对话耗时变长。
- 缺少连接超时、异常断开、客户端主动断连检测,无效请求持续占用资源。
普通简单 demo 只能跑少量测试并发,真正面向 C 端高并发场景,必须做三件事:SSE 长连接生命周期管理、业务隔离异步线程池、请求全链路耗时优化 + 背压保护。
核心区分:SSE 是服务端向客户端推送的长连接;调用大模型 API 属于 IO 密集型远程调用;两者不能混用一套线程。
二、整体架构思路
- SSE 连接池:统一管理所有客户端 SseEmitter,记录会话 ID、用户 ID、创建时间、状态,自动清理超时、断开、异常的连接,防止句柄泄露。
- 双线程池隔离
- IO 线程池:专门调用大模型第三方 API,网络 IO 密集,线程数配置偏大。
- SSE 推送线程池:专门负责向客户端推送流式 chunk 数据,轻量任务。
- 全链路耗时优化手段
- 请求排队 + 拒绝策略,超过阈值直接返回繁忙,拒绝新请求保护服务。
- 客户端断开感知,立刻中断大模型调用,释放资源,不做无效等待。
- 设置大模型调用超时,避免长慢请求占满线程。
- 异步解耦,controller 只接收请求,实际模型调用交给业务线程池。
- 熔断降级:大模型服务报错、超时占比过高,触发局部熔断,减少无效调用。
三、自定义 SSE 连接池实现
SseEmitter 原生没有自带连接池,默认不做管理,如果客户端浏览器关闭页面,服务端不会自动回收,极易泄露。我们维护本地内存池管理会话,生产环境集群场景可以结合 Redis 做会话元数据同步。
会话实体 SseSession.java
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.time.LocalDateTime;
/**
* SSE会话元数据
*/
public class SseSession {
/** 会话唯一id */
private String sessionId;
/** 用户id */
private Long userId;
/** sse发射器 */
private SseEmitter emitter;
/** 创建时间 */
private LocalDateTime createTime;
/** 是否已经结束 */
private volatile boolean completed;
public SseSession(String sessionId, Long userId, SseEmitter emitter) {
this.sessionId = sessionId;
this.userId = userId;
this.emitter = emitter;
this.createTime = LocalDateTime.now();
this.completed = false;
// 注册回调:正常完成 / 超时 / 异常
emitter.onCompletion(() -> this.completed = true);
emitter.onTimeout(() -> this.completed = true);
emitter.onError(e -> this.completed = true);
}
/**
* 发送sse消息
*/
public void send(Object data) throws IOException {
if (!completed) {
emitter.send(data);
}
}
/**
* 关闭连接
*/
public void complete() {
if (!completed) {
emitter.complete();
this.completed = true;
}
}
// getter setter省略
}
SSE 连接池管理器 SseEmitterPool.java
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.*;
@Component
public class SseEmitterPool {
/** 存储全部活跃sse会话 */
private final Map<String, SseSession> sessionMap = new ConcurrentHashMap<>();
/** 定时清理过期会话的调度线程 */
private ScheduledExecutorService cleanScheduler;
/** 会话最大存活时间,单位秒,超时自动回收 */
private static final long SESSION_MAX_ALIVE_SECOND = 300;
@PostConstruct
public void init() {
cleanScheduler = Executors.newSingleThreadScheduledExecutor();
// 每30秒扫描一次,清理过期、已断开会话
cleanScheduler.scheduleAtFixedRate(this::cleanInvalidSession, 30, 30, TimeUnit.SECONDS);
}
/**
* 添加会话到池子
*/
public void addSession(SseSession session) {
sessionMap.put(session.getSessionId(), session);
}
/**
* 获取会话
*/
public SseSession getSession(String sessionId) {
return sessionMap.get(sessionId);
}
/**
* 移除并且关闭会话
*/
public void removeSession(String sessionId) {
SseSession session = sessionMap.remove(sessionId);
if (session != null) {
session.complete();
}
}
/**
* 广播消息(可选)
*/
public void broadcast(String msg) {
for (SseSession session : sessionMap.values()) {
try {
session.send(msg);
} catch (IOException e) {
removeSession(session.getSessionId());
}
}
}
/**
* 清理无效过期会话
*/
private void cleanInvalidSession() {
long nowSecond = System.currentTimeMillis() / 1000;
sessionMap.entrySet().removeIf(entry -> {
SseSession session = entry.getValue();
// 已经标记完成 或者会话超时
boolean invalid = session.isCompleted() ||
(nowSecond - session.getCreateTime().getSecond()) > SESSION_MAX_ALIVE_SECOND;
if (invalid) {
session.complete();
return true;
}
return false;
});
}
@PreDestroy
public void destroy() {
cleanScheduler.shutdown();
// 应用关闭,全部断开
sessionMap.values().forEach(SseSession::complete);
sessionMap.clear();
}
public int getActiveCount() {
return sessionMap.size();
}
}
四、AI 业务自定义异步线程池设计
大模型调用属于 IO 密集任务,不要使用 Spring 默认线程池。IO 密集型线程数可以设置大一些;SSE 推送是轻量任务单独隔离。 重要:一定要指定拒绝策略,并发打满不能无限创建线程。
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.concurrent.*;
@Configuration
public class AiThreadPoolConfig {
/**
* 大模型远程调用IO线程池:IO密集,等待第三方模型接口返回
*/
@Bean("aiModelIoThreadPool")
public Executor aiModelIoThreadPool() {
ThreadPoolExecutor executor = new ThreadPoolExecutor(
20,
120,
60L,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(500),
new ThreadPoolExecutor.CallerRunsPolicy()
);
return executor;
}
/**
* SSE消息推送线程池:轻量推送任务
*/
@Bean("ssePushThreadPool")
public Executor ssePushThreadPool() {
ThreadPoolExecutor executor = new ThreadPoolExecutor(
10,
40,
30L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(200),
new ThreadPoolExecutor.AbortPolicy()
);
return executor;
}
}
参数说明:
- IO 密集任务,最大线程数可以调高,因为大部分时间线程在等待网络返回,不占用 CPU。
- 队列设置有界,防止任务无限堆积内存溢出。
CallerRunsPolicy:队列满了,交给调用者线程执行,起到限流背压效果,不会直接抛异常。
五、改造 SSE Controller 整合连接池 + 线程池
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.UUID;
import java.util.concurrent.Executor;
@RestController
@RequestMapping("/ai")
public class AiSseController {
@Autowired
private SseEmitterPool sseEmitterPool;
@Autowired
@Qualifier("aiModelIoThreadPool")
private Executor aiModelIoThreadPool;
@Autowired
@Qualifier("ssePushThreadPool")
private Executor ssePushThreadPool;
/**
* 流式对话接口
*/
@GetMapping(value = "/stream/chat", produces = "text/event-stream;charset=utf-8")
public SseEmitter chatStream(@RequestParam Long userId, @RequestParam String question) {
// 设置超时时间
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);
String sessionId = UUID.randomUUID().toString();
SseSession sseSession = new SseSession(sessionId, userId, emitter);
// 存入连接池
sseEmitterPool.addSession(sseSession);
// 交给AI IO线程池执行模型调用逻辑,不阻塞http主线程
aiModelIoThreadPool.execute(() -> {
try {
// 模拟调用大模型流式接口,真实业务替换为你的大模型sdk调用
String[] chunks = {"你好",",","这是","流式回答","测试。"};
for (String chunk : chunks) {
// 判断连接是否已经断开,客户端关掉页面直接终止任务,节约资源
if(sseSession.isCompleted()){
break;
}
ssePushThreadPool.execute(()->{
try {
sseSession.send(chunk);
} catch (IOException e) {
//发送失败,移除会话
sseEmitterPool.removeSession(sessionId);
}
});
Thread.sleep(150);
}
//结束标记
sseSession.send("[DONE]");
} catch (Exception e) {
try {
sseSession.send("error:服务繁忙");
} catch (IOException ignored) {
}
}finally {
sseEmitterPool.removeSession(sessionId);
}
});
return emitter;
}
}
六、高并发耗时优化关键手段
1. 客户端断开快速感知,及时终止模型调用
很多线上浪费来自:浏览器关掉页面,服务端还在持续请求大模型接口、继续生成 token。
在每一轮循环推送 chunk 的时候,判断
sseSession.isCompleted(),一旦 true 立刻终止循环,中断大模型 http 请求,释放线程。
2. 请求排队、背压保护
线程池队列写满之后,拒绝策略不要直接抛异常。 业务上可以做:活跃 SSE 会话数量监控,超过阈值直接返回服务繁忙,拒绝新来请求,保护整个服务不雪崩。
//示例:超过最大允许并发连接数,直接拒绝
if(sseEmitterPool.getActiveCount() > 800){
throw new RuntimeException("当前访问量过大,请稍后重试");
}
3. HTTP 客户端优化(大模型调用)
调用大模型三方 API,不要使用默认 RestTemplate,使用 Okhttp3 或者 HttpClient5,配置连接池,设置合理 connect、read 超时。
如果 http 没有设置超时,第三方大模型卡住,会直接占满业务线程池。
4. 超时全链路管控
- SseEmitter 设置总超时;
- http 调用大模型设置 IO 超时;
- 业务线程池队列有界;
- SSE 连接池定时回收过期会话。
5. 监控指标埋点(生产必备)
需要监控的指标:
- SSE 活跃连接数;
- AI 线程池活跃线程、队列大小;
- 请求总耗时,p50/p95/p99;
- SSE 异常断开次数、超时次数。 对接 Prometheus + Grafana 或者 Skywalking 做告警。
七、生产环境注意点
- 集群部署问题:上面 SSE 连接池是本地内存模式。如果多实例部署,一个用户的 SSE 连接落在实例 A,后续消息无法跨实例推送。
解决方案:网关层会话粘性,或者使用 Redis + 消息队列(RabbitMQ/RocketMQ)实现跨实例消息分发。
- Nginx 反向代理 SSE 必须配置
proxy_cache off;
proxy_buffering off;
proxy_read_timeout 300s;
如果开启 proxy_buffering,nginx 会缓存 sse 输出,前端收不到流式实时效果。
- 文件句柄调优 大量 SSE 长连接,操作系统 file descriptors 要调高,否则会报 too many open files。
- 线程池参数不要拍脑袋,压测后调参。IO 密集和 CPU 密集任务线程池严格隔离,不要共用。
八、压测预期效果
经过这套改造之后:
- SSE 长连接不会泄露,自动回收无效会话;
- 线程资源可控,不会并发上来直接 OOM;
- 无效请求及时中断,减少大模型接口浪费;
- 可以支撑上万 SSE 长连接,配合网关限流,实现百万级吞吐能力。
注意:百万级吞吐指全链路 qps,受限于下游大模型接口能力,如果大模型本身 QPS 有限,要在网关层增加限流和排队。
九、后续扩展方向
- 接入 Sentinel,针对 AI 接口做流量控制、熔断降级;
- Redis 存储对话上下文,实现会话跨实例;
- Token 统计、计费拦截;
- 增加链路 traceId,全链路追踪每一条 AI 对话请求。
更多推荐

所有评论(0)