ChatGPT-SDK 流式应答会话设计实现:以 MyBatis 会话模型封装 OpenAI 事件驱动应答
文档教程后端【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载导读本文聚焦于《ChatGPT 微服务应用体系构建》中 chatgpt-sdk 组件的关键一环——流式应答会话的设计与实现。在统一IOpenAiApi接口、统一OpenAiSession会话的标准之下通过事件监听EventSourceListener接收 OpenAI 的增量应答数据使调用方在统一会话工厂中获得会话服务后依据接口入参差异走不同的请求处理逻辑从而实现 ChatGPT 打字机式的渐进展示效果。读完本文你将掌握 OpenAI-SDK 中会话工厂、会话接口、事件源工厂的三角协作模型以及流式应答从请求构建、拦截鉴权到事件回调的完整落地链路。一、本章诉求为什么要把流式应答收敛到统一会话ChatGPT 的官方 HTTP 接口众多且包含各类配置项如果每个使用方都各自对接就会出现口口相传 文档提示式的低效协作。本章的核心诉求是以IOpenAiApi统一接口、OpenAiSession统一会话这 2 个标准封装流式应答操作流式应答以事件方式接收应答消息。这种设计带来的直接收益是在统一的会话工厂中获得会话接口服务以后根据接口入参的不同做不同的请求处理。对于使用方来说接口即文档——代码标准本身成为比文档更可靠的说明减少了沟通成本与文档维护成本。这一思路与 MyBatis 的会话模型一脉相承。MyBatis 中SqlSessionFactory创建SqlSessionSqlSession统一承载对 Mapper 的操作类比到 OpenAI SDK就是OpenAiSessionFactory创建OpenAiSessionOpenAiSession统一承载对 OpenAI 各类接口的调用文生文、文生图、向量计算、语音识别等。会话模型作为框架的根基一旦搭建好后续扩展各项功能便有迹可循而不是等新功能涌入后再把代码填充得杂乱无章。二、流程设计参照 MyBatis 会话模型丰富 OpenAiSession本章的整个流程为丰富OpenAiSession会话服务接口增加流式回答的事件监听处理此过程的实现以 MyBatis 的会话模型为参照。在实现理念上一个需求的落地被拆分为三个部分架构是骨架定义系统边界与模块职责设计是方法选择合适的设计模式组织代码结构代码是填材料完成具体逻辑实现。如果没有设计方法设计模式的运用就相当于把代码的材料直接扔进架构里久而久之代码会越来越混乱。所以本章的重点不只是功能实现还包括如何在会话这个流程下把流式的事件应答处理巧妙地封装进统一的会话接口内——这是架构设计能力的一部分也正是第 1 节建立的工程结构所要承接的内容工程搭建与 okhttp3/Retrofit2 封装详见第1节ChatGPT-SDK组件工程简单功能实现。整个 SDK 的模块划分与使用链路如下Configuration配置apiHost、apiKey、超时、拦截器、IOpenAiApi ↓ 注入 OpenAiSessionFactory会话工厂 ↓ openSession() OpenAiSession统一会话chatCompletions 等各类接口 ↓ 内部 EventSource.Factory事件源工厂okhttp3 提供 ↓ EventSourceListener事件监听onEvent/onFailure/onClosed 接收流式数据三、会话模型的三角协作工厂创建会话、会话承载接口1. 会话工厂屏蔽使用细节会话工厂是整个 SDK 的使用入口。使用方只需要创建工厂对象即可开启会话对接使用。从仓库中的同类实现chatglm-sdk 同款会话模型可以看到DefaultOpenAiSessionFactory#openSession的标准组装过程Override public OpenAiSession openSession() { // 1. 日志配置 HttpLoggingInterceptor httpLoggingInterceptor new HttpLoggingInterceptor(); httpLoggingInterceptor.setLevel(configuration.getLevel()); // 2. 开启 Http 客户端 OkHttpClient okHttpClient new OkHttpClient .Builder() .addInterceptor(httpLoggingInterceptor) .addInterceptor(new OpenAiHTTPInterceptor(configuration)) .connectTimeout(configuration.getConnectTimeout(), TimeUnit.SECONDS) .writeTimeout(configuration.getWriteTimeout(), TimeUnit.SECONDS) .readTimeout(configuration.getReadTimeout(), TimeUnit.SECONDS) .build(); configuration.setOkHttpClient(okHttpClient); // 3. 创建 API 服务 IOpenAiApi openAiApi new Retrofit.Builder() .baseUrl(configuration.getApiHost()) .client(okHttpClient) .addCallAdapterFactory(RxJava2CallAdapterFactory.create()) .addConverterFactory(JacksonConverterFactory.create()) .build().create(IOpenAiApi.class); configuration.setOpenAiApi(openAiApi); return new DefaultOpenAiSession(configuration); }这段代码的核心价值在于把 OkHttpClient 的组装、Retrofit 代理对象的创建、拦截器的注入全部封装在工厂内使用方只需一行factory.openSession()即可获得会话。其中configuration.getApiHost()作为 Retrofit 的baseUrl决定了请求的根地址OpenAiHTTPInterceptor负责在每个 HTTP 请求上补充ApiKey、Token等必要参数超时配置connectTimeout/writeTimeout/readTimeout来自Configuration的可配置项。这正是工厂屏蔽使用细节、外部最少知道原则的体现——会话模型的使用方不需要关心 HTTP 客户端如何构建只需要关心会话提供什么能力。2. 统一接口IOpenAiApi 与 OpenAiSession 双标准SDK 内部存在两层标准IOpenAiApi基于 Retrofit2 注解描述的 OpenAI HTTP 接口层将 HTTP API 转化为 Java 接口通过注解声明请求路径、请求方式与参数序列化方式Retrofit 的JacksonConverterFactory负责 JSON 转换OpenAiSession面向使用方的统一会话层把IOpenAiApi的底层细节全部隐藏对外只暴露语义清晰的业务方法如chatCompletions。这种双层结构的好处是底层接口负责 HTTP 协议细节上层会话负责业务语义与参数组织两层之间通过Configuration持有IOpenAiApi实例完成桥接。后续无论对接文生图、向量计算还是语音转文字等新接口都只需在IOpenAiApi增加方法、在OpenAiSession暴露对应会话方法即可第3节完善实现各类常用接口正是沿着这条路径继续补全各类接口。四、流式应答的核心事件监听封装流式应答与普通同步调用的本质区别在于响应不是一次性返回而是以 SSEServer-Sent Events数据块的形式持续到达。因此会话接口必须把接收数据的能力暴露为事件监听器。1. 会话接口以 EventSourceListener 接收流式数据在仓库后续章节第4节支持多渠道对话中可以看到该设计落地的最终形态——cn.bugstack.chatgpt.session.OpenAiSession接口/** * 问答模型 GPT-3.5/4.0 流式反馈 * * param apiHostByUser 自定义host * param apiKeyByUser 自定义Key * param chatCompletionRequest 请求信息 * param eventSourceListener 实现监听通过监听的 onEvent 方法接收数据 * return 应答结果 */ EventSource chatCompletions(String apiHostByUser, String apiKeyByUser, ChatCompletionRequest chatCompletionRequest, EventSourceListener eventSourceListener) throws JsonProcessingException;可以看到流式会话方法的标准形态请求参数对象ChatCompletionRequest 事件监听器EventSourceListener返回EventSource对象——该对象可随时取消应答。EventSourceListener是 okhttp3 提供的 SSE 事件回调接口正是本章以事件实现方式接收应答消息的载体。2. 会话实现构建请求并交给事件源工厂cn.bugstack.chatgpt.session.defaults.DefaultOpenAiSession中的chatCompletions实现完整展示了流式请求的构建链路public EventSource chatCompletions(String apiHostByUser, String apiKeyByUser, ChatCompletionRequest chatCompletionRequest, EventSourceListener eventSourceListener) throws JsonProcessingException { // 核心参数校验不对用户的传参做更改只返回错误信息。 if (!chatCompletionRequest.isStream()) { throw new RuntimeException(illegal parameter stream is false!); } // 动态设置 Host、Key便于用户传递自己的信息 String apiHost Constants.NULL.equals(apiHostByUser) ? configuration.getApiHost() : apiHostByUser; String apiKey Constants.NULL.equals(apiKeyByUser) ? configuration.getApiKey() : apiKeyByUser; // 构建请求信息 Request request new Request.Builder() // url: https://api.openai.com/v1/chat/completions - 通过 IOpenAiApi 配置的 POST 接口用这样的方式从统一的地方获取配置信息 .url(apiHost.concat(IOpenAiApi.v1_chat_completions)) .addHeader(apiKey, apiKey) // 封装请求参数信息如果使用了 Fastjson 也可以替换 ObjectMapper 转换对象 .post(RequestBody.create(MediaType.parse(ContentType.JSON.getValue()), new ObjectMapper().writeValueAsString(chatCompletionRequest))) .build(); // 返回结果信息EventSource 对象可以取消应答 return factory.newEventSource(request, eventSourceListener); }这段实现有几个关键点值得展开参数校验先行stream必须为true才允许进入流式应答流程避免非流式参数误入事件监听链路属于防御式编程统一从IOpenAiApi获取接口路径常量IOpenAiApi.v1_chat_completions对应/v1/chat/completions接口路径信息只维护在一处避免散落各处导致失配请求头携带 apiKey通过addHeader(apiKey, apiKey)把 Key 挂在原始请求上供后续拦截器读取使用factory.newEventSource(...)是流式的关键一步这里的factory是 okhttp3 的EventSource.Factory由Configuration#createRequestFactory创建它负责建立 SSE 长连接并把数据块分发给EventSourceListener。3. 拦截器把 ApiKey 装配成鉴权头cn.bugstack.chatgpt.interceptor.OpenAiInterceptor负责在请求发出前完成最终鉴权头组装public Response intercept(Chain chain) throws IOException { // 1. 获取原始 Request Request original chain.request(); // 2. 读取 apiKey优先使用自己传递的 apiKey String apiKeyByUser original.header(apiKey); String apiKey Constants.NULL.equals(apiKeyByUser) ? apiKeyBySystem : apiKeyByUser; // 3. 构建 Request Request request original.newBuilder() .url(original.url()) .header(Header.AUTHORIZATION.getValue(), Bearer apiKey) .header(Header.CONTENT_TYPE.getValue(), ContentType.JSON.getValue()) .method(original.method(), original.body()) .build(); // 4. 返回执行结果 return chain.proceed(request); }拦截器的逻辑清晰从原始请求头读取apiKey会话实现中写入的那个头优先使用调用方传入的 Key否则回落到系统初始化时配置的默认 Key最终拼装为Authorization: Bearer apiKey标准鉴权头。这样流式会话既可以用系统默认账号体验也支持每个用户绑定自己的 Key——这一能力在第 4 节中被扩展为多渠道对话的核心机制。4. 事件监听打字机效果的源头使用方通过实现EventSourceListener接收流式数据块EventSource eventSource openAiSession.chatCompletions(chatCompletion, new EventSourceListener() { Override public void onEvent(EventSource eventSource, String id, String type, String data) { // 每个数据块chunk携带增量内容前端逐块填充即可形成打字机效果 log.info(测试结果 id:{} type:{} data:{}, id, type, data); } Override public void onFailure(EventSource eventSource, Throwable t, Response response) { log.error(失败 code:{} message:{}, response.code(), response.message()); } });流式应答消息中每次onEvent回调的data是一个 JSON 数据块形如{choices:[{delta:{content:1}}]}其中delta.content就是本次的增量文本当收到finish_reason: stop后服务端最终发送[DONE]标记整个流结束。前端拿到这些增量块逐段填充到消息区就呈现出 ChatGPT 式的渐进打字效果——这也是整个《ChatGPT 微服务应用体系构建》项目中 API 层异步响应接口ResponseBodyEmitter与 WEB 层 ReadableStream 填充的数据来源可对照 第4节工程重构和流式异步响应接口实现 与 第8节流式接口对接。五、单元测试与功能验证1. 会话初始化Before public void test_OpenAiSessionFactory() { // 1. 配置文件 [联系小傅哥获取key] Configuration configuration new Configuration(); configuration.setApiHost(https://api.openai.com/); configuration.setApiKey(sk-xxxxxxxxxxxxxxxxxxxxxxxx); // 2. 会话工厂 OpenAiSessionFactory factory new DefaultOpenAiSessionFactory(configuration); // 3. 开启会话 this.openAiSession factory.openSession(); }Configuration是会话配置的载体至少需要设置两项apiHostOpenAI 接口根地址例如https://api.openai.com/apiKey官网申请的密钥。2. 流式问答测试public void test_chat_completions_stream() throws JsonProcessingException, InterruptedException { // 1. 创建参数 ChatCompletionRequest chatCompletion ChatCompletionRequest .builder() .stream(true) .messages(Collections.singletonList(Message.builder() .role(Constants.Role.USER) .content(11) .build())) .model(ChatCompletionRequest.Model.GPT_3_5_TURBO.getCode()) .maxTokens(1024) .build(); // 2. 发起请求 EventSource eventSource openAiSession.chatCompletions(chatCompletion, new EventSourceListener() { Override public void onEvent(EventSource eventSource, String id, String type, String data) { log.info(测试结果 id:{} type:{} data:{}, id, type, data); } Override public void onFailure(EventSource eventSource, Throwable t, Response response) { log.error(失败 code:{} message:{}, response.code(), response.message()); } }); // 等待 new CountDownLatch(1).await(); }测试要点stream(true)是触发流式应答的开关与服务端 SSE 模式对应model指定模型例如GPT_3_5_TURBOgpt-3.5-turbomessages携带对话上下文roleuser、content11为本次提问内容new CountDownLatch(1).await()阻塞主线程等待事件回调驱动流程避免测试方法提前退出。3. 测试结果流式应答的日志输出如下数据块逐条到达测试结果{id:chatcmpl-85SM...,object:chat.completion.chunk,created:1696311335, model:gpt-3.5-turbo-0613,choices:[{index:0,delta:{role:assistant,content:},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,object:chat.completion.chunk, choices:[{index:0,delta:{content:1},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content:},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content:1},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content: equals},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content: },finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content:2},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{content:.},finish_reason:null}]} 测试结果{id:chatcmpl-85SM...,choices:[{index:0,delta:{},finish_reason:stop}]} 测试结果[DONE]可以看到delta.content依次输出1、、1、equals、、2、.——正是11 equals 2.的逐字增量最后一个数据块携带finish_reason:stop随后[DONE]标记流结束。这种逐块到达的数据形态就是前端打字机效果的数据基础。六、小结本章通过统一接口 统一会话 事件监听三层设计把 OpenAI 流式应答完整收敛到 SDK 的会话模型中架构层面以OpenAiSessionFactory → OpenAiSession的会话模型承接所有 OpenAI 接口架构即骨架设计层面流式应答不直接暴露 HTTP 细节而是以EventSourceListener事件监听方式对外提供增量数据设计即方法代码层面DefaultOpenAiSession#chatCompletions完成参数校验、请求构建与事件源创建拦截器补齐鉴权头代码即材料。从源码结构看这套会话模型具备极强的可扩展性后续的文生图、向量计算、语音转文字等接口以及多渠道动态apiHost/apiKey支持都是在这套骨架上的自然延伸见第3节完善实现各类常用接口与第4节支持多渠道对话。而架构、设计、代码三部分的正确定位才是避免代码库随着功能涌入而腐化的根本——这也是把 SDK 写好、把任何 HTTP 服务封装好的通用方法论。完整的项目体系说明可参考《OpenAi 大模型应用服务体系构建》与OpenAi 项目复盘总结。赞分享文档教程后端【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载相关推荐DeepSeek Harness 事件溯源会话架构以仅追加日志作为唯一真源的回放式会话设计DeepSeek Harness 事件溯源会话架构以仅追加日志作为唯一真源的回放式会话设计 DeepSeek Harnessdsh以“Everything人工智能AI AgentAgent 框架DeepSeekTypeChat事件驱动架构响应式对话系统设计TypeChat事件驱动架构响应式对话系统设计 1. 痛点直击当LLM输出失控时 你是否曾遭遇这些困境对话系统将少放辣误解为不要辣用户输入半份大模型AI 应用后端oh-my-opencode-slim 会话考古/reflect --sessions 跨会话反思模式的设计与实现解析oh my opencode slim 会话考古/reflect sessions 跨会话反思模式的设计与实现解析 导读 本篇文章以 oh my openco人工智能AI AgentAgent 编排AI 技能上一篇WeUI.js核心组件详解从alert到uploader的完整使用教程下一篇时间尺度自适应原理实战Granite-TimeSeries-FlowState-R1-NPU的scale_factor如何适配任意采样率创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考