在大模型(LLM)应用、实时通知、进度条推送等场景中,SSE(Server-Sent Events) 已经成为替代传统轮询和 WebSocket 的首选方案。它基于 HTTP 协议,实现简单,且天然支持防火墙和代理。
本文将手把手教你如何在 Java 中实现 SSE 服务端(流式输出)以及作为客户端调用 SSE 流,并提供生产环境的避坑指南。
一、 什么是 SSE?
SSE(服务器发送事件)是一种允许服务器向客户端单向推送数据的 HTTP 技术。
单向通信:只能服务器向客户端推送,客户端不能向服务器发送数据(区别于 WebSocket)。
基于 HTTP:使用标准的 HTTP 协议,无需额外握手,兼容性好。
文本格式:默认传输纯文本,通常配合 JSON 使用。
自动重连:浏览器原生
EventSource支持断线自动重连。
二、 服务端实现:Spring Boot 输出 SSE 流
在 Spring Boot 中,实现 SSE 的核心类是 SseEmitter。
1. 引入依赖
如果你使用的是 Spring Boot Web,默认已包含所需依赖。
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>2. 编写 SSE Controller
我们需要返回 SseEmitter 对象,并在异步线程中不断向其中写入数据。
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@RestController
@RequestMapping("/sse")
public class SseController {
// 建议使用自定义线程池,避免使用默认的 ForkJoinPool
private final ExecutorService executor = Executors.newCachedThreadPool();
@GetMapping("/stream")
public SseEmitter streamData() {
// 设置超时时间为 60 秒(0 表示不超时,但生产环境建议设置合理值)
SseEmitter emitter = new SseEmitter(60000L);
// 注册回调事件
emitter.onCompletion(() -> System.out.println("SSE 连接完成"));
emitter.onTimeout(() -> System.out.println("SSE 连接超时"));
emitter.onError((e) -> System.out.println("SSE 连接异常: " + e.getMessage()));
// 异步发送数据,避免阻塞 Tomcat 工作线程
executor.execute(() -> {
try {
for (int i = 1; i <= 5; i++) {
// 模拟耗时操作(如大模型推理)
Thread.sleep(1000);
// 构建 SSE 事件
String data = "这是第 " + i + " 条流式消息";
// 发送数据。name 是事件名(可选),data 是数据内容
emitter.send(SseEmitter.event()
.name("message")
.data(data));
}
// 数据发送完毕,关闭连接
emitter.complete();
} catch (IOException | InterruptedException e) {
// 客户端主动断开连接时,会抛出异常,需要捕获并标记错误
emitter.completeWithError(e);
}
});
return emitter;
}
}3. 测试服务端
启动项目后,在浏览器访问 http://localhost:8080/sse/stream,或者使用 curl 测试:
curl -N http://localhost:8080/sse/stream你会看到数据每隔一秒逐行输出:
event: message
data: 这是第 1 条流式消息
event: message
data: 这是第 2 条流式消息
...三、 客户端调用:Java 消费 SSE 流
在实际业务中,我们经常需要让 Java 后端作为客户端,去调用第三方大模型(如 OpenAI、通义千问)的 SSE 接口。
这里提供两种主流方案:Java 11+ 原生 HttpClient 和 OkHttp。
方案 A:使用 OkHttp(强烈推荐,最优雅)
OkHttp 提供了专门的 okhttp-sse 扩展库,完美封装了 SSE 协议的解析。
1. 引入依赖:
<dependency>
<groupId>com.squareup.okhttp3</groupId>
<artifactId>okhttp</artifactId>
<version>4.12.0</version>
</dependency>
<dependency>
<groupId>com.squareup.okhttp3</groupId>
<artifactId>okhttp-sse</artifactId>
<version>4.12.0</version>
</dependency>2. 编写调用代码:
import okhttp3.OkHttpClient;
import okhttp3.Request;
import okhttp3.Response;
import okhttp3.sse.EventSource;
import okhttp3.sse.EventSourceListener;
import okhttp3.sse.EventSources;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
import java.util.concurrent.TimeUnit;
public class OkHttpSseClient {
public static void main(String[] args) {
// 1. 配置 OkHttpClient,注意设置读取超时时间为 0(不超时)
OkHttpClient client = new OkHttpClient.Builder()
.readTimeout(0, TimeUnit.MILLISECONDS)
.build();
// 2. 构建请求
Request request = new Request.Builder()
.url("http://localhost:8080/sse/stream")
.header("Accept", "text/event-stream")
.get()
.build();
// 3. 创建 EventSource 监听器
EventSourceListener listener = new EventSourceListener() {
@Override
public void onOpen(@NotNull EventSource eventSource, @NotNull Response response) {
System.out.println("连接已建立");
}
@Override
public void onEvent(@NotNull EventSource eventSource, @Nullable String id, @Nullable String type, @NotNull String data) {
// 核心:接收到的 data 就是服务端发送的内容
System.out.println("收到流式数据 -> type: " + type + ", data: " + data);
}
@Override
public void onClosed(@NotNull EventSource eventSource) {
System.out.println("连接已关闭");
}
@Override
public void onFailure(@NotNull EventSource eventSource, @Nullable Throwable t, @Nullable Response response) {
System.out.println("连接异常: " + (t != null ? t.getMessage() : "unknown"));
}
};
// 4. 创建 EventSource 并开始监听
EventSource eventSource = EventSources.createFactory(client).newEventSource(request, listener);
// 注意:如果是 main 方法测试,需要阻塞主线程,否则程序会直接退出
try { Thread.sleep(10000); } catch (InterruptedException e) {}
}
}方案 B:使用 Java 11+ 原生 HttpClient
如果你不想引入第三方库,可以使用 Java 11 引入的 HttpClient。但需要注意,它按行读取,需要你自己处理 SSE 的换行和协议头。
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.Duration;
public class Java11SseClient {
public static void main(String[] args) throws Exception {
HttpClient client = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(10))
.build();
HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create("http://localhost:8080/sse/stream"))
.header("Accept", "text/event-stream")
.GET()
.build();
// 使用 ofLines() 按行流式读取响应体
client.sendAsync(request, HttpResponse.BodyHandlers.ofLines())
.thenAccept(response -> {
response.body().forEach(line -> {
// 简单的 SSE 解析:过滤掉 event: 和空行,只取 data:
if (line.startsWith("data:")) {
String data = line.substring(5).trim();
System.out.println("收到数据: " + data);
}
});
})
.join(); // 阻塞等待完成
}
}四、 进阶:模拟大模型 JSON 流式输出
大模型(如 ChatGPT)返回的 SSE 数据通常是 JSON 格式,且包含特殊的结束标志 [DONE]。
1. 服务端模拟大模型返回格式
// 在 Controller 的循环中修改发送逻辑
String jsonContent = "{\"choices\": [{\"delta\": {\"content\": \"字\" + i + \"}}]}";
emitter.send(SseEmitter.event().data(jsonContent));
// 循环结束后,发送结束标志
emitter.send(SseEmitter.event().data("[DONE]"));
emitter.complete();2. 客户端解析 JSON 流
在使用 OkHttp 客户端时,在 onEvent 中解析 JSON:
@Override
public void onEvent(@NotNull EventSource eventSource, @Nullable String id, @Nullable String type, @NotNull String data) {
// 1. 判断是否结束
if ("[DONE]".equals(data)) {
System.out.println("\n--- 输出结束 ---");
return;
}
// 2. 解析 JSON (使用 Jackson 或 Fastjson)
// 这里为了演示,直接截取 content 的值
// 实际开发中请使用 ObjectMapper 解析
System.out.print(extractContent(data));
}
private String extractContent(String json) {
// 简单的字符串截取,实际请用 JSON 库
int start = json.indexOf("\"content\": \"") + 12;
int end = json.indexOf("\"}}", start);
return json.substring(start, end);
}五、 生产环境避坑指南(必看)
在实际项目中,SSE 往往会遇到一些“玄学”问题,请务必注意以下几点:
1. Nginx 代理导致流式输出失效(最常见)
Nginx 默认会开启缓冲(Buffering),这会导致服务端发送的数据被攒在 Nginx 内存中,直到缓冲区满了才发给客户端,完全破坏了流式效果。 解决办法:在 Nginx 配置中关闭缓冲:
location /api/sse/ {
proxy_pass http://backend;
proxy_buffering off; # 必须关闭缓冲!
proxy_cache off; # 关闭缓存
proxy_read_timeout 300s; # 延长读取超时时间,防止长连接被 Nginx 掐断
}2. SseEmitter 的线程安全问题
Spring 官方文档明确指出:SseEmitter 不是线程安全的。 如果你的业务逻辑中,有多个线程同时调用同一个 emitter.send(),会导致数据交错甚至抛出异常。 解决办法:
方案一:在调用
send时加synchronized锁。方案二(推荐):使用一个单线程的线程池,或者使用
ConcurrentLinkedQueue将消息入队,由单个线程循环消费并发送。
3. 客户端断开导致的异常刷屏
当用户关闭浏览器或取消请求时,服务端的 emitter.send() 会抛出 IOException: Broken pipe 或 ClientAbortException。 解决办法:在 catch 块中捕获异常,并调用 emitter.completeWithError(e) 或 emitter.complete(),同时不要打印 Error 级别的日志,降级为 Debug 或 Warn,否则日志会被刷屏。
4. 连接数耗尽问题
SSE 是长连接,每个 SSE 请求都会占用 Tomcat 的一个工作线程(或 NIO 连接)。如果并发量极大(如几万人同时请求),可能会耗尽连接数。 解决办法:
调大 Tomcat 的
server.tomcat.max-connections。对于超高并发场景,建议引入 WebSocket,或者使用响应式框架(Spring WebFlux)来实现 SSE。
六、 总结
服务端:使用 Spring Boot 的
SseEmitter,注意异步发送和超时设置。客户端:推荐使用
OkHttp的EventSource,解析优雅且自带重连机制。运维:务必检查 Nginx 的
proxy_buffering off配置。大模型场景:注意处理
[DONE]结束标志和 JSON 流式解析。
掌握了以上知识,你就可以在 Java 项目中轻松实现类似 ChatGPT 的“打字机”流式输出效果了!