WebClient 流式解析与 SSE 接口
约 540 字大约 2 分钟
布欧-Lewyon
2026-06-08
首页 › Agent Spring Boot › WebFlux 与流式 API
核心链路
agent-web-flux 的 ChatService.streamChat() 做三件事:
- 组装
stream: true请求体 WebClient拉上游 SSE,得到Flux<String>(每行一个 JSON chunk)- 解析
choices[0].delta.content,拼成下游Flux<String>推给浏览器
return webClient.post()
.uri("/chat/completions")
.bodyValue(requestBody)
.retrieve()
.bodyToFlux(String.class)
.filter(line -> !line.isBlank() && !"[DONE]".equals(line.trim()))
.flatMap(data -> {
String content = extractContent(data);
return (content != null && !content.isEmpty())
? Flux.just(content)
: Flux.empty();
})
.concatWith(Flux.just("[DONE]"))
.onErrorResume(e -> { /* 错误转 SSE 文本 */ });JSON 解析
DeepSeek 流式 chunk 示例:
{"choices":[{"delta":{"content":"你好"},"index":0}]}extractContent 用 Jackson 读取:
private String extractContent(String json) {
JsonNode root = objectMapper.readTree(json);
JsonNode choices = root.get("choices");
if (choices != null && choices.isArray() && !choices.isEmpty()) {
JsonNode delta = choices.get(0).get("delta");
if (delta != null) {
JsonNode content = delta.get("content");
if (content != null && content.isTextual()) {
return content.asText();
}
}
}
return null;
}错误友好降级
.onErrorResume(e -> {
String errMsg = e instanceof WebClientResponseException wcre
? "API 错误: " + wcre.getStatusCode() + " " + wcre.getResponseBodyAsString()
: "调用失败: " + e.getMessage();
return Flux.just(errMsg);
});把异常变成 SSE 文本,前端能展示,而不是直接断连。
前端消费(fetch + ReadableStream)
frontend/src/api/chat.ts 不用 EventSource,以便后续扩展 POST:
const response = await fetch(`/api/chat/stream?message=${encodeURIComponent(message)}`);
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 });
const lines = buffer.split('\n');
buffer = lines.pop() ?? '';
for (const line of lines) {
if (line.startsWith('data:')) {
const data = line.slice(5).trim();
if (data && data !== '[DONE]') onChunk(data);
}
}
}App.tsx 里把 chunk 追加到 assistant 消息,配合 AbortController 实现「停止生成」。
流式时序图
学习点
| 知识点 | 要点 |
|---|---|
bodyToFlux(String.class) | 按 SSE 行订阅上游 |
flatMap + Flux.empty() | 跳过无 content 的 heartbeat chunk |
concatWith("[DONE]") | 显式结束,前端好收尾 |
| 假流式 | 禁止拿完整答案再 split 模拟 |
易错点
- Nginx 缓冲:生产需
proxy_buffering off,否则 SSE 一次性吐出。 - Boot 4 Jackson:import 是
tools.jackson,不是com.fasterxml.jackson.databind(注解仍可能用com.fasterxml.jackson.annotation)。 - 空 chunk:首包可能只有
role,delta.content为空,要过滤。
小结
- 流式核心在
ChatService:bodyToFlux→ 解析 delta →concatWith([DONE])。 - 前端用 ReadableStream 解析
data:行,支持中止。 - 错误用
onErrorResume转成可见 SSE 消息。
上一节:项目初始化
下一节:响应式 API 与前端
