Spring Boot 3 SSE 流式输出:AI 对话打字机效果的完整实现(附 Vue 3 案例)

作者:忆笙智云官方 | 发布时间:2026-05-05 14:00 | 更新时间:2026-06-05 14:00

Spring Boot 3 SSE流式输出实现原理

引言

在AI对话场景中,大语言模型(LLM)的响应通常需要数秒甚至数十秒才能完整生成。如果等待全部内容生成完毕再返回,用户体验极差——用户面对长时间空白页面,不知道系统是否在正常工作。

SSE(Server-Sent Events)协议允许服务端向客户端单向推送数据流,是实现AI对话"打字机效果"的核心技术。本文将深入分析SSE协议原理、Spring Boot 3中的实现方案,以及AI对话流式输出的完整架构。

一、SSE协议原理

1.1 SSE vs WebSocket vs 轮询

特性 HTTP轮询 WebSocket SSE
通信方向 客户端→服务端 双向 服务端→客户端
协议 HTTP WS HTTP
连接开销 每次新建连接 一次握手 一次握手
断线重连 不需要 需手动实现 浏览器自动
数据格式 自定义 自定义 文本(text/event-stream)
浏览器支持 全部 现代浏览器 现代浏览器
适用场景 低频查询 实时双向通信 服务端推送

SSE的优势在于:基于标准HTTP协议,无需额外握手,浏览器原生支持断线重连,天然适合服务端向客户端单向推送数据的场景。

1.2 SSE消息格式

SSE遵循text/event-stream MIME类型,消息格式如下:

field: value


支持的字段:

id: 1
event: message
data: {"content": "你", "index": 0}

id: 2
event: message
data: {"content": "好", "index": 1}

id: 3
event: done
data: [DONE]

二、Spring Boot 3 SSE实现

2.1 SseEmitter核心类

Spring Boot 3提供了SseEmitter类,封装了SSE协议的细节:

sequenceDiagram
    participant Client as 客户端
    participant Controller as Spring Controller
    participant Emitter as SseEmitter
    participant Service as AI服务

    Client->>Controller: GET /ai/chat/stream
    Controller->>Emitter: 创建SseEmitter(timeout)
    Controller-->>Client: 返回SseEmitter(Content-Type: text/event-stream)

    loop 流式输出
        Service->>Emitter: emitter.send(data)
        Emitter-->>Client: 推送SSE事件
    end

    Service->>Emitter: emitter.complete()
    Emitter-->>Client: 关闭连接

2.2 基础实现

/**
 * AI对话控制器
 * 提供SSE流式对话接口
 */
@RestController
@RequestMapping("/ai/chat")
public class AiChatController {

    @Autowired
    private AiChatService aiChatService;

    /** SSE超时时间:5分钟 */
    private static final long SSE_TIMEOUT = 5 * 60 * 1000L;

    /**
     * 流式对话接口
     * @param message 用户消息
     * @return SseEmitter
     */
    @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public SseEmitter streamChat(@RequestParam String message) {
        // 创建SseEmitter,设置超时时间
        SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);

        // 注册超时回调
        emitter.onTimeout(() -> {
            emitter.complete();
        });

        // 注册错误回调
        emitter.onError(ex -> {
            emitter.complete();
        });

        // 异步执行AI对话
        CompletableFuture.runAsync(() -> {
            try {
                aiChatService.streamChat(message, emitter);
                emitter.complete();  // 正常完成
            } catch (Exception e) {
                emitter.completeWithError(e);  // 异常完成
            }
        });

        return emitter;
    }
}

2.3 AI服务流式调用

/**
 * AI对话服务
 * 调用大模型API并流式输出
 */
@Service
public class AiChatService {

    @Autowired
    private LlmClient llmClient;

    /**
     * 流式对话
     * @param message 用户消息
     * @param emitter SSE发射器
     */
    public void streamChat(String message, SseEmitter emitter) throws IOException {
        // 构建AI请求
        ChatRequest request = ChatRequest.builder()
                .message(message)
                .stream(true)  // 启用流式输出
                .build();

        // 调用大模型API,获取流式响应
        Flux<ChatChunk> responseFlux = llmClient.streamChat(request);

        // 订阅流式响应,逐块推送给客户端
        responseFlux.subscribe(
            chunk -> {
                try {
                    // 构建SSE事件
                    SseEmitter.SseEventBuilder event = SseEmitter.event()
                            .id(String.valueOf(chunk.getIndex()))
                            .data(chunk.getContent(), MediaType.APPLICATION_JSON);
                    emitter.send(event);
                } catch (IOException e) {
                    emitter.completeWithError(e);
                }
            },
            error -> {
                try {
                    // 发送错误事件
                    emitter.send(SseEmitter.event()
                            .name("error")
                            .data("AI服务异常"));
                } catch (IOException ignored) {
                }
                emitter.completeWithError(error);
            },
            () -> {
                try {
                    // 发送完成事件
                    emitter.send(SseEmitter.event()
                            .name("done")
                            .data("[DONE]"));
                } catch (IOException ignored) {
                }
            }
        );
    }
}

