摘要
ai 对话接口如果一直等模型完整生成后再返回,用户会明显感觉“卡住了”。尤其是回答较长、模型推理较慢、网络链路较远时,接口可能需要数秒甚至十几秒才能返回第一段内容。流式输出的核心价值,就是让用户尽早看到模型正在生成内容,从而把等待过程变成连续反馈。
在 java 后端项目中,流式 ai 对话通常使用 sse 或 websocket 实现。sse 更简单,适合“服务器持续向浏览器推送文本片段”的场景;websocket 更适合双向实时通信和复杂协作场景。对大多数 ai 聊天接口来说,sse 已经足够。
本文基于 spring boot、spring webflux 和 spring ai,介绍如何实现一个可运行的流式 ai 对话接口。内容包括流式响应原理、sse 协议、spring ai 的 chatclient.stream()、后端接口设计、前端消费方式、异常处理、取消生成、上下文保存、token 统计和生产实践。
说明:spring ai 和 spring boot 都在持续迭代,实际项目要锁定版本,并以当前使用版本的官方文档为准。本文重点放在设计思路和常见实现方式。
一、背景与问题
1. 非流式接口的用户体验问题
前一篇文章已经实现了基于 spring ai 的基本模型调用。最简单的接口通常长这样:
@postmapping("/chat")
public chatresponse chat(@requestbody chatrequest request) {
string answer = chatclient.prompt()
.user(request.message())
.call()
.content();
return new chatresponse(answer);
}
这种写法的问题是:后端必须等模型生成完整结果后,才能把响应返回给前端。
调用链路如下:
浏览器提交问题
↓
spring boot 接收请求
↓
调用大模型
↓
模型生成完整答案
↓
后端一次性返回
↓
浏览器展示内容
如果模型 8 秒后才完成输出,那么用户前 8 秒什么都看不到。
2. 流式输出的体验差异
流式输出把完整回答拆成许多小片段:
第 0.3 秒:你好 第 0.6 秒:,这个问题 第 0.9 秒:可以从 第 1.2 秒:三个方面理解 ...
前端可以边接收边渲染:
用户提问 ↓ 后端建立流式连接 ↓ 模型生成一个片段 ↓ 后端推送一个片段 ↓ 前端追加到页面 ↓ 直到生成完成
用户感受到的是“模型正在回答”,而不是“接口一直没反应”。
3. 为什么 java 项目常用 sse
ai 聊天的主流程通常是:
客户端发送一次问题 服务端持续返回生成片段
这是一种单向持续推送。sse 正好适合这个场景。
| 方案 | 特点 | 适合场景 |
|---|---|---|
| 普通 http | 一次请求,一次响应 | 短回答、结构化任务 |
| sse | 一次请求,服务端持续推送 | ai 对话、日志流、任务进度 |
| websocket | 双向长连接 | 协同编辑、实时游戏、多端通信 |
sse 的优势是:
- 基于 http,接入简单;
- 浏览器原生支持
eventsource; - 服务端实现成本低;
- 适合文本流;
- 更容易和网关、日志、鉴权体系结合。
但 sse 也有限制:
- 浏览器原生
eventsource只支持 get; - 如果要 post 请求体,通常需要使用
fetch读取流; - 不适合复杂双向通信;
- 代理和网关可能需要关闭响应缓冲。
二、核心概念
1. token、片段和完整回答
大模型并不是一次性生成完整答案,而是逐步生成 token。token 可以理解为模型处理文本的基本单位,它可能是一个字、一个词、一个标点或词的一部分。
流式接口通常不会逐个 token 暴露给业务,而是返回一个个文本片段:
flux<string> ├─ "你好" ├─ "," ├─ "可以" ├─ "这样" └─ "理解..."
后端需要把这些片段组织成浏览器可以消费的事件。
2. sse 是什么
sse,全称 server-sent events,是一种服务端向客户端持续推送事件的协议。响应的 content-type 通常是:
text/event-stream
一个事件可以长这样:
event: message data: 你好 event: message data: ,这是一个流式回答 event: done data: [done]
每个事件之间用空行分隔。
3. spring webflux 中的流
spring webflux 使用 reactor 的 flux 表示 0 到 n 个异步数据项:
flux<string> stream = flux.just("a", "b", "c");
对于 sse,后端可以返回:
flux<serversentevent<string>>
spring 会持续把 flux 中的数据写入 http 响应。
4. spring ai 的流式接口
根据 spring ai 官方文档,chatclient 支持同步和流式两种调用方式。流式文本可以这样获取:
flux<string> output = chatclient.prompt()
.user("tell me a joke")
.stream()
.content();
如果需要更完整的响应元数据,也可以流式获取 chatresponse:
flux<chatresponse> response = chatclient.prompt()
.user("tell me a joke")
.stream()
.chatresponse();
实际业务中,最常见的是先使用 content() 拿文本片段,再封装成 sse 事件。
三、工作原理
1. 整体架构
一个基础流式对话架构如下:
浏览器 │ │ post /api/ai/chat/stream ▼ spring boot api ├─ 参数校验 ├─ 鉴权 ├─ 会话加载 ├─ prompt 组装 ├─ spring ai 流式调用 ├─ sse 事件封装 ├─ 错误处理 └─ 消息落库 │ ▼ 大模型服务
注意,流式对话不是“把模型流直接透传给前端”这么简单。生产接口通常还要处理:
- 用户是否有权限访问会话;
- 输入是否超过长度限制;
- 是否需要注入历史上下文;
- 是否需要敏感词或提示词安全检查;
- 流式内容如何保存;
- 用户中途断开时如何取消模型调用;
- 生成失败后如何通知前端;
- 如何统计 token、延迟和费用。
2. 后端事件类型
建议不要只返回纯文本片段,而是设计事件类型:
| 事件 | 含义 |
|---|---|
start | 服务端已开始处理 |
delta | 模型生成的文本片段 |
error | 生成失败 |
done | 本轮生成完成 |
usage | token、耗时、模型等统计信息 |
sse 数据可以是 json:
{
"conversationid": "c_1001",
"messageid": "m_2001",
"content": "你好"
}这样前端更容易处理状态。
3. 为什么要保存完整回答
流式输出过程中,前端看到的是片段,但后端最终仍然需要保存完整回答。
原因包括:
- 下轮对话需要上下文;
- 用户刷新页面要能看到历史;
- 后台需要做质量评估;
- 出错时需要排查;
- 需要统计生成内容长度;
- 需要支持会话导出。
常见做法是在服务端聚合文本:
flux<string> ↓ doonnext 追加到 stringbuilder ↓ dooncomplete 保存 assistant 消息
但要注意线程安全、异常场景和用户主动取消。
四、实战示例
1. 创建依赖
示例使用 spring boot、spring webflux 和 spring ai。不同版本 starter 名称可能变化,实际项目以当前文档为准。
maven 示例:
<dependencies>
<dependency>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter-webflux</artifactid>
</dependency>
<dependency>
<groupid>org.springframework.ai</groupid>
<artifactid>spring-ai-starter-model-openai</artifactid>
</dependency>
</dependencies>配置示例:
spring:
ai:
openai:
api-key: ${openai_api_key}
chat:
options:
model: ${ai_chat_model:gpt-4.1-mini}
temperature: 0.7真实项目不要把密钥写进配置文件,应使用环境变量、密钥管理系统或配置中心。
2. 定义请求对象
public record streamchatrequest(
string conversationid,
string message
) {
}
定义事件数据:
public record streamchatevent(
string conversationid,
string messageid,
string type,
string content
) {
public static streamchatevent start(string conversationid, string messageid) {
return new streamchatevent(conversationid, messageid, "start", "");
}
public static streamchatevent delta(string conversationid, string messageid, string content) {
return new streamchatevent(conversationid, messageid, "delta", content);
}
public static streamchatevent done(string conversationid, string messageid) {
return new streamchatevent(conversationid, messageid, "done", "");
}
public static streamchatevent error(string conversationid, string messageid, string content) {
return new streamchatevent(conversationid, messageid, "error", content);
}
}
3. 配置 chatclient
@configuration
public class aiconfig {
@bean
chatclient chatclient(chatclient.builder builder) {
return builder
.defaultsystem("""
你是一个 java 后端技术助手。
回答要准确、简洁,并优先给出可落地的工程建议。
如果用户的问题缺少关键上下文,先提出澄清问题。
""")
.build();
}
}
系统提示词可以放在配置中心或 prompt 模板中,不建议散落在 controller 里。
4. 实现流式服务
@service
public class aistreamservice {
private final chatclient chatclient;
public aistreamservice(chatclient chatclient) {
this.chatclient = chatclient;
}
public flux<streamchatevent> stream(streamchatrequest request) {
string conversationid = request.conversationid();
string assistantmessageid = uuid.randomuuid().tostring();
stringbuilder fullcontent = new stringbuilder();
flux<streamchatevent> start = flux.just(
streamchatevent.start(conversationid, assistantmessageid)
);
flux<streamchatevent> content = chatclient.prompt()
.user(request.message())
.stream()
.content()
.map(chunk -> {
fullcontent.append(chunk);
return streamchatevent.delta(conversationid, assistantmessageid, chunk);
})
.dooncomplete(() -> {
string answer = fullcontent.tostring();
// todo 保存 assistant 消息
// messagerepository.saveassistantmessage(conversationid, assistantmessageid, answer);
});
flux<streamchatevent> done = flux.just(
streamchatevent.done(conversationid, assistantmessageid)
);
return flux.concat(start, content, done)
.onerrorresume(ex -> flux.just(
streamchatevent.error(conversationid, assistantmessageid, "生成失败,请稍后重试")
));
}
}
这里有几个关键点:
stream().content()返回文本片段流;map中把片段转成业务事件;stringbuilder聚合完整回答;dooncomplete保存最终内容;onerrorresume把异常转为错误事件。
5. controller 返回 sse
@restcontroller
@requestmapping("/api/ai")
public class aistreamcontroller {
private final aistreamservice aistreamservice;
public aistreamcontroller(aistreamservice aistreamservice) {
this.aistreamservice = aistreamservice;
}
@postmapping(
value = "/chat/stream",
produces = mediatype.text_event_stream_value
)
public flux<serversentevent<streamchatevent>> stream(
@requestbody streamchatrequest request
) {
return aistreamservice.stream(request)
.map(event -> serversentevent.<streamchatevent>builder()
.event(event.type())
.data(event)
.build());
}
}
produces = mediatype.text_event_stream_value 表示接口会持续输出 sse。
6. 前端使用 fetch 读取流
如果接口使用 post,请不要直接用 eventsource,可以用 fetch 读取流:
async function streamchat(message) {
const response = await fetch("/api/ai/chat/stream", {
method: "post",
headers: {
"content-type": "application/json"
},
body: json.stringify({
conversationid: "c_1001",
message
})
});
const reader = response.body.getreader();
const decoder = new textdecoder("utf-8");
let buffer = "";
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const events = buffer.split("\n\n");
buffer = events.pop() || "";
for (const rawevent of events) {
const dataline = rawevent
.split("\n")
.find(line => line.startswith("data:"));
if (!dataline) continue;
const json = dataline.substring(5).trim();
const event = json.parse(json);
if (event.type === "delta") {
appendmessagetext(event.content);
}
if (event.type === "done") {
markmessagedone();
}
if (event.type === "error") {
showerror(event.content);
}
}
}
}
更完整的前端还要处理:
- loading 状态;
- 取消生成;
- 网络断开;
- 重试;
- markdown 渲染;
- 代码块高亮;
- xss 防护。
7. 支持用户取消生成
浏览器端可以使用 abortcontroller:
const controller = new abortcontroller();
fetch("/api/ai/chat/stream", {
method: "post",
signal: controller.signal,
headers: {
"content-type": "application/json"
},
body: json.stringify({ conversationid, message })
});
function stopgeneration() {
controller.abort();
}后端可以在流取消时记录状态:
return aistreamservice.stream(request)
.dooncancel(() -> {
// todo 标记本轮消息被用户中断
})
.map(event -> serversentevent.<streamchatevent>builder()
.event(event.type())
.data(event)
.build());取消生成很重要,因为用户可能发现问题问错了,或者模型输出太长。
8. 保存用户消息和助手消息
对话系统通常需要两类消息:
user 用户输入 assistant 模型回答
可以设计消息表:
create table ai_message (
id bigint primary key,
conversation_id bigint not null,
role varchar(32) not null,
content text not null,
status varchar(32) not null,
model varchar(128),
created_at timestamp not null
);流式接口推荐流程:
1. 收到请求后先保存 user 消息 2. 创建 assistant 消息,状态为 generating 3. 流式返回 delta 4. 完成后更新 assistant 内容和状态 5. 失败或取消时更新状态
这样即使前端刷新页面,也能恢复会话状态。
五、常见问题与实践建议
1. 为什么本地能流式,线上变成一次性返回
常见原因是代理或网关开启了响应缓冲。
例如 nginx 需要关注:
location /api/ai/chat/stream {
proxy_pass http://backend;
proxy_http_version 1.1;
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 300s;
}否则后端虽然分片输出,代理可能等缓冲区满或响应结束后才统一返回。
2. 不要在流式接口里做耗时阻塞操作
流式链路应该尽量保持顺畅。不要在每个 token 片段里同步写数据库、调用远程接口或做复杂计算。
推荐:
- 片段只用于推送;
- 完整回答结束后再保存;
- 日志和统计异步处理;
- 高成本检查放在生成前或生成后。
3. 如何处理模型输出 markdown
ai 对话常包含 markdown。前端可以边接收边渲染,但要注意:
- 流式片段可能截断 markdown 语法;
- 代码块可能在中途不完整;
- html 内容必须做 xss 过滤;
- 表格和公式最好等完整后再增强渲染。
简单做法是:流式阶段先按纯文本追加,完成后再统一 markdown 渲染。
4. 是否每个片段都要落库
通常不建议每个片段都落库。
原因:
- 写入频率过高;
- 数据库压力大;
- 内容碎片化;
- 后续查询不方便。
更常见做法是:
流式过程中内存聚合 完成后保存完整消息 异常或取消时保存当前已生成内容
5. 上下文不能无限拼接
多轮对话要加载历史消息,但不能无限拼接。否则会导致:
- prompt 过长;
- 成本上升;
- 响应变慢;
- 超过模型上下文限制;
- 历史噪声影响回答。
建议策略:
- 最近 n 轮完整保留;
- 更早历史做摘要;
- 重要事实单独存储;
- 每次调用前估算上下文长度;
- 超限时优先移除低价值上下文。
6. 错误事件要明确
流式接口一旦开始写出响应,http 状态码就很难再表达业务失败。应通过 sse 事件通知前端:
{
"type": "error",
"content": "模型服务暂时不可用,请稍后重试"
}前端不要只依赖 http 状态码,也要监听业务事件。
六、进阶思考
1. 流式接口的生产级能力清单
一个生产级流式 ai 对话接口通常需要:
| 能力 | 说明 |
|---|---|
| 鉴权 | 用户只能访问自己的会话 |
| 限流 | 控制用户、租户、ip 的请求频率 |
| 并发控制 | 防止同一会话同时生成多条回答 |
| 取消生成 | 用户可以中断当前输出 |
| 内容安全 | 输入和输出都需要安全检查 |
| 上下文管理 | 控制历史消息和知识库内容 |
| 失败恢复 | 失败、取消、超时都要有消息状态 |
| 观测指标 | 记录首 token 延迟、总耗时、token 数 |
| 成本统计 | 按用户、租户、模型统计费用 |
| 网关配置 | 关闭缓冲,设置合理超时 |
2. 首 token 延迟比总耗时更影响体验
ai 流式接口有两个关键指标:
| 指标 | 含义 |
|---|---|
| 首 token 延迟 | 从请求发出到第一个片段返回的时间 |
| 总生成耗时 | 从请求发出到回答完成的时间 |
用户对首 token 延迟非常敏感。即使完整回答需要 20 秒,只要 1 秒内看到内容开始生成,体验也会好很多。
优化方向:
- prompt 不要过长;
- 知识库召回要快;
- 模型路由要合理;
- 避免生成前做过多同步操作;
- 常用系统提示词可以缓存;
- 长任务可以先返回
start事件。
3. sse 和 websocket 如何选择
建议:
| 场景 | 推荐 |
|---|---|
| 普通 ai 聊天 | sse |
| 知识库问答 | sse |
| 代码生成流式输出 | sse |
| 多人协作编辑 | websocket |
| 语音实时对话 | websocket 或专门实时协议 |
| 前端需要频繁发送控制消息 | websocket |
不要为了“高级”而默认上 websocket。简单的单向流,sse 更容易维护。
4. 如何和 rag 结合
如果流式回答基于知识库,流程通常变成:
用户提问 ↓ 改写查询 ↓ 向量召回 ↓ 重排序 ↓ 构造 prompt ↓ 流式生成答案 ↓ 返回引用来源
引用来源建议作为单独事件返回:
{
"type": "references",
"content": [
{
"title": "售后政策",
"url": "/docs/after-sales"
}
]
}不要把引用信息只混在自然语言里,否则前端很难结构化展示。
结论
流式 ai 对话的核心不是“把字符串一段段返回”,而是围绕用户体验、服务稳定性和工程可维护性设计一条完整链路。
在 spring boot 中,可以使用 spring webflux 返回 flux<serversentevent<?>>,再结合 spring ai 的 chatclient.stream().content() 获取模型生成片段。这个组合足以支撑大多数 ai 聊天、知识库问答和代码助手场景。
进入生产环境时,还需要补齐鉴权、限流、取消生成、上下文管理、消息落库、网关缓冲配置、错误事件、token 统计和监控指标。只有把这些工程细节处理好,流式对话才会从 demo 变成真正稳定可用的 ai 服务。
以上就是在springboot中实现流式ai对话功能的详细内容,更多关于springboot流式ai对话的资料请关注代码网其它相关文章!
发表评论