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

推荐订阅源

Cisco Talos Blog
Cisco Talos Blog
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
Google Online Security Blog
Google Online Security Blog
博客园 - Franky
Hugging Face - Blog
Hugging Face - Blog
Security Archives - TechRepublic
Security Archives - TechRepublic
博客园 - 司徒正美
N
News and Events Feed by Topic
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
WordPress大学
WordPress大学
博客园 - 三生石上(FineUI控件)
Help Net Security
Help Net Security
N
News and Events Feed by Topic
O
OpenAI News
L
LangChain Blog
F
Full Disclosure
A
About on SuperTechFans
The GitHub Blog
The GitHub Blog
GbyAI
GbyAI
Cloudbric
Cloudbric
W
WeLiveSecurity
Application and Cybersecurity Blog
Application and Cybersecurity Blog
罗磊的独立博客
Attack and Defense Labs
Attack and Defense Labs
PCI Perspectives
PCI Perspectives
TaoSecurity Blog
TaoSecurity Blog
AI
AI
有赞技术团队
有赞技术团队
酷 壳 – CoolShell
酷 壳 – CoolShell
C
CXSECURITY Database RSS Feed - CXSecurity.com
C
Cisco Blogs
D
Darknet – Hacking Tools, Hacker News & Cyber Security
Apple Machine Learning Research
Apple Machine Learning Research
C
CERT Recently Published Vulnerability Notes
T
The Exploit Database - CXSecurity.com
T
Threatpost
P
Palo Alto Networks Blog
G
GRAHAM CLULEY
Last Week in AI
Last Week in AI
雷峰网
雷峰网
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
C
Cyber Attacks, Cyber Crime and Cyber Security
博客园 - 聂微东
P
Proofpoint News Feed
Latest news
Latest news
S
SegmentFault 最新的问题
J
Java Code Geeks
T
Threat Research - Cisco Blogs
H
Help Net Security
P
Privacy International News Feed

博客园 - work hard work smart

使用 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 AI ChatClient 深度解析:优雅构建大模型应用的利器 Spring AI Alibaba DashScopeChatModel 实战 OpenSandbox 实战指南:为 AI Agent 构建安全的代码执行沙箱 在本机启动 LangGraph 开发服务器:完整指南 DeepAgents中Backend的奥秘:让AI Agent拥有文件操作能力 为什么选择 Go 开发 Web 接口?从入门到实践 智能搜索DeepAgent笔记 RAG学习笔记2--系统查询流程 RAG学习笔记1--系统文件导入流程 百炼 WebSearch 快速入门指南 Python 连接 MongoDB 完整指南:从连接配置到增删改查实战 一行命令搞定 MongoDB 开发环境:Docker Compose 部署 + 可视化管理 使用 Attu 可视化管理 Milvus 向量数据库 在 Windows Docker 中快速安装 Milvus 2.5.6 minio使用 Spark 集群搭建 hadoop集群安装 Spring AI Alibaba 入门实战 Windows 安装 OpenClaw 实战指南 MyBatis 核心流程和原理 Idea中安装Claude code插件 Spark 编程 使用Matplotlib 绘制直方图 Flink安装部署 Flume安装 查找导致cpu过高的代码方法 JVisualVM监控远程Java进程 jmap jacoco多模块生成java单元测试报告实践 arthas 使用demo LockSupport Exchanger CyclicBarrier CountDownLatch 手把手教你用python开始第一个机器学习项目
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。


参考资源