AI 对话耗时优化、自定义 SSE 连接池、异步线程池设计,支撑高并发百万级吞吐调优

文章标签:#SpringBoot3 #SSE #AI 大模型 #线程池 #高并发 #性能调优 #后端实战

适合人群:Java 后端、AI 应用开发、大模型服务网关、微服务性能调优、架构师

阅读收益

  • 定位 AI 流式对话场景常见性能瓶颈,理清 SSE 高并发下的痛点
  • 掌握自定义 SSE 连接池实现,管理长连接生命周期,避免连接泄露
  • 学会 AI 业务专属异步线程池设计,区分 IO 密集 / CPU 密集任务
  • 掌握超时、排队、熔断、背压策略,实现百万级吞吐生产落地
  • 拿到可直接复制改造的生产级 Java 代码,规避线上踩坑

一、背景痛点

基于 SpringBoot SSE 做大模型流式问答,很多同学上线后遇到这些线上问题:

  1. 用户并发上来,SSE 长连接不断累积,连接没有正常关闭,发生连接泄露,文件句柄耗尽,服务报错。
  2. 使用默认SimpleAsyncTaskExecutor,无边界线程,并发暴涨直接 OOM,线程疯狂创建销毁,CPU 飙升。
  3. 大模型接口响应慢,大量请求堆积,没有排队限流策略,雪崩效应。
  4. SSE 推送事件和调用大模型接口共用同一个线程,IO 阻塞导致流式输出卡顿,对话耗时变长。
  5. 缺少连接超时、异常断开、客户端主动断连检测,无效请求持续占用资源。

普通简单 demo 只能跑少量测试并发,真正面向 C 端高并发场景,必须做三件事:SSE 长连接生命周期管理、业务隔离异步线程池、请求全链路耗时优化 + 背压保护

核心区分:SSE 是服务端向客户端推送的长连接;调用大模型 API 属于 IO 密集型远程调用;两者不能混用一套线程。

二、整体架构思路

  1. SSE 连接池:统一管理所有客户端 SseEmitter,记录会话 ID、用户 ID、创建时间、状态,自动清理超时、断开、异常的连接,防止句柄泄露。
  2. 双线程池隔离
    • IO 线程池:专门调用大模型第三方 API,网络 IO 密集,线程数配置偏大。
    • SSE 推送线程池:专门负责向客户端推送流式 chunk 数据,轻量任务。
  3. 全链路耗时优化手段
    • 请求排队 + 拒绝策略,超过阈值直接返回繁忙,拒绝新请求保护服务。
    • 客户端断开感知,立刻中断大模型调用,释放资源,不做无效等待。
    • 设置大模型调用超时,避免长慢请求占满线程。
    • 异步解耦,controller 只接收请求,实际模型调用交给业务线程池。
  4. 熔断降级:大模型服务报错、超时占比过高,触发局部熔断,减少无效调用。

三、自定义 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;
    }
}

参数说明:

  1. IO 密集任务,最大线程数可以调高,因为大部分时间线程在等待网络返回,不占用 CPU。
  2. 队列设置有界,防止任务无限堆积内存溢出。
  3. 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. 监控指标埋点(生产必备)

需要监控的指标:

  1. SSE 活跃连接数;
  2. AI 线程池活跃线程、队列大小;
  3. 请求总耗时,p50/p95/p99;
  4. SSE 异常断开次数、超时次数。 对接 Prometheus + Grafana 或者 Skywalking 做告警。

七、生产环境注意点

  1. 集群部署问题:上面 SSE 连接池是本地内存模式。如果多实例部署,一个用户的 SSE 连接落在实例 A,后续消息无法跨实例推送。

解决方案:网关层会话粘性,或者使用 Redis + 消息队列(RabbitMQ/RocketMQ)实现跨实例消息分发。

  1. Nginx 反向代理 SSE 必须配置
proxy_cache off;
proxy_buffering off;
proxy_read_timeout 300s;

如果开启 proxy_buffering,nginx 会缓存 sse 输出,前端收不到流式实时效果。

  1. 文件句柄调优 大量 SSE 长连接,操作系统 file descriptors 要调高,否则会报 too many open files。
  2. 线程池参数不要拍脑袋,压测后调参。IO 密集和 CPU 密集任务线程池严格隔离,不要共用。

八、压测预期效果

经过这套改造之后:

  • SSE 长连接不会泄露,自动回收无效会话;
  • 线程资源可控,不会并发上来直接 OOM;
  • 无效请求及时中断,减少大模型接口浪费;
  • 可以支撑上万 SSE 长连接,配合网关限流,实现百万级吞吐能力。

注意:百万级吞吐指全链路 qps,受限于下游大模型接口能力,如果大模型本身 QPS 有限,要在网关层增加限流和排队。

九、后续扩展方向

  1. 接入 Sentinel,针对 AI 接口做流量控制、熔断降级;
  2. Redis 存储对话上下文,实现会话跨实例;
  3. Token 统计、计费拦截;
  4. 增加链路 traceId,全链路追踪每一条 AI 对话请求。
Logo

欢迎加入 MCP 技术社区!与志同道合者携手前行,一同解锁 MCP 技术的无限可能!

更多推荐