- 文档
- 教程
- 后端
【免费下载链接】CodeGuide
:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
导读
本文聚焦于《ChatGPT 微服务应用体系构建》中 chatgpt-sdk 组件的关键一环——流式应答会话的设计与实现。在统一IOpenAiApi接口、统一OpenAiSession会话的标准之下,通过事件监听(EventSourceListener)接收 OpenAI 的增量应答数据,使调用方在统一会话工厂中获得会话服务后,依据接口入参差异走不同的请求处理逻辑,从而实现 ChatGPT 打字机式的渐进展示效果。读完本文,你将掌握 OpenAI-SDK 中会话工厂、会话接口、事件源工厂的三角协作模型,以及流式应答从请求构建、拦截鉴权到事件回调的完整落地链路。
一、本章诉求:为什么要把流式应答收敛到统一会话
ChatGPT 的官方 HTTP 接口众多,且包含各类配置项,如果每个使用方都各自对接,就会出现"口口相传 + 文档提示"式的低效协作。本章的核心诉求是:以IOpenAiApi统一接口、OpenAiSession统一会话这 2 个标准,封装流式应答操作,流式应答以事件方式接收应答消息。
这种设计带来的直接收益是:在统一的会话工厂中获得会话接口服务以后,根据接口入参的不同做不同的请求处理。对于使用方来说,接口即文档——代码标准本身成为比文档更可靠的说明,减少了沟通成本与文档维护成本。
这一思路与 MyBatis 的会话模型一脉相承。MyBatis 中SqlSessionFactory创建SqlSession,SqlSession统一承载对 Mapper 的操作;类比到 OpenAI SDK,就是OpenAiSessionFactory创建OpenAiSession,OpenAiSession统一承载对 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节:完善实现各类常用接口正是沿着这条路径继续补全各类接口)。
四、流式应答的核心:事件监听封装
流式应答与普通同步调用的本质区别在于:响应不是一次性返回,而是以 SSE(Server-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是会话配置的载体,至少需要设置两项:
apiHost:OpenAI 接口根地址,例如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("1+1") .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_TURBO(gpt-3.5-turbo);messages携带对话上下文,role=user、content="1+1"为本次提问内容;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、.——正是"1+1 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核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
相关推荐
DeepSeek Harness 事件溯源会话架构:以仅追加日志作为唯一真源的回放式会话设计
DeepSeek Harness 事件溯源会话架构:以仅追加日志作为唯一真源的回放式会话设计 DeepSeek Harness(dsh)以“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 技能
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考