三、前端EventSource对接

3.1 原生EventSource

/**
 * 使用原生EventSource接收SSE流
 */
function startSSEChat(message, onMessage, onDone, onError) {
    const url = `/api/ai/chat/stream?message=${encodeURIComponent(message)}`;
    const eventSource = new EventSource(url);

    // 监听消息事件
    eventSource.addEventListener('message', (event) => {
        const data = JSON.parse(event.data);
        onMessage(data.content);
    });

    // 监听完成事件
    eventSource.addEventListener('done', (event) => {
        eventSource.close();
        onDone();
    });

    // 监听错误事件
    eventSource.addEventListener('error', (event) => {
        eventSource.close();
        onError(event);
    });

    // 连接错误处理
    eventSource.onerror = () => {
        eventSource.close();
        onError(new Error('SSE连接异常'));
    };

    return eventSource;
}

3.2 Vue 3组合式封装

/**
 * AI对话SSE组合式函数
 * 封装SSE连接、消息接收、状态管理
 */
import { ref } from 'vue';

export function useAiChat() {
    const messages = ref<ChatMessage[]>([]);
    const isStreaming = ref(false);
    const currentContent = ref('');
    let eventSource: EventSource | null = null;

    /**
     * 发送消息并接收流式响应
     */
    function sendMessage(content: string) {
        if (isStreaming.value) return;

        // 添加用户消息
        messages.value.push({ role: 'user', content });
        isStreaming.value = true;
        currentContent.value = '';

        const url = `/api/ai/chat/stream?message=${encodeURIComponent(content)}`;
        eventSource = new EventSource(url);

        eventSource.addEventListener('message', (event) => {
            const data = JSON.parse(event.data);
            currentContent.value += data.content;
        });

        eventSource.addEventListener('done', () => {
            // 流式输出完成,将完整AI消息加入列表
            messages.value.push({
                role: 'assistant',
                content: currentContent.value
            });
            currentContent.value = '';
            isStreaming.value = false;
            eventSource?.close();
        });

        eventSource.onerror = () => {
            isStreaming.value = false;
            eventSource?.close();
        };
    }

    /**
     * 停止生成
     */
    function stopGeneration() {
        eventSource?.close();
        if (currentContent.value) {
            messages.value.push({
                role: 'assistant',
                content: currentContent.value
            });
        }
        currentContent.value = '';
        isStreaming.value = false;
    }

    return { messages, isStreaming, currentContent, sendMessage, stopGeneration };
}

四、进阶:支持POST请求的SSE

原生EventSource只支持GET请求,无法携带请求体。在实际项目中,对话请求通常需要携带复杂的参数(会话ID、上下文、模型配置等),需要使用fetch API实现POST方式的SSE。

4.1 服务端:基于ResponseBodyEmitter

/**
 * 支持POST请求的流式对话接口
 */
