Java 中 SSE (Server-Sent Events) 流式输出与调用完全指南

Java 中 SSE (Server-Sent Events) 流式输出与调用完全指南

在大模型(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+ 原生 HttpClientOkHttp

方案 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 pipeClientAbortException解决办法:在 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,注意异步发送和超时设置。

  • 客户端:推荐使用 OkHttpEventSource,解析优雅且自带重连机制。

  • 运维:务必检查 Nginx 的 proxy_buffering off 配置。

  • 大模型场景:注意处理 [DONE] 结束标志和 JSON 流式解析。

掌握了以上知识,你就可以在 Java 项目中轻松实现类似 ChatGPT 的“打字机”流式输出效果了!

开源文件分享系统 FileCodeBox 安装教程 2026-07-03
七月的一些事情 2026-07-28

评论区