在调用支持流式输出的 AI 接口时,网络波动可能导致连接意外中断。本文介绍如何在 Python 中底层解析 Server-Sent Events (SSE) 数据流,并实现网络抖动时的断线重连与数据校验。
为什么需要自定义 SSE 解析与重连?
在前面的文章中,我们介绍过使用 stream=True 实现流式输出。但在实际生产环境中,流式请求往往比普通请求更容易受到网络波动的影响:
如果直接使用高层 SDK,连接一断,整个生成过程就会抛出异常中断,用户界面或日志中只看到“半截”回答。
为了让流式应用更加健壮,我们需要了解底层的 Server-Sent Events (SSE) 协议,并编写具备断线重连和增量解析能力的客户端。
一、理解 SSE 协议的底层格式
Server-Sent Events 是一种基于 HTTP 的服务器推送标准。AI 接口的流式响应(如 OpenAI 兼容的 /v1/chat/completions)正是基于这种格式。
一个典型的 SSE 响应流长这样:
HTTP/1.1 200 OKContent-Type: text/event-streamdata: {"id":"chatcmpl-123","choices":[{"delta":{"content":"你好"}}]}data: {"id":"chatcmpl-123","choices":[{"delta":{"content":"世界"}}]}data: [DONE]
核心规则很简单:
- 当流结束时,会发送特定的结束标记(如
data: [DONE])。
二、用 Python requests 库实现基础 SSE 解析
如果不使用 OpenAI SDK,而是用底层的 requests 库以流式读取(stream=True),我们可以手动解析每一行数据:
import jsonimport requestsurl = "https://your-api-domain.com/v1/chat/completions"headers = { "Authorization": "Bearer your-api-key", "Content-Type": "application/json",}payload = { "model": "your-model-name", "messages": [{"role": "user", "content": "写一首短诗"}], "stream": True,}response = requests.post(url, headers=headers, json=payload, stream=True, timeout=30)response.raise_for_status()# 逐行读取流式响应for line in response.iter_lines(decode_unicode=True): if not line: continue # 跳过空行 # SSE 消息固定以 "data: " 开头 if line.startswith("data: "): data_str = line[len("data: "):].strip() # 检查是否为结束标记 if data_str == "[DONE]": break try: packet = json.loads(data_str) content = packet["choices"][0]["delta"].get("content") if content: print(content, end="", flush=True) except json.JSONDecodeError: # 兼容偶尔出现的空包或心跳注释 pass
三、处理网络抖动:实现带重试的流式消费器
在真实网络环境下,requests.post 可能会因为 ChunkedEncodingError 或 ConnectionError 中途报错。
我们可以加入重试与断点续传(如果接口支持 Last-Event-ID 或上下文允许重新发起)的思路。对于大多数大模型对话来说,最简单的降级容错是:当连接意外中断时,捕获异常并尝试重新发起请求,或者向用户返回友好提示。
下面是一个带异常捕获与重试的 SSE 消费函数:
import timeimport jsonimport requestsfrom typing import Generatordef stream_with_retry(url: str, headers: dict, payload: dict, max_retries: int = 3) -> Generator[str, None, None]: retries = 0 while retries < max_retries: try: with requests.post(url, headers=headers, json=payload, stream=True, timeout=25) as response: response.raise_for_status() for line in response.iter_lines(decode_unicode=True): if not line: continue if line.startswith("data: "): data_str = line[len("data: "):].strip() if data_str == "[DONE]": return packet = json.loads(data_str) content = packet["choices"][0]["delta"].get("content") if content: yield content # 正常结束跳出循环 break except (requests.exceptions.ChunkedEncodingError, requests.exceptions.ConnectionError, requests.exceptions.Timeout) as e: retries += 1 print(f"\n[网络波动] 流式连接中断 ({e}),正在进行第 {retries}/{max_retries} 次重试...") if retries >= max_retries: raise RuntimeError("流式传输多次重试后彻底失败") from e time.sleep(2 * retries)
四、增量数据校验与过滤
在流式输出过程中,由于网络分片的原因,偶、尔可能会收到格式不完整的 JSON 片段。虽然 iter_lines 通常能按行切分,但在高并发或特殊代理环境下,依然建议在解析时做好防御性编程:
def safe_parse_chunk(line: str) -> str | None: if not line.startswith("data: "): return None data_str = line[6:].strip() if data_str == "[DONE]": return None try: obj = json.loads(data_str) # 兼容不同厂商可能存在的字段差异 choices = obj.get("choices") if not choices: return None delta = choices[0].get("delta", {}) return delta.get("content") except (json.JSONDecodeError, KeyError, IndexError): # 过滤掉无效或畸形的片段,防止程序崩溃 return None
五、在实际项目中的应用建议
- 选择高层封装优先:如果使用的是官方的
openai SDK(如 client.chat.completions.create(..., stream=True)),SDK 内部已经帮我们处理了大部分底层的行解析。只有在需要特殊代理、自定义网关或绕过 SDK 限制时,才需要自己写底层的 requests 解析。 - 区分断线重试与业务失败
- 只有当发生底层网络中断、超时、连接重置时,才进行流式重试。
- 如果接口已经返回了 HTTP 400、401、429 等状态码,说明请求本身有误或被限流,此时不应该盲目重试。
- 前端配合:如果后端支持流式重试,前端也需要能够识别重试信号,或者在界面上平滑拼接多次收到的文本片段。
六、结语
流式输出(SSE)提升了用户体验,但也带来了更高的网络稳定性要求。通过掌握底层行解析逻辑、捕获网络异常、编写带重试的生成器,我们可以让 AI 应用在面对复杂网络环境时表现得更加稳健。
对于 Python 开发者来说,理解这些底层协议不仅有助于排查各种诡异的“断流”问题,也能在构建自定义 AI 网关和聚合工具时游刃有余。
免责声明
本文内容仅用于技术交流与经验分享,具体实现请结合项目实际网络环境调整。