☰
ChatGPT-SDK 流式应答会话设计实现:以 MyBatis 会话模型封装 OpenAI 事件驱动应答
2026/9/25 5:59:19 网站建设 项目流程
  • 文档
  • 教程
  • 后端

【免费下载链接】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创建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 的会话模型中:

  1. 架构层面:以OpenAiSessionFactory → OpenAiSession的会话模型承接所有 OpenAI 接口,架构即骨架;
  2. 设计层面:流式应答不直接暴露 HTTP 细节,而是以EventSourceListener事件监听方式对外提供增量数据,设计即方法;
  3. 代码层面:DefaultOpenAiSession#chatCompletions完成参数校验、请求构建与事件源创建,拦截器补齐鉴权头,代码即材料。

从源码结构看,这套会话模型具备极强的可扩展性:后续的文生图、向量计算、语音转文字等接口,以及多渠道动态apiHost/apiKey支持,都是在这套骨架上的自然延伸(见第3节:完善实现各类常用接口与第4节:支持多渠道对话)。而"架构、设计、代码"三部分的正确定位,才是避免代码库随着功能涌入而腐化的根本——这也是把 SDK 写好、把任何 HTTP 服务封装好的通用方法论。完整的项目体系说明可参考《OpenAi 大模型应用服务体系构建》与OpenAi 项目复盘总结。

  • 文档
  • 教程
  • 后端

【免费下载链接】CodeGuide

:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!

项目地址:https://gitcode.com/gh_mirrors/code/CodeGuide
点击查看免费下载

相关推荐

上一篇:WeUI.js核心组件详解:从alert到uploader的完整使用教程
下一篇:时间尺度自适应原理实战:Granite-TimeSeries-FlowState-R1-NPU的scale_factor如何适配任意采样率

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询