1、支持流式调用

在 AI 基础对话 中将 call() 方法,改为 stream() 方法即可。

/**
     * AI 基础对话(支持多轮对话记忆)
     * @param message
     * @param chatId
     * @return
     */
    public String doChat(String message,String chatId){
        ChatResponse response = chatClient
                .prompt()
                .user(message)
                .advisors(spec -> spec.param(CHAT_MEMORY_CONVERSATION_ID_KEY,chatId))
                .call()
                .chatResponse();
        String content = response.getResult().getOutput().getText();
        log.info("content: {}", content);
        log.info("chatId: {}", chatId);
        return content;
    }

改为 stream() 方法

/**
     * AI 基础对话(支持多轮对话记忆,SSE流式传输)
     * @param message
     * @param chatId
     * @return
     */
    public Flux<String> doChatByStream(String message,String chatId){
        return  chatClient
                .prompt()
                .user(message)
                .advisors(spec -> spec.param(CHAT_MEMORY_CONVERSATION_ID_KEY, chatId))
                .stream()
                .content();
    }

注意:不要直接‏使用 ChatResponse 作؜为返回类型,因为这会导致返回内容膨​胀,影响传输效率。所以上述代码中我‌们使用 content 方法,只返‏回 AI 输出的文本信息。

2、开发同步接口

@RestController
@RequestMapping("/ai")
public class AiController {

    @Resource
    private LoveApp loveApp;

   

    @GetMapping("/love_app/chat/sync")
    public String doChatWithLoveAppSync(String message, String chatId) {
        return loveApp.doChat(message, chatId);
    }
}

3、开发 SSE 流式接口

编写基‏于 SSE 的流式؜输出接口,有两种常见方式:

1、返回‏ Flux 响应式؜对象,并且添加 S​SE 对应的 Me‌diaType:

/**
     * 返回Flux 响应式؜对象,并且添加 SSE 对应的 MediaType:
     * @param message
     * @param chatId
     * @return
     */
    @GetMapping(value = "/love_app/chat/sse", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public Flux<String> doChatWithLoveAppSSE(String message, String chatId) {
        return loveApp.doChatByStream(message, chatId);
    }

2、使用 ‏SSEEmiter,؜通过 send 方法​持续向 SseEmi‌tter 发送消息

/**
     * 使用 SSEEmiter,؜通过 send 方法持续向 SseEmitter 发送消息
     * @param message
     * @param chatId
     * @return
     */
    @GetMapping("/love_app/chat/sse/emitter")
    public SseEmitter doChatWithLoveAppSSEEmitter(String message, String chatId) {
        // 3 分钟超时
        SseEmitter emitter = new SseEmitter(180000L);
        //获取 Flux 数据流并直接订阅
        loveApp.doChatByStream(message, chatId)
                .subscribe(
                        //处理每条消息
                        chunk->{
                            try {
                                emitter.send(chunk);
                            } catch (IOException e) {
                                emitter.completeWithError(e);
                            }
                        },
                        //处理错误
                        emitter::completeWithError,
                        //处理完成
                        emitter::complete
                );
        //返回 emitter
        return emitter;

    }

为啥要写 emitter.completeWithError(e);

如果不写可能导致两个问题:

  1. 服务端这边可能资源没释放,导致内存泄漏、线程泄漏(因为这个异步请求一直挂着)。
  2. 客户端不知道你已经崩了,还在傻傻等。

所以一定要写emitter.completeWithError(e); ,这样可以让服务端优雅退出 SSE 推送,并把出错信息同步给客户端,同时释放资源。

4、总结:两种写法的区别

特点SseEmitter 手动流式推送Flux + Spring WebFlux 自动流式
返回值SseEmitter(自己控制推送)Flux<String>(Spring WebFlux 响应式原生流)
谁负责流控制手写 emitter.send(),完全自己掌控什么时候发、怎么发Spring WebFlux 框架接管了流推送,用户只负责返回 Flux
MediaType没写 produces,需要手工组装消息格式(SSE 格式)自动设置 MediaType.TEXT_EVENT_STREAM_VALUE,浏览器原生支持
线程模型Spring MVC 的 Servlet 模型,可能卡主线程 / 阻塞WebFlux 的响应式非阻塞线程模型,天然支持高并发、低资源消耗
适用场景需要精细控制:如定制 SSE 格式、心跳、重连 ID、分片组装等纯流式输出、简单直接、无需额外控制
复杂度写得多,手工处理异常、完成、超时框架帮你搞定,写得少,易维护

5、测试

debug 模式测试,可以看出 AI 的输出是像 “打字机” 一样,一段一段输出,而不是之前所有消息都拼接好后才输出。
在这里插入图片描述

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