本篇功能由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>
  1. 点击按钮后,通过创建EventSource的方式和后端建立长连接,低级API和高级API的分别有不同的path
  2. 建立连接后,收到后台数据时,和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;
        }
    }
}
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