惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

V
V2EX
宝玉的分享
宝玉的分享
Jina AI
Jina AI
IT之家
IT之家
博客园 - Franky
MyScale Blog
MyScale Blog
Y
Y Combinator Blog
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
I
InfoQ
雷峰网
雷峰网
WordPress大学
WordPress大学
Microsoft Security Blog
Microsoft Security Blog
Google DeepMind News
Google DeepMind News
美团技术团队
S
SegmentFault 最新的问题
罗磊的独立博客
博客园 - 聂微东
大猫的无限游戏
大猫的无限游戏
H
Help Net Security
D
Docker
博客园 - 司徒正美
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
阮一峰的网络日志
阮一峰的网络日志
M
MIT News - Artificial intelligence

博客园 - work hard work smart

Java 面试1 Java 常见面试问题 手撕java常用代码 WebStorm 创建react工程 构建企业级 Text-to-SQL Agent:基于 LangGraph 的智能数据查询系统设计 Harness 工程:驾驭 AI Agent 的工程化艺术 DeepAgents 多智能体架构实战:从设计模式到后端选型 Vue 自定义组件完全指南:从零构建待办事项应用 使用 LangChain + Hugging Face 构建文本向量化服务 SQLAlchemy 使用详解 Python 中使用 Elasticsearch 的完整指南 Qdrant 向量数据库使用指南 OpenEvals 快速入门:LLM 评估指南 DeepEval 快速入门:LLM 应用评估指南 LangSmith 批量评估完全指南 Qwen-Agent 入门指南:快速构建智能体应用 LangSmith 集成实战:从追踪到评估的完整指南 初识 go-zero:一款让你写后端更规范、更高效的 Go 微服务框架 RAG 中为什么需要 Rerank,以及如何使用 Rerank LangChain4j RAG 核心组件与组合方式 如何使用 Elasticsearch 进行全文检索和向量检索 MinerU Docker 部署指南 5 分钟上手:为 Cline 配置一个免费的 MCP 天气服务 Neo4j 图数据库安装与 Spring Boot 集成实战指南 LangFuse 实战指南:用 @observe 三行代码给 LLM 应用加上全链路追踪 Function Call 深度解析:让大模型从"嘴炮"到"实干"的技术革命 Spring AI 提示词模板实战:告别硬编码,实现提示词工程化管理 LangChain4j 实战指南:用 Java 轻松构建 AI 应用 Spring AI 对话短期记忆实战:让大模型拥有"记忆力" Spring AI 提示词工程实战:让大模型更懂你的意图
Spring 中 SSE 流式输出的多种实现方式详解
work hard work smart · 2026-05-30 · via 博客园 - work hard work smart

摘要:Server-Sent Events (SSE) 是一种基于 HTTP 的服务器推送技术,在 AI 大模型流式响应场景中广泛应用。本文结合实际代码,详细介绍 Spring 框架中实现 SSE 的三种主流方式。

什么是 SSE?

SSE(Server-Sent Events)是一种允许服务器向浏览器推送实时数据的 Web 标准技术。它基于 HTTP 协议,使用 text/event-stream 内容类型,具有以下特点:

  • 单向通信:仅支持服务器向客户端推送数据
  • 基于 HTTP:无需特殊协议,兼容现有基础设施
  • 自动重连:内置断线重连机制
  • 事件 ID:支持事件标识和恢复

与 WebSocket 相比,SSE 更简单易用,适合服务器推送场景(如 AI 对话流式输出、实时通知、数据监控等)。


实现方式一:SseEmitter(Spring MVC 原生支持)

SseEmitter 是 Spring Web MVC 提供的 SSE 支持类,常用于处理大模型的流式响应。

代码示例

@GetMapping("/streamEvents")
public SseEmitter streamServerEvents() {
    SseEmitter emitter = new SseEmitter(60_000L);
    Executors.newVirtualThreadPerTaskExecutor().submit(() -> {
        try {
            for (int i = 0; i < 50; i++) {
                emitter.send("Event Data: " + i);
                Thread.sleep(500);
            }
        } catch (IOException | InterruptedException e) {
            emitter.completeWithError(e);
        } finally {
            emitter.complete();
        }
    });
    return emitter;
}

核心要点

  1. 超时设置new SseEmitter(60_000L) 设置 60 秒超时,避免连接长时间占用
  2. 异步执行:使用虚拟线程(Java 21+)在后台发送数据,不阻塞主线程
  3. 异常处理:通过 completeWithError() 处理异常,complete() 正常结束
  4. 适用场景:传统 Spring MVC 项目、简单流式推送

优缺点

优点 缺点
简单易用,无需额外依赖 手动管理线程和生命周期
兼容性好,Spring 4.2+ 支持 错误处理需要手动编码
支持自定义事件名称和 ID 不适合复杂响应式场景

实现方式二:StreamingResponseBody(灵活控制响应流)

StreamingResponseBody 允许直接操作 OutputStream,适合自定义流式响应格式。

代码示例

@GetMapping("/streamResponse")
public ResponseEntity<StreamingResponseBody> streamChatResponse() {
    
    StreamingResponseBody body = outputStream -> {
        for (int i = 0; i < 50; i++) {
            String data = "Chunk: " + i + "\n";
            outputStream.write(data.getBytes(StandardCharsets.UTF_8));
            outputStream.flush();
            try {
                Thread.sleep(500); // 模拟处理延迟
            } catch (InterruptedException e) {
                throw new RuntimeException(e);
            }
        }
    };

    return ResponseEntity.ok()
            .header(HttpHeaders.CONTENT_TYPE, MediaType.TEXT_EVENT_STREAM_VALUE)
            .body(body);
}

