1. 项目概述
最近在开发一个基于大语言模型的聊天应用时,遇到了一个常见的需求:如何实现类似ChatGPT那样的流式响应效果?传统的HTTP请求-响应模式需要等待整个响应生成完毕才能返回给客户端,这在处理大语言模型的长文本生成时会导致明显的延迟。经过调研,我最终选择了Spring框架提供的Server-Sent Events(SSE)技术来实现流式传输,结合LangChain4j的大语言模型集成能力,打造了一个完整的流式聊天解决方案。
这个方案的核心优势在于:
- 实时性:每个token生成后立即推送到前端,无需等待完整响应
- 低延迟:相比轮询方案,SSE基于HTTP长连接,减少了不必要的网络开销
- 兼容性:SSE是标准Web技术,所有现代浏览器都支持
- 资源效率:服务端可以更好地控制连接和资源释放
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 技术选型与架构设计
2.1 技术栈组成
项目采用的技术栈主要包括:
- Spring Boot 3.5.0:提供基础的Web框架和SSE支持
- LangChain4j 1.9.1:Java版LangChain,集成大语言模型能力
- 通义千问:通过OpenAI兼容接口接入的大语言模型
- Server-Sent Events:实现服务端到客户端的单向实时通信
2.2 架构设计思路
整个系统采用典型的三层架构:
- 表现层:HTML前端页面 + SSE客户端实现
- 业务逻辑层:处理流式聊天请求,协调模型调用
- 数据访问层:LangChain4j与通义千问模型的交互
关键设计决策:
- 同时提供低级API和高级API两种实现方式,满足不同场景需求
- 使用SseEmitter作为SSE的技术实现载体
- 设置30分钟的超时时间,适应长对话场景
- 前端采用EventSource API处理SSE连接
3. 环境准备与项目配置
3.1 项目初始化
首先创建一个标准的Spring Boot项目,使用Maven进行依赖管理。关键依赖包括:
xml复制<dependencies>
<!-- Spring Boot基础依赖 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- LangChain4j核心库 -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j</artifactId>
<version>1.9.1</version>
</dependency>
<!-- OpenAI兼容接口支持 -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-open-ai-spring-boot-starter</artifactId>
<version>1.9.1-beta17</version>
</dependency>
<!-- 响应式编程支持 -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-reactor</artifactId>
<version>1.9.1-beta17</version>
</dependency>
</dependencies>
3.2 配置文件设置
在application.yml中配置LangChain4j与通义千问的连接参数:
yaml复制langchain4j:
open-ai:
chat-model:
api-key: your-api-key
model-name: qwen-plus
base-url: https://dashscope.aliyuncs.com/compatible-mode/v1
注意:实际部署时应将API密钥存储在安全的地方,如环境变量或密钥管理服务,不要直接硬编码在配置文件中。
4. 前端实现
4.1 HTML页面结构
前端页面主要包含三个核心区域:
- 输入区:文本输入框和发送按钮
- 状态区:显示当前连接状态
- 响应区:实时显示模型返回的内容
html复制<div class="container">
<h1>流式聊天演示</h1>
<div class="input-area">
<textarea id="promptInput" placeholder="请输入您的问题..."></textarea>
<button id="sendBtn" onclick="sendMessage()">低级API发送</button>
<button id="highLevelSendBtn" onclick="sendHighLevelMessage()">高级API发送</button>
</div>
<div class="status">
<span id="statusText">就绪</span>
<span id="loading" class="loading" style="display: none;"></span>
</div>
<div class="response-area" id="responseArea"></div>
</div>
4.2 JavaScript事件处理
前端使用EventSource API处理SSE连接,关键逻辑包括:
javascript复制function sendMessage() {
const prompt = document.getElementById('promptInput').value.trim();
if (!prompt) return;
disableAllButtons();
document.getElementById('statusText').textContent = '正在发送请求...';
document.getElementById('loading').style.display = 'inline-block';
document.getElementById('responseArea').textContent = '';
eventSource = new EventSource(`/api/streamingChat/low-level-sse-streaming-chat?prompt=${encodeURIComponent(prompt)}&userId=1`);
eventSource.onmessage = function(event) {
const responseArea = document.getElementById('responseArea');
responseArea.textContent += event.data;
responseArea.scrollTop = responseArea.scrollHeight;
};
eventSource.onerror = function(error) {
console.error('SSE Error:', error);
handleConnectionError();
};
}
提示:在实际项目中,应该添加更多的错误处理和重试逻辑,确保网络波动时的用户体验。
5. 低级API实现
5.1 配置StreamingChatModel
首先需要配置LangChain4j的流式聊天模型:
java复制@Configuration
@ConfigurationProperties(prefix = "langchain4j.open-ai.chat-model")
public class LangChain4jConfig {
private String apiKey;
private String modelName;
private String baseUrl;
@Bean
public StreamingChatModel streamingChatModel() {
return OpenAiStreamingChatModel.builder()
.apiKey(apiKey)
.modelName(modelName)
.baseUrl(baseUrl)
.build();
}
}
5.2 服务层实现
服务层负责处理流式聊天的核心逻辑:
java复制@Service
@Slf4j
public class StreamingChatService {
@Autowired
private StreamingChatModel streamingChatModel;
public SseEmitter lowLevelStreamingChat(String prompt) {
SseEmitter emitter = new SseEmitter(30 * 60 * 1000L);
streamingChatModel.chat(prompt, new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
try {
emitter.send(partialResponse);
} catch (Exception e) {
emitter.completeWithError(e);
}
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
emitter.complete();
}
@Override
public void onError(Throwable throwable) {
emitter.completeWithError(throwable);
}
});
return emitter;
}
}
5.3 控制器层
控制器层提供REST接口:
java复制@RestController
@RequestMapping("/api/streamingChat")
public class StreamingChatController {
@Autowired
private StreamingChatService streamingChatService;
@GetMapping("/low-level-sse-streaming-chat")
public SseEmitter lowLevelStreamingChat(@RequestParam String prompt, @RequestParam int userId) {
return streamingChatService.lowLevelStreamingChat(prompt);
}
}
6. 高级API实现
6.1 定义StreamingAssistant接口
高级API首先需要定义一个服务接口:
java复制public interface StreamingAssistant {
@SystemMessage("你是资深中国历史学者,回答问题的风格是清晰简洁")
TokenStream chat(@UserMessage String message);
}
6.2 配置StreamingAssistant Bean
在配置类中添加:
java复制@Bean
public StreamingAssistant streamingAssistant(StreamingChatModel streamingChatModel) {
return AiServices.builder(StreamingAssistant.class)
.streamingChatModel(streamingChatModel)
.chatMemory(MessageWindowChatMemory.withMaxMessages(10))
.build();
}
6.3 服务层实现
高级API的服务实现:
java复制public SseEmitter highLevelStreamingChat(String prompt) {
SseEmitter emitter = new SseEmitter(30 * 60 * 1000L);
try {
TokenStream tokenStream = streamingAssistant.chat(prompt);
tokenStream.onPartialResponse(token -> {
try {
emitter.send(token);
} catch (Exception e) {
emitter.completeWithError(e);
}
});
tokenStream.onCompleteResponse(completeResponse -> {
emitter.complete();
});
tokenStream.onError(throwable -> {
emitter.completeWithError(throwable);
});
tokenStream.start();
} catch (Exception e) {
emitter.completeWithError(e);
}
return emitter;
}
7. 性能优化与问题排查
7.1 性能优化建议
- 连接池管理:对于高并发场景,考虑使用连接池管理SSE连接
- 心跳机制:定期发送心跳包保持连接活跃
- 超时调整:根据实际场景调整超时时间
- 背压处理:当客户端处理速度跟不上服务端时,应有相应的背压策略
7.2 常见问题排查
-
连接意外断开
- 检查服务端日志是否有异常
- 确认网络环境是否稳定
- 验证客户端EventSource的实现是否正确
-
响应延迟
- 检查模型API的响应时间
- 确认服务端是否有阻塞操作
- 检查网络延迟
-
内存泄漏
- 确保所有SSE连接都能正确关闭
- 监控JVM内存使用情况
- 定期检查并释放闲置连接
8. 扩展思考
8.1 功能扩展方向
- 对话历史管理:集成聊天记忆功能,实现多轮对话
- 流式文件处理:扩展支持流式处理文件内容
- 多模型切换:支持动态切换不同的大语言模型
- 速率限制:实现基于用户或IP的速率限制
8.2 技术深度探索
- WebSocket对比:与WebSocket方案进行性能对比
- HTTP/2优势:探索HTTP/2对SSE性能的提升
- 边缘计算:考虑在边缘节点部署流式处理逻辑
- 自适应流控:根据网络状况动态调整流式传输速率
在实际项目中,我发现流式传输的实现细节会显著影响最终用户体验。特别是在移动网络环境下,需要考虑更复杂的网络状况处理。建议在正式上线前进行充分的压力测试和异常场景测试,确保系统在各种条件下都能稳定运行。
