LangChain4j实战五:响应流式传输
·
本篇功能由Spring 提供的Server-Sent Events 实现,本质是HTTP 长连接 + 流式文本
1. 准备工作
1.1. 构建项目,添加pom文件
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>3.5.0</version>
<relativePath/>
</parent>
<groupId>cn.cjc</groupId>
<artifactId>springboot-ai</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- https://mvnrepository.com/artifact/dev.langchain4j/langchain4j -->
<!-- LangChain4j 核心库 -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j</artifactId>
<version>1.9.1</version>
</dependency>
<!-- AiService依赖 -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-spring-boot-starter</artifactId>
<version>1.0.1-beta6</version>
</dependency>
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-core</artifactId>
<version>1.9.1</version>
</dependency>
<!-- LangChain4j OpenAI支持(可用于通义千问的OpenAI兼容接口) -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-open-ai-spring-boot-starter</artifactId>
<version>1.9.1-beta17</version>
</dependency>
<!-- 集成原生阿里云通义千问 (DashScope) -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-community-dashscope-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>
<!-- mapdb-->
<dependency>
<groupId>org.mapdb</groupId>
<artifactId>mapdb</artifactId>
<version>3.0.9</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.22</version>
</dependency>
</dependencies>
</project>
1.2. 配置文件
spring:
application:
name: springboot-ai
main:
allow-bean-definition-overriding: true
langchain4j:
open-ai:
chat-model:
api-key: ******
model-name: qwen-plus
base-url: https://dashscope.aliyuncs.com/compatible-mode/v1
2. 前端html
- 新增文件/resources/static/index.html
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>流式聊天演示</title>
<style>
body {
font-family: Arial, sans-serif;
max-width: 800px;
margin: 0 auto;
padding: 20px;
background-color: #f5f5f5;
}
.container {
background-color: white;
border-radius: 8px;
padding: 20px;
box-shadow: 0 2px 10px rgba(0,0,0,0.1);
}
h1 {
text-align: center;
color: #333;
}
.input-area {
margin-bottom: 20px;
}
textarea {
width: 100%;
height: 100px;
padding: 10px;
border: 1px solid #ddd;
border-radius: 4px;
font-size: 16px;
resize: vertical;
}
button {
background-color: #4CAF50;
color: white;
border: none;
padding: 10px 20px;
font-size: 16px;
border-radius: 4px;
cursor: pointer;
margin-top: 10px;
}
button:hover {
background-color: #45a049;
}
button:disabled {
background-color: #cccccc;
cursor: not-allowed;
}
.response-area {
border: 1px solid #ddd;
border-radius: 4px;
padding: 15px;
min-height: 200px;
background-color: #f9f9f9;
font-family: 'Courier New', monospace;
white-space: pre-wrap;
}
.status {
margin-top: 10px;
color: #666;
font-size: 14px;
}
.loading {
display: inline-block;
width: 20px;
height: 20px;
border: 3px solid #f3f3f3;
border-top: 3px solid #4CAF50;
border-radius: 50%;
animation: spin 1s linear infinite;
}
@keyframes spin {
0% { transform: rotate(0deg); }
100% { transform: rotate(360deg); }
}
</style>
</head>
<body>
<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>
<script>
let eventSource = null;
function disableAllButtons() {
document.getElementById('sendBtn').disabled = true;
document.getElementById('highLevelSendBtn').disabled = true;
document.getElementById('promptInput').disabled = true;
}
function enableAllButtons() {
document.getElementById('sendBtn').disabled = false;
document.getElementById('highLevelSendBtn').disabled = false;
document.getElementById('promptInput').disabled = false;
}
function sendMessage() {
const prompt = document.getElementById('promptInput').value.trim();
if (!prompt) {
alert('请输入问题!');
return;
}
// 禁用按钮和输入框
disableAllButtons();
document.getElementById('statusText').textContent = '正在发送低级API请求...';
document.getElementById('loading').style.display = 'inline-block';
document.getElementById('responseArea').textContent = '';
// 使用Server-Sent Events建立连接(低级API)
eventSource = new EventSource(`/api/streamingChat/low-level-sse-streaming-chat?prompt=${encodeURIComponent(prompt)}&userId=1`);
setupEventSourceListeners();
}
function sendHighLevelMessage() {
const prompt = document.getElementById('promptInput').value.trim();
if (!prompt) {
alert('请输入问题!');
return;
}
// 禁用按钮和输入框
disableAllButtons();
document.getElementById('statusText').textContent = '正在发送高级API请求...';
document.getElementById('loading').style.display = 'inline-block';
document.getElementById('responseArea').textContent = '';
// 使用Server-Sent Events建立连接(高级API)
eventSource = new EventSource(`/api/streamingChat/high-level-sse-streaming-chat?prompt=${encodeURIComponent(prompt)}&userId=1`);
setupEventSourceListeners();
}
function setupEventSourceListeners() {
// 处理接收到的消息
eventSource.onmessage = function(event) {
const data = event.data;
if (data) {
const responseArea = document.getElementById('responseArea');
responseArea.textContent += data;
responseArea.scrollTop = responseArea.scrollHeight; // 自动滚动到底部
}
};
// 处理错误
eventSource.onerror = function(error) {
console.error('SSE Error:', error);
document.getElementById('statusText').textContent = '连接错误';
document.getElementById('loading').style.display = 'none';
enableAllButtons();
if (eventSource.readyState !== EventSource.CLOSED) {
eventSource.close();
}
};
// 处理连接关闭
eventSource.onclose = function() {
document.getElementById('statusText').textContent = '会话结束';
document.getElementById('loading').style.display = 'none';
enableAllButtons();
};
}
// 监听窗口关闭事件,确保关闭连接
window.addEventListener('beforeunload', function() {
if (eventSource && eventSource.readyState !== EventSource.CLOSED) {
eventSource.close();
}
});
</script>
</body>
</html>
- 点击按钮后,通过创建EventSource的方式和后端建立长连接,低级API和高级API的分别有不同的path
- 建立连接后,收到后台数据时,和responseArea的已有内容拼接起来再重新展示,效果就是只要后台数据不断,前端的内容就不断增加
3. 低级API
3.1. LangChain4j配置类
@Data
@Configuration
@ConfigurationProperties(prefix = "langchain4j.open-ai.chat-model")
public class LangChain4jConfig {
private String apiKey;
private String modelName;
private String baseUrl;
/**
* 创建并配置StreamingChatModel实例(使用通义千问的OpenAI兼容接口)
*/
@Bean
public StreamingChatModel streamingChatModel() {
return OpenAiStreamingChatModel.builder()
.apiKey(apiKey)
.modelName(modelName)
.baseUrl(baseUrl)
.build();
}
}
3.2. 服务类
streamingChatModel.chat方法的第二个入参,是StreamingChatResponseHandler的实现类,每当大模型返回token的时候,onPartialResponse方法就会被调用,传入最新的内容,所以只要在onPartialResponse方法中把新的内容给到前端即可,也就是执行emitter.send
package cn.cjc.ai.service.impl;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@Service
@Slf4j
public class StreamingChatService {
@Autowired
private StreamingChatModel streamingChatModel;
/**
* SSE流式聊天方法(用于网页实时显示)
* @param prompt 提示词
* @return SseEmitter实例,用于向客户端发送流式响应
*/
public SseEmitter lowLevelStreamingChat(String prompt) {
// 创建SseEmitter,设置超时时间为30分钟
SseEmitter emitter = new SseEmitter(30 * 60 * 1000L);
// 使用streamingChatModel进行流式聊天
streamingChatModel.chat(prompt, new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
if (partialResponse == null || partialResponse.trim().isEmpty()) {
log.warn("空的部分响应,忽略");
return;
}
try {
// 发送部分响应到客户端
emitter.send(partialResponse);
log.info("partialResponse: {}", partialResponse);
} catch (Exception e) {
log.error("发送部分响应失败: {}", e.getMessage(), e);
emitter.completeWithError(e);
}
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
try {
// 发送完成标志
emitter.send(SseEmitter.event().name("complete"));
// 完成响应
emitter.complete();
log.info("completeResponse: {}", completeResponse);
} catch (Exception e) {
log.error("发送完成响应失败: {}", e.getMessage(), e);
emitter.completeWithError(e);
}
}
@Override
public void onError(Throwable throwable) {
log.error("流式聊天发生错误: {}", throwable.getMessage(), throwable);
emitter.completeWithError(throwable);
}
});
// 设置超时处理
emitter.onTimeout(() -> {
log.error("SSE连接超时");
emitter.completeWithError(new RuntimeException("连接超时"));
});
// 设置完成处理
emitter.onCompletion(() -> {
log.info("SSE连接完成");
});
return emitter;
}
}
3.3. controller类
- SseEmitter是Spring 提供的 Server-Sent Events 实现的常规做法,这样就能做到把数据实时推送到前端
@RestController
@RequestMapping("/api/streamingChat")
@Slf4j
public class StreamingChatController {
@Autowired
private StreamingChatService streamingChatService;
/**
* SSE流式聊天接口(用于网页实时显示)
* @param prompt 提示词
* @param userId 用户ID
* @return SseEmitter实例,用于向客户端发送流式响应
*/
@GetMapping("/low-level-sse-streaming-chat")
public SseEmitter lowLevelStreamingChat(@RequestParam(name = "prompt") String prompt,
@RequestParam(name = "userId") int userId) {
log.info("收到来自用户[{}]的请求, 提示词: {}", userId, prompt);
try {
// 调用QwenService的流式聊天方法
return streamingChatService.lowLevelStreamingChat(prompt);
} catch (Exception e) {
// 捕获异常并返回错误的SseEmitter
SseEmitter emitter = new SseEmitter();
emitter.completeWithError(e);
return emitter;
}
}
}
4. 高级API
4.1. StreamingAssistant类
- 返回值必须是
TokenStream接口,该接口的onPartialResponse方法用来注册大模型返回token时候的回调
import dev.langchain4j.service.SystemMessage;
import dev.langchain4j.service.TokenStream;
import dev.langchain4j.service.UserMessage;
/**
* 流式助手接口
*/
public interface StreamingAssistant {
@SystemMessage("你是资深中国历史学者,回答问题的风格是清晰简洁")
TokenStream chat(@UserMessage String message);
}
4.2. 基于配置类新增实例
- 在配置类中增加一个bean的实例化,AiServices.builder是典型的高级API创建方式,还要注意绑定模型实例的方法是streamingChatModel
@Bean
public StreamingAssistant streamingAssistant(StreamingChatModel streamingChatModel) {
return AiServices.builder(StreamingAssistant.class)
.streamingChatModel(streamingChatModel)
.chatMemory(MessageWindowChatMemory.withMaxMessages(10))
.build();
}
4.3. 服务类
- 返回值和之前的低级API一致都是SseEmitter,而具体实现流式传输的关键,就是通过
streamingAssistant.chat拿到TokenStream实例后,用该实例的onPartialResponse方法完成每次token生成的响应
package cn.cjc.ai.service.impl;
import cn.cjc.ai.service.StreamingAssistant;
import dev.langchain4j.service.TokenStream;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@Service
@Slf4j
public class StreamingChatService {
@Autowired
private StreamingAssistant streamingAssistant;
/**
* 基于高级API的SSE流式聊天方法
* @param prompt 提示词
* @return SseEmitter实例,用于向客户端发送流式响应
*/
public SseEmitter highLevelStreamingChat(String prompt) {
// 创建SseEmitter,设置超时时间为30分钟
SseEmitter emitter = new SseEmitter(30 * 60 * 1000L);
try {
// 使用高级API的TokenStream
TokenStream tokenStream = streamingAssistant.chat(prompt);
// 注册回调函数
tokenStream.onPartialResponse(token -> {
if (token == null || token.trim().isEmpty()) {
log.warn("空的token,忽略");
return;
}
try {
// 发送token到客户端
emitter.send(token);
log.info("token: {}", token);
} catch (Exception e) {
log.error("发送token失败: {}", e.getMessage(), e);
emitter.completeWithError(e);
}
});
tokenStream.onError(throwable -> {
log.error("高级API流式聊天发生错误: {}", throwable.getMessage(), throwable);
emitter.completeWithError(throwable);
});
tokenStream.onCompleteResponse(completeResponse -> {
try {
// 发送完成标志
emitter.send(SseEmitter.event().name("complete"));
// 完成响应
emitter.complete();
log.info("高级API流式聊天完成,完整响应: {}", completeResponse);
} catch (Exception e) {
log.error("发送完成响应失败: {}", e.getMessage(), e);
emitter.completeWithError(e);
}
});
// 启动流处理
tokenStream.start();
} catch (Exception e) {
log.error("创建高级API流式响应失败: {}", e.getMessage(), e);
emitter.completeWithError(e);
}
// 设置超时处理
emitter.onTimeout(() -> {
log.error("高级API SSE连接超时");
emitter.completeWithError(new RuntimeException("连接超时"));
});
// 设置完成处理
emitter.onCompletion(() -> {
log.info("高级API SSE连接完成");
});
return emitter;
}
}
4.4. controller类
@RestController
@RequestMapping("/api/streamingChat")
@Slf4j
public class StreamingChatController {
@Autowired
private StreamingChatService streamingChatService;
/**
* 基于高级API的SSE流式聊天接口(用于网页实时显示)
* @param prompt 提示词
* @param userId 用户ID
* @return SseEmitter实例,用于向客户端发送流式响应
*/
@GetMapping("/high-level-sse-streaming-chat")
public SseEmitter highLevelStreamingChat(@RequestParam(name = "prompt") String prompt,
@RequestParam(name = "userId") int userId) {
log.info("收到来自用户[{}]的高级API请求, 提示词: {}", userId, prompt);
try {
// 调用QwenService的高级API流式聊天方法
return streamingChatService.highLevelStreamingChat(prompt);
} catch (Exception e) {
// 捕获异常并返回错误的SseEmitter
SseEmitter emitter = new SseEmitter();
emitter.completeWithError(e);
return emitter;
}
}
}
更多推荐



所有评论(0)