核心要点

  1. 直接写流:通过 outputStream.write() 直接写入字节数据
  2. 手动刷新:每次写入后调用 flush() 确保数据立即发送
  3. Content-Type:需手动设置 text/event-stream 标识 SSE 流
  4. 适用场景:需要精细控制输出格式、文件流式下载

优缺点

优点 缺点
完全控制输出格式 需要手动处理所有细节
适合大数据量传输 不符合 SSE 标准格式(需自行实现)
可用于非 SSE 场景 缺乏事件 ID、重试等 SSE 特性

注意事项

这种方式严格来说不是标准 SSE,而是普通的 HTTP 流式响应。如果要实现标准 SSE,需要按照以下格式输出:

data: {"message": "Hello"}
id: 1
event: customEvent
retry: 3000


实现方式三:Flux(响应式编程 - Spring WebFlux)

Flux 是 Spring WebFlux(基于 Project Reactor)提供的响应式类型,适合高并发的流式调用场景。

代码示例

@GetMapping(value = "/reactiveStream")
public Flux<String> reactiveDataStream() {
    return Flux.interval(Duration.ofSeconds(1))
            .take(20)
            .map(seq -> "Reactive Event: " + seq);
}

核心要点

  1. 声明式编程:通过链式调用定义数据流,无需手动管理线程
  2. 背压支持:自动处理消费者速度慢于生产者的情况
  3. 非阻塞:基于事件循环,资源占用更少
  4. 适用场景:高并发场景、响应式微服务、实时数据流处理

进阶示例:结合大模型 API

@GetMapping("/aiStream")
public Flux<String> streamAIResponse(@RequestParam String question) {
    return webClient.post()
            .uri("https://api.openai.com/v1/chat/completions")
            .bodyValue(buildRequest(question))
            .retrieve()
            .bodyToFlux(String.class)
            .map(this::parseSSE);
}

优缺点

优点 缺点
天然支持响应式流 学习曲线陡峭
自动背压和资源管理 需要引入 WebFlux 依赖
组合操作符丰富 与传统 MVC 不兼容

三种方式对比

特性 SseEmitter StreamingResponseBody Flux
编程模型 命令式 命令式 响应式
线程管理 手动 手动 自动
SSE 标准支持 ✅ 完整支持 ⚠️ 需手动实现 ✅ 完整支持
背压支持
适用框架 Spring MVC Spring MVC Spring WebFlux
学习成本
推荐场景 简单推送 自定义流式输出 高并发响应式应用

实际应用:调用大模型流式 API

以下是一个完整的示例,展示如何转发大模型的 SSE 响应:

@RequestMapping("/mockStreamResponse")
public String mockStreamResponse() {
    String requestBody = """
            {
                "model": "qwen-plus",
                "messages": [
                    {
                        "role": "system",
                        "content": "You are a professional technical writer."
                    },
                    {
                        "role": "user",
                        "content": "请简要说明Python在数据科学中的优势"
                    }
                ],
                "stream": true
            }
            """;
    
    HttpClient client = HttpClient.newHttpClient();
    HttpRequest request = HttpRequest.newBuilder()
            .uri(URI.create(API_URL))
            .header("Content-Type", "application/json")
            .header("Authorization", "Bearer " + API_KEY)
            .header("X-DashScope-SSE", "enable")
            .POST(HttpRequest.BodyPublishers.ofString(requestBody))
            .build();

    try {
        HttpResponse<String> response = client.send(
                request, HttpResponse.BodyHandlers.ofString());
        return response.body();
    } catch (IOException | InterruptedException e) {
        throw new RuntimeException(e);
    }
}

关键配置

  • X-DashScope-SSE: enable:启用阿里云 DashScope 的 SSE 支持
  • stream: true:请求体中开启流式模式
  • 直接返回原始 SSE 数据,由前端解析

前端调用示例

// 方式一:使用 EventSource(推荐)
const eventSource = new EventSource('/streamEvents');

eventSource.onmessage = (event) => {
    console.log('Received:', event.data);
};

eventSource.onerror = (error) => {
    console.error('Connection error:', error);
    eventSource.close();
};

// 方式二:使用 fetch(更灵活)
const response = await fetch('/reactiveStream');
const reader = response.body.getReader();

while (true) {
    const { done, value } = await reader.read();
    if (done) break;
    console.log('Chunk:', new TextDecoder().decode(value));
}

最佳实践

1. 合理设置超时时间

// 根据业务场景调整超时
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L); // 5分钟

2. 处理客户端断开连接

emitter.onCompletion(() -> {
    log.info("SSE connection completed");
});

emitter.onTimeout(() -> {
    log.warn("SSE connection timeout");
});

3. 使用结构化数据

// 发送 JSON 格式数据
emitter.send(SseEmitter.event()
    .name("message")
    .data(Map.of("content", "Hello", "timestamp", System.currentTimeMillis())));

4. 监控和日志

emitter.onError((ex) -> {
    log.error("SSE error: {}", ex.getMessage(), ex);
    metrics.increment("sse.errors");
});

总结

Spring 提供了多种方式实现 SSE 流式输出,选择时应考虑:

  • 项目技术栈:MVC 还是 WebFlux
  • 功能需求:是否需要背压、事件 ID 等特性
  • 团队熟悉度:响应式编程的学习成本
  • 性能要求:高并发场景推荐 WebFlux + Flux

在 Spring 项目中,推荐使用 FluxSseEmitter,具体选择取决于你使用的是 Spring WebFlux 还是 Spring Web MVC。


参考资源