@PostMapping(value = "/stream/post", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamChatPost(@RequestBody ChatRequestDTO request) {
    SseEmitter emitter = new SseEmitter(SSE_TIMEOUT);

    emitter.onTimeout(emitter::complete);
    emitter.onError(ex -> emitter.complete());

    CompletableFuture.runAsync(() -> {
        try {
            aiChatService.streamChat(request.getMessage(), emitter);
            emitter.complete();
        } catch (Exception e) {
            emitter.completeWithError(e);
        }
    });

    return emitter;
}

4.2 前端:fetch + ReadableStream

/**
 * 使用fetch API实现POST方式的SSE接收
 * 支持携带请求体和自定义Header
 */
async function fetchSSE(
    url: string,
    body: object,
    onMessage: (data: string) => void,
    onDone: () => void,
    onError: (error: Error) => void
) {
    const response = await fetch(url, {
        method: 'POST',
        headers: {
            'Content-Type': 'application/json',
            'Authorization': `Bearer ${getToken()}`
        },
        body: JSON.stringify(body)
    });

    if (!response.ok) {
        onError(new Error(`HTTP ${response.status}`));
        return;
    }

    const reader = response.body!.getReader();
    const decoder = new TextDecoder();
    let buffer = '';

    while (true) {
        const { done, value } = await reader.read();
        if (done) break;

        buffer += decoder.decode(value, { stream: true });

        // 按SSE格式解析(双换行分隔事件)
        const events = buffer.split('

');
        buffer = events.pop() || '';  // 保留未完成的部分

        for (const event of events) {
            const lines = event.split('
');
            let eventType = 'message';
            let data = '';

            for (const line of lines) {
                if (line.startsWith('event:')) {
                    eventType = line.substring(6).trim();
                } else if (line.startsWith('data:')) {
                    data = line.substring(5).trim();
                }
            }

            if (eventType === 'done') {
                onDone();
                return;
            }
            if (eventType === 'error') {
                onError(new Error(data));
                return;
            }
            if (data) {
                onMessage(data);
            }
        }
    }
    onDone();
}

五、完整时序图

sequenceDiagram
    actor User as 用户
    participant FE as 前端Vue
    participant API as 后端API
    participant Emitter as SseEmitter
    participant LLM as 大模型API

    User->>FE: 输入消息并发送
    FE->>API: POST /ai/chat/stream (SSE请求)

    API->>Emitter: 创建SseEmitter(5min)
    API-->>FE: 返回SseEmitter (text/event-stream)

    API->>LLM: 发起流式请求 (stream=true)

    loop 流式响应
        LLM-->>API: ChatChunk: "你"
        API->>Emitter: send(event: message, data: "你")
        Emitter-->>FE: SSE: data: "你"
        FE->>FE: 追加显示"你"

        LLM-->>API: ChatChunk: "好"
        API->>Emitter: send(event: message, data: "好")
        Emitter-->>FE: SSE: data: "好"
        FE->>FE: 追加显示"好"

        LLM-->>API: ChatChunk: "!"
        API->>Emitter: send(event: message, data: "!")
        Emitter-->>FE: SSE: data: "!"
        FE->>FE: 追加显示"!"
    end

    LLM-->>API: 流结束
    API->>Emitter: send(event: done, data: [DONE])
    Emitter-->>FE: SSE: event: done
    API->>Emitter: complete()

    FE->>FE: 关闭EventSource
    FE->>User: 显示完整AI回复

六、关键问题与解决方案

6.1 连接超时

SSE连接可能因网络问题或代理超时断开。解决方案:

// 设置合理的超时时间
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L);

// 定期发送心跳包保持连接
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(() -> {
    try {
        emitter.send(SseEmitter.event().name("heartbeat").data(""));
    } catch (IOException e) {
        scheduler.shutdown();
    }
}, 0, 30, TimeUnit.SECONDS);

6.2 背压(Backpressure)控制

当客户端处理速度慢于服务端推送速度时,需要背压控制:

// 使用ResponseBodyEmitter的alreadyComplete状态判断
if (emitter instanceof SseEmitter) {
    try {
        emitter.send(data);
    } catch (IllegalStateException e) {
        // 连接已关闭,停止发送
        return;
    }
}

6.3 Nginx代理配置

Nginx默认会缓冲代理响应,导致SSE消息无法实时推送:

# Nginx SSE代理配置
location /ai/chat/ {
    proxy_pass http://backend;
    proxy_http_version 1.1;

    # 禁用响应缓冲,确保SSE实时推送
    proxy_buffering off;
    proxy_cache off;

    # 关闭连接超时
    proxy_read_timeout 300s;
    proxy_send_timeout 300s;

    # 不使用chunked编码
    proxy_set_header Connection '';
}

结论与建议

核心要点

  1. SSE是AI对话流式输出的最佳选择:基于HTTP、浏览器原生支持、自动重连
  2. Spring Boot 3的SseEmitter封装了SSE协议细节,使用简单
  3. POST方式SSE需要前端使用fetch + ReadableStream,但功能更强大

实践建议

  1. 超时设置:SSE超时时间应大于大模型最长响应时间,建议5分钟以上
  2. 心跳机制:长时间无数据推送时发送心跳包,防止中间代理关闭连接
  3. 错误处理:前端需要处理网络断开、服务端异常等场景,提供重试机制
  4. Nginx配置:必须关闭proxy_buffering,否则SSE消息会被缓冲
  5. 资源释放:SSE连接关闭时,务必取消大模型API的流式订阅,释放资源

相关资源