ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

从零实现 C++ AI 大模型接入 SDK(五):DeepSeek 流式响应 sendMessageStream() 完整实现

从零实现 C++ AI 大模型接入 SDK(五):DeepSeek 流式响应 sendMessageStream() 完整实现 目录前言一、从全量响应到流式响应1.1 为什么这里使用 SSE二、sendMessageStream()从请求到流式接收2.1 先告诉 DeepSeek我要流式返回2.2 为什么不再直接使用 client.Post()2.3 response_handler先判断 HTTP 层是否成功2.4 content_receiver真正的流式入口三、网络 chunk 为什么不能直接解析3.1 chunk 不等于一条完整 SSE event3.2 buffer把任意 chunk 重新拼回完整 SSE event3.3 为什么 buffer 要一直保留“半截事件”3.4 同时兼容 \n\n 和 \r\n\r\n四、从 SSE event 到模型增量文本4.1 先从完整 event 中提取 data payload4.2 payload 先判断 [DONE]再解析 JSON4.3 为什么流式响应读取 choices[0].delta.content五、回调、完整结果与结束处理5.1 callback 的 bool 参数表示什么5.2 为什么还要有 fullResponse5.3 notifyDone() 为什么单独封装5.4 请求结束后还要区分三类问题六、把 sendMessageStream() 的完整链路串起来七、测试流式响应7.1 用项目测试验证 callback 和最终返回值7.2 最小 Demo直接看见“边生成边输出”写在最后前言系列从零实现 C AI 大模型接入 SDK第五篇AI-CHAT-SDKhttps://gitee.com/kuang-zhenting/my_ai_cpp_project上一篇我们已经把DeepSeekProvider::sendMessage()的全量调用链跑通了把Message转成 JSON通过 HTTPS 请求 DeepSeek最后从choices[0].message.content里拿到完整回复。但全量响应有一个很明显的问题。如果模型要生成一段比较长的内容用户看到的过程会是发送问题 ↓ 等待 ↓ 模型全部生成完成 ↓ 一次性显示整段回答真正的聊天应用更希望发送问题 ↓ 生成一点 ↓ 马上显示一点 ↓ 继续生成 ↓ 继续显示所以这一篇继续沿着DeepSeekProvider往下实现上一篇已经预留好的sendMessageStream(...)这次最麻烦的地方并不是给请求体多加一句stream: true而是响应体不再是一份完整 JSON。我们会不断收到原始网络数据块再从这些数据块里恢复完整 SSE 事件提取data:识别[DONE]解析choices[0].delta.content最后通过回调把新增文本一段一段交给上层。这一篇真正要解决的就是下面这条链一、从全量响应到流式响应先把上一篇和这一篇放在一起看。DeepSeek 的请求路径仍然是POST /v1/chat/completions请求里仍然有model messages temperature max_tokensProvider 的初始化方式、API Key、Base URL 也没有因为流式响应而重新设计一套。真正变化的是响应体怎样被客户端消费。上一篇的sendMessage()更像发送请求 ↓ 等待完整响应体 ↓ 解析一份完整 JSON ↓ 得到 std::string这一篇的sendMessageStream()则变成发送请求 ↓ 收到一部分响应体 ↓ 立即处理 ↓ 再收到一部分 ↓ 继续处理 ↓ 直到本轮结束这里有一个很容易误解的点流式响应描述的是“响应体分段到达、分段处理”不等于当前函数自动变成异步函数。当前源码最后仍然会调用auto result client.send(req);当前线程会等待这次请求结束只不过等待期间cpp-httplib会不断调用我们的content_receiver把已经到达的响应体数据交出来而不是等整个 body 全部收完再统一返回。1.1 为什么这里使用 SSESSE 的全称是Server-Sent Events。 是一种基于 HTTP 的“服务器向客户端单向推送”技术。SSE 让服务器可以不断给浏览器发消息但浏览器不能通过这条连接给服务器发消息。对于一轮普通大模型对话来说通信方向通常是客户端提交一次问题 ↓ 模型开始生成 ↓ 服务器持续向客户端返回增量文本 ↓ 本轮结束也就是说问题发出去以后接下来主要是服务器不断向客户端推送这一轮的新内容。这和当前需要的模式非常匹配。我们并不需要为了“实时”就把 Provider 改成 WebSocket。简单区分一下对比项SSE / HTTP 流式响应WebSocket当前需要的通信方式一次请求后持续接收双方都可以随时主动发送是否继续沿用 HTTP 请求语义是建连后进入 WebSocket 通信适合当前模型回复很合适能实现但不是当前 Provider 的方案当前项目的 DeepSeekProvider使用未使用对这一篇来说SSE 最重要的不是背协议概念而是看懂它返回的数据长什么样。一个简化后的事件可能是data: {choices:[{delta:{content:你}}]}下一条事件又可能是data: {choices:[{delta:{content:好}}]}最后还会出现结束标记data: [DONE]所以和上一篇相比JSON 路径也发生了变化。全量响应读取choices[0].message.content流式响应读取choices[0].delta.content这里的delta可以理解成“这一小次新增加的内容”。二、sendMessageStream()从请求到流式接收真正进入实现后可以把这一部分看成两步先把请求切换到流式模式再让cpp-httplib在响应体还没有接收完成时持续把数据交给我们。2.1 先告诉 DeepSeek我要流式返回sendMessageStream()的前半部分和上一篇sendMessage()很像。函数签名来自ILLMProvidervirtual std::string sendMessageStream( const std::vectorMessage messages, const std::mapstd::string, std::string request_param, std::functionvoid(const std::string , bool) callback) 0;DeepSeekProvider对它进行实现std::string DeepSeekProvider::sendMessageStream( const std::vectorMessage messages, const std::mapstd::string, std::string request_param, std::functionvoid(const std::string , bool) callback)前面的工作仍然包括检查 Provider 是否可用 ↓ 把 Message 转成 messages 数组 ↓ 读取 temperature / max_tokens ↓ 构造 JSON 请求体真正新增的是request_body[stream] true;请求体大致会变成{ model: deepseek-chat, messages: [ { role: user, content: 请解释一下什么是 SSE。 } ], temperature: 0.7, max_tokens: 2048, stream: true }请求头也增加了{Accept, text/event-stream}所以当前源码的写法是httplib::Headers headers { {Authorization, Bearer _api_key}, {Content-Type, application/json}, {Accept, text/event-stream} };到这里我们只是把请求切换到了流式模式。真正复杂的事情发生在“怎样接收响应体”这一侧。2.2 为什么不再直接使用 client.Post()上一篇发送全量请求时可以直接写auto response client.Post( /v1/chat/completions, headers, json_string, application/json);因为那时候只需要等请求结束再统一处理response-status response-body流式响应不一样。我们需要在 body还没有收完的时候就介入所以当前源码改成手动构造httplib::Request req; req.method POST; req.path /v1/chat/completions; req.headers headers; req.body json_string;然后给这个Request安装两个非常重要的回调response_handler content_receiver它们负责的事情完全不同。2.3 response_handler先判断 HTTP 层是否成功在项目源码中req.response_handler [](const httplib::Response response) { statusCode response.status; if (statusCode ! 200) { gotError true; errorMsg HTTP error: std::to_string(statusCode); ERR({}, errorMsg); return true; } return true; };这个回调关注的是响应头里的 HTTP 状态码例如200 请求正常 401 API Key 有问题 404 接口路径不对 429 请求过于频繁这里有一个很值得注意的处理即使状态码不是200代码仍然return true;不是说错误可以忽略而是为了继续读取服务器返回的错误正文。如果这里直接终止我们可能只能得到HTTP error: 401却拿不到服务器 body 里更具体的错误描述。所以当前实现采用响应头发现非 200 ↓ gotError true ↓ 继续接收 body ↓ content_receiver 收集错误正文 ↓ 请求结束后统一输出2.4 content_receiver真正的流式入口当前源码的核心回调是req.content_receiver [](const char *data, size_t len, uint64_t /*offset*/, uint64_t /*totalLength*/) { // ... return true; };只要新的响应体数据到达这个回调就可能再次被调用。最重要的两个参数是data 当前数据块的起始地址 len 当前数据块的字节数而offset、totalLength在当前实现中暂时没有使用。如果前面的 HTTP 状态已经表明请求失败当前这一次收到的就不是正常 SSE 数据而是错误正文if (gotError) { errorMsg.append(data, len); return true; }正常情况下才进入真正的 SSE 处理。三、网络 chunk 为什么不能直接解析这是整个流式实现里最值得花时间理解的地方。3.1 chunk 不等于一条完整 SSE eventcontent_receiver每次拿到一块data len这块数据通常可以叫做一个network chunk。但网络层怎么切 chunk并不需要和 SSE 的事件边界对齐。例如逻辑上的 SSE 数据本来是data: event-1\n\n data: event-2\n\n data: event-3\n\n网络完全可能这样交给我们第一次data: eve 第二次nt-1\n\ndata: event- 第三次2\n\ndata: event-3\n\n也可能一次回调就同时带来两三条完整事件。所以这两句话一定要分清chunk 是网络这一次交给我们的数据块。event 才是 SSE 业务上应该完整解析的一条事件。如果我们拿到一个 chunk 就立刻当 JSON 解析会出现一种非常难查的问题有时候能解析有时候突然 JSON 报错换个网络环境表现又不一样。因为真正的问题不是 JSON 内容错了而是这次拿到的 JSON 只有半截。这就是buffer必须存在的原因。3.2 buffer把任意 chunk 重新拼回完整 SSE event当前源码在发送请求前准备了一个缓冲区std::string buffer;每次content_receiver被调用先把当前数据块追加进去buffer.append(data, len);这里特意没有写buffer data;因为data是一段由len明确指定长度的内存并不能假设它一定以\0结尾。所以最稳妥的写法就是buffer.append(data, len);明确告诉std::string只追加这len个字节。3.3 为什么 buffer 要一直保留“半截事件”假设当前 buffer 是data: {event1}\n\ndata: {event2 的一半那么我们真正能处理的只有data: {event1}event2还没有收完整必须继续留在 buffer 中。等下一次 chunk 到达以后event2 的另一半\n\n再拼成完整事件。所以当前代码使用一个循环while (true) { // 尝试从 buffer 中找完整事件 }只要 buffer 里还能找到完整事件就继续处理一旦只剩半截就退出循环等待下一次网络数据。3.4 同时兼容 \n\n 和 \r\n\r\nSSE 事件之间通过空行分隔。当前实现同时查找const size_t lfpos buffer.find(\n\n); const size_t crlfpos buffer.find(\r\n\r\n);然后取先出现的那一个作为本次事件边界。找到以后std::string event buffer.substr(0, eventEnd); buffer.erase(0, eventEnd separatorLength);这两句配合起来做了两件事substr() ↓ 取出一条完整 event buffer.erase() ↓ 把已经处理的 event 分隔符从 buffer 删除剩下的内容继续留给下一轮处理。到这里我们才真正从“网络数据块”恢复成“可以按 SSE 语义处理的完整事件”。四、从 SSE event 到模型增量文本前面解决的是“怎样拿到一条完整事件”。接下来才进入 SSE 和 DeepSeek 返回格式本身先从 event 中提取data:再判断[DONE]或解析 JSON最后取出真正新增的模型文本。4.1 先从完整 event 中提取 data payload拿到一条完整 SSE event 后还不能直接送进 JSONCpp。SSE 事件本身可能包含注释行 空行 data: ...当前源码会先逐行处理再把真正的data:内容提取出来。先准备std::string payload; std::istringstream eventStream(event); std::string line;然后逐行读取while (std::getline(eventStream, line)) { // ... }先处理\r\n留下的\r如果事件使用\r\n换行std::getline()按\n切开以后行尾还可能剩一个\r。当前源码先做if (!line.empty() line.back() \r) { line.pop_back(); }这里先判断line.empty()是为了避免对空字符串直接调用back()。空行和注释行直接跳过SSE 中以:开头的行可以是注释例如心跳信息。当前实现if (line.empty() || line[0] :) { continue; }也就是说这些内容不会进入 JSON 解析。rfind(..., 0)是在判断“是否以data:开头”接下来有一句很容易第一次看到时不理解if (line.rfind(data:, 0) 0)这里不是在“从右边随便找一个data:”。第二个参数传0后它表达的是data: 是否正好从下标 0 开始出现也就是在判断这一行是不是以 data: 开头然后去掉前缀std::string value line.substr(5);如果data:后面还有一个空格再去掉if (!value.empty() value.front() ) { value.erase(value.begin()); }最后追加到payload value;如果同一个 event 中出现多行data:当前实现会按照读取顺序把这些 value 依次追加到payload。处理完成以后后面的代码只关心payload不再关心 SSE 行前缀。4.2 payload 先判断 [DONE]再解析 JSON拿到payload后会出现两类情况。第一类是[DONE]第二类是正常 JSON{ choices: [ { delta: { content: 你好 } } ] }所以当前源码先判断if (payload [DONE]) { streamFinish true; notifyDone(); continue; }[DONE]不是 JSON。如果不先判断就把它交给 JSONCpp自然会解析失败。对于普通 payload才进入Json::Value chunk; Json::CharReaderBuilder readerBuilder; std::string errors; std::istringstream jsonStream(payload); if (!Json::parseFromStream(readerBuilder, jsonStream, chunk, errors)) { WARN(DeepSeek SSE JSON parse error: {}, errors); continue; }一条事件解析失败时当前实现打印警告并继续下一条而不是直接让整个程序崩掉。4.3 为什么流式响应读取 choices[0].delta.contentJSON 解析成功以后接下来终于到了真正的模型文本。当前源码检查if (chunk.isMember(choices) chunk[choices].isArray() !chunk[choices].empty() chunk[choices][0].isMember(delta) chunk[choices][0][delta].isMember(content)) { const std::string content chunk[choices][0][delta][content].asString(); fullResponse content; if (!content.empty() callback) { callback(content, false); } }这里有两个重点。为什么总是choices[0]对于当前这条聊天调用链我们取choices数组中的第一个候选结果chunk[choices][0]所以后面的路径自然是choices[0] ↓ delta ↓ content为什么必须一层一层检查不能直接写chunk[choices][0][delta][content].asString();然后假设每一条流式事件都有content。因为流式过程中第一块可能只有角色信息中间块通常携带content最后一块可能只有结束原因。所以“收到了一条合法的流式 JSON”不等于“这一条一定包含新文本”。当前源码逐层判断字段就是为了避免对不存在的节点直接访问。五、回调、完整结果与结束处理模型文本已经能一段段解析出来以后还要解决另一个问题这些增量内容怎样交给上层以及正常结束、HTTP 错误、网络错误和异常断流怎样统一收口。5.1 callback 的 bool 参数表示什么sendMessageStream()的回调类型是std::functionvoid(const std::string , bool) callback两个参数分别承担第一个参数本次新增文本 第二个参数这一轮流是否结束所以普通增量到来时callback(content, false);上层可以立刻把content显示出来。而流真正结束时callback(, true);这时候重点已经不是又来了一段文字而是告诉上层这一轮回复结束了。5.2 为什么还要有 fullResponse既然文本已经通过 callback 一段段交出去了为什么函数还要std::string fullResponse;并不断fullResponse content;因为这两个出口的职责不同callback ↓ 负责实时交付 return fullResponse ↓ 负责函数结束后拿到完整回答所以sendMessageStream()最终仍然返回return fullResponse;这也给我们提供了一个非常直接的测试办法把 callback 中收到的所有 chunk 拼起来最终应该和fullResponse相同。5.3 notifyDone() 为什么单独封装当前源码没有在每个结束分支里都直接写callback(, true);而是先定义auto notifyDone []() { if (!doneCallbackSent callback) { callback(, true); doneCallbackSent true; } };原因是结束不只可能来自正常的[DONE]。网络失败、HTTP 错误、连接提前结束都需要确保上层不再继续等待。如果每个分支各发一次结束回调又很容易重复通知。所以用doneCallbackSent统一做一次防重复保护。5.4 请求结束后还要区分三类问题流式接收结束并不代表一定是正常结束。当前源码在auto result client.send(req);返回以后还会继续检查。第一类transport / 网络层错误例如DNS 解析失败 连接超时 TLS 握手失败 读取超时这类错误可能连 HTTP 响应头都没有收到。当前实现if (!result) { const auto error result.error(); ERR(Network error: {}, httplib::to_string(error)); notifyDone(); return ; }这里使用httplib::to_string(error)把httplib::Error转成可打印字符串。第二类HTTP 错误如果已经拿到 HTTP 响应但状态不是200前面的response_handler会设置gotError true;同时content_receiver会继续把错误正文追加到errorMsg。请求结束后统一处理if (gotError) { ERR(DeepSeek stream request failed: {}, errorMsg); notifyDone(); return ; }第三类连接结束了但没有收到[DONE]还有一种情况更隐蔽HTTP 建连成功 ↓ 也收到过一些模型内容 ↓ 连接却中途结束 ↓ 从来没收到 data: [DONE]这种情况下不能把连接断了当成业务正常结束。当前源码专门检查if (!streamFinish) { WARN(Stream ended without [DONE] marker); notifyDone(); }所以这里有两个不同的结束网络连接结束 ≠ 模型流正常结束只有看到[DONE]streamFinish才会被设置成true。六、把 sendMessageStream() 的完整链路串起来前面已经把几个容易混淆的层次拆开了现在把整条调用链重新合起来看整个sendMessageStream()可以理解为真正的难点其实集中在中间这一小段chunk → buffer → event → payload → delta.content只要把这五层分清楚后面的 ChatGPT、Gemini 等 Provider 即使协议细节不同理解流式接收时也不会再把“网络数据块”和“模型文本片段”混成同一个东西。七、测试流式响应实现完成以后至少要验证两件事一是项目测试能否证明 callback 没有漏掉增量内容二是最小 Demo 能否直观看到“边生成边输出”的效果。7.1 用项目测试验证 callback 和最终返回值项目里已经有对应测试TEST(DeepSeekProviderTest, SendMessageStream)先从环境变量读取 Keyconst char *apiKey std::getenv(DEEPSEEK_KEY_API); ASSERT_NE(apiKey, nullptr);初始化 Providerstd::mapstd::string, std::string config { {api_key, apiKey} }; auto provider std::make_sharedai_chat_sdk::DeepSeekProvider(); ASSERT_NE(provider, nullptr); ASSERT_TRUE(provider-initModel(config)); ASSERT_TRUE(provider-isAvailable());准备多轮消息std::vectorai_chat_sdk::Message messages { {user, 你好}, {assistant, 你好很高兴见到你。有什么可以帮助你的吗}, {user, 请解释一下什么是 SSE。} };然后在 callback 中累积实际收到的增量文本std::string streamedText; auto writeChunk [](const std::string chunk, bool isDone) { if (!chunk.empty()) { streamedText chunk; INFO(chunk: {}, chunk); } if (isDone) { INFO([DONE]); } };调用const std::string fullData provider-sendMessageStream(messages, requestParam, writeChunk);最后最关键的不是只检查“有没有返回”而是比较两条路径EXPECT_FALSE(fullData.empty()); EXPECT_FALSE(streamedText.empty()); EXPECT_EQ(fullData, streamedText);也就是callback 一段段收到并拼起来的 streamedText sendMessageStream() 最终返回的 fullData这能直接验证我们没有在中间漏掉某一段delta.content。7.2 最小 Demo直接看见“边生成边输出”单元测试主要验证正确性。如果想更直观地看流式效果可以写一个最小 Demo把 callback 中的文本直接输出到终端。示例目录examples/ └── deepseek_stream_demo/ ├── main.cpp └── CMakeLists.txt核心代码如下#include cstdlib #include iostream #include map #include memory #include string #include vector #include spdlog/spdlog.h #include ai_chat_sdk/DeepSeekProvider.h #include ai_chat_sdk/common.h #include ai_chat_sdk/util/myLOG.h int main() { myLOG::Logger::init_logger( DeepSeekStreamDemo, stdout, spdlog::level::info); const char *api_key std::getenv(DEEPSEEK_KEY_API); if (api_key nullptr || std::string(api_key).empty()) { std::cerr 请先设置环境变量 DEEPSEEK_KEY_API std::endl; return 1; } auto provider std::make_sharedai_chat_sdk::DeepSeekProvider(); std::mapstd::string, std::string config { {api_key, api_key} }; if (!provider-initModel(config)) { std::cerr DeepSeekProvider 初始化失败 std::endl; return 1; } std::vectorai_chat_sdk::Message messages { {user, 你好}, {assistant, 你好很高兴见到你。有什么可以帮助你的吗}, {user, 请用几句话解释一下什么是 SSE。} }; std::mapstd::string, std::string request_param { {temperature, 0.7}, {max_tokens, 2048} }; std::string streamed_text; auto write_chunk [](const std::string chunk, bool is_done) { if (!chunk.empty()) { streamed_text chunk; std::cout chunk std::flush; } if (is_done) { std::cout \n[DONE] std::endl; } }; std::cout DeepSeek 流式回复\n; const std::string full_data provider-sendMessageStream(messages, request_param, write_chunk); if (full_data.empty() || streamed_text.empty()) { std::cerr 没有得到有效流式回复 std::endl; return 1; } if (full_data ! streamed_text) { std::cerr fullData 与 streamedText 不一致 std::endl; return 1; } std::cout \n验证通过fullData 与 streamedText 一致 std::endl; return 0; }对应的CMakeLists.txt可以继续沿用前面“安装 SDK 后再由示例工程链接”的方式cmake_minimum_required(VERSION 3.10) project(DeepSeekStreamDemo LANGUAGES CXX) set(CMAKE_CXX_STANDARD 17) set(CMAKE_CXX_STANDARD_REQUIRED ON) add_executable(deepseek_stream_demo main.cpp) set(SDK_INSTALL_DIR ${CMAKE_CURRENT_LIST_DIR}/../../SDK/install) target_include_directories(deepseek_stream_demo PRIVATE ${SDK_INSTALL_DIR}/include ) target_compile_definitions(deepseek_stream_demo PRIVATE CPPHTTPLIB_OPENSSL_SUPPORT ) find_package(OpenSSL REQUIRED) find_package(Threads REQUIRED) find_package(SQLite3 REQUIRED) find_package(jsoncpp REQUIRED) find_package(spdlog REQUIRED) find_package(fmt REQUIRED) target_link_libraries(deepseek_stream_demo PRIVATE ${SDK_INSTALL_DIR}/lib/libai_chat_sdk.a jsoncpp_lib spdlog::spdlog fmt::fmt OpenSSL::SSL OpenSSL::Crypto SQLite::SQLite3 Threads::Threads )把示例目录放到项目根目录的examples/deepseek_stream_demo/后可以编译cd examples/deepseek_stream_demo mkdir -p build cd build cmake .. cmake --build . -j当前终端如果还没有 DeepSeek Key再设置export DEEPSEEK_KEY_API你的 DeepSeek API Key运行./deepseek_stream_democallback 中这一句std::cout chunk std::flush;会让已经收到的新内容尽快显示到终端而不是等整个回答结束后再一次性打印。写在最后上一篇完成的是sendMessage(...)这一篇完成的是sendMessageStream(...)现在DeepSeekProvider已经具备两条完整调用链全量响应 messages → HTTP → 完整 JSON → message.content → std::string以及流式响应 messages → HTTP → network chunk → buffer → SSE event → data payload → delta.content → callback → fullResponse这些需要注意的地方不是某一行字符串处理代码而是这几个层次不能混网络 chunk ≠ SSE event ≠ data payload ≠ 模型增量文本 delta.content把这一层理清以后后面再接入其他模型时我们就已经有了一套可以对照的 Provider 实现。下一篇继续沿着统一ILLMProvider往下走接入ChatGPTProvider看看另一套模型 API 怎样继续适配到同一套 SDK 接口中。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进