拿我自己的项目举个例,线上跑着一个基于 LangChain 调 Ollama 的对话服务,模型是 qwen2.5:7b。某天开始,用户反馈对话到一半突然没反应了,刷新页面后能看到之前的半句话,但继续发消息它就罢工。看日志,Ollama 那边有连接被重置的报错,LangChain 侧倒是没有任何异常抛出,就是流式输出静默停止。更诡异的是,偶发重试后竟然出现了两句回复拼接在一起的情况,像鬼打墙。后来我逐层拆解,发现根子不在 Ollama,也不在模型,而在客户端请求的生命周期处理上。今天就把这段排查经历掰开揉碎,讲给同样掉过坑的开发者听。
一、先别怪网络,看看请求是怎么活的
一个完整的 LLM 流式请求,从客户端视角看,大致可以拆成五个阶段:
- 请求发起:把用户的消息封装成 payload,交给 HTTP 层。
- 连接建立:TCP 握手、TLS 协商,然后发送 HTTP 请求头。
- 流式响应开始:服务端返回 200,
Content-Type: text/event-stream,然后逐行吐出 SSE 格式的数据。 - 流式数据持续到达:客户端边收边解析,交给回调函数处理。
- 流结束或异常终止:正常结束会收到
[DONE],异常则是连接断开、超时、被重置。
大部分开发者只关注“发起请求”和“处理 content”,却忽略了连接是活物,它有存活时间,会被各种因素掐死。
在我们这个案例里,LangChain 的 Ollama 类内部用的是 httpx 的流式接口。如果你不设置超时,默认情况下连接可以一直挂着,但这恰恰是第一个陷阱:没有空闲超时,不代表没有物理层超时。服务器或中间的反向代理(比如 Nginx)会隔一段时间就断开空闲连接,而客户端可能对此一无所知。
我画了一条生命线来帮助自己理解,你也可以在脑海里过一遍:
- 用户点击发送 → 请求进入队列
- LangChain 拿到用户消息 → 调用 Ollama
/api/chat接口,设置stream=True - Ollama 开始逐 token 返回
- 每个 token 到达后,LangChain 触发
on_llm_new_token回调 - 用户看到文字逐步生成
- 突然某个时间点,连接被服务器关闭,但客户端还在等下一个 token
问题就出在“客户端还在等”和“服务器已经断开了”之间。LangChain 的默认行为是,如果连接中途 EOF,它会认为生成结束,然后正常触发 on_llm_end。也就是说,它把半截话当成完整的话交付了。你的业务层看到会话没有异常,但用户看到的是“话说到一半没了”。
二、问题重现:模拟一个会断连的 Ollama 服务
为了彻底弄清楚,我写了一个模拟 Ollama 的 HTTP 服务,故意在生成 20 个 token 后断开连接。这样可以在本地复现线上问题,而不需要真的去折磨那台 GPU 服务器。
技术栈:Python 3.11 + FastAPI + LangChain 0.2
先写一个假的 Ollama SSE 接口:
# mock_ollama.py
from fastapi import FastAPI, Request
from fastapi.responses import StreamingResponse
import asyncio
import uuid
app = FastAPI()
@app.post("/api/chat")
async def chat(request: Request):
"""模拟Ollama的聊天接口,流式返回,第20个token后直接断开TCP连接"""
body = await request.json()
prompt = body.get("messages", [])[-1]["content"] if body.get("messages") else ""
async def event_generator():
# 先发送一个固定前缀
yield "data: {\"message\": {\"role\": \"assistant\", \"content\": \"\"}}\n\n"
# 模拟生成20个token,每个token之间间隔0.05秒
for i in range(20):
token = f"第{i+1}块-"
# 模拟Ollama的SSE输出格式
payload = f"data: {{\"message\": {{\"role\": \"assistant\", \"content\": \"{token}\"}}}}\n\n"
yield payload
await asyncio.sleep(0.05)
# 关键!第20个token后,直接关闭连接,不发送[DONE]
# 通过抛出GeneratorExit或者直接return,FastAPI会关闭连接
print(">>> 模拟服务主动断开连接")
# 这里直接return掉,但是客户端会收到EOF,而不是[DONE]
yield "data: [DONE]\n\n" # 注释掉这一行,就可以模拟中途断连
# 我们注释掉上面的[DONE],只让生成器结束
return StreamingResponse(event_generator(), media_type="text/event-stream")
注意上面代码里,我把 yield "[DONE]" 这一行注释掉了。这意味着生成器自然结束时,FastAPI 会正常关闭流,但由于没有 [DONE] 标记,客户端解析时就会认为连接是异常终止的。但更阴险的是,某些 HTTP 客户端在这种情况下会抛出 RemoteProtocolError 或直接返回空 content。
然后写一个使用 LangChain 的客户端,看看会发生什么:
# langchain_client.py
from langchain_community.chat_models import ChatOllama
from langchain_core.callbacks import BaseCallbackHandler
from langchain_core.messages import HumanMessage
import sys
class DebugHandler(BaseCallbackHandler):
"""用于追踪回调事件的自定义处理器"""
def on_llm_start(self, serialized, prompts, **kwargs):
print("[回调] 开始请求 LLM,prompts=", prompts)
def on_llm_new_token(self, token, **kwargs):
print("[回调] 收到新 token:", token)
def on_llm_end(self, response, **kwargs):
print("[回调] LLM 生成结束,完整输出:", response.generations[0][0].text)
def on_llm_error(self, error, **kwargs):
print("[回调] LLM 出错:", error)
# 这里如果不重新抛出,LangChain会吞掉异常
raise error # 注释掉这一行会怎么样?
# 创建一个不设置超时的Ollama连接
llm = ChatOllama(
base_url="http://localhost:8000", # 指向mock服务
model="fake-model",
streaming=True,
temperature=0,
)
handler = DebugHandler()
llm.callbacks = [handler]
# 发起一次请求
print("=== 开始请求 ===")
try:
result = llm.invoke([HumanMessage(content="给我讲个故事")])
print("=== 返回的result ===")
print(result.content)
except Exception as e:
print("!!! 捕获到异常:", type(e).__name__, str(e))
sys.exit(1)
运行这个客户端,输出会是什么?让我们走一遍:
=== 开始请求 ===
[回调] 开始请求 LLM,prompts= ['给我讲个故事']
[回调] 收到新 token:
[回调] 收到新 token: 第1块-
...
[回调] 收到新 token: 第20块-
[回调] LLM 生成结束,完整输出: 第1块-第2块-...第20块-
=== 返回的result ===
第1块-第2块-...第20块-
看到了吗?LangChain 没有抛异常,也没有触发 on_llm_error,而是正常地触发了 on_llm_end,把缺少 [DONE] 的连接终止当成了正常结束。这就是“静默断流”的真相。
三、为什么重试会让事情变得更糟
原始代码里我写了一版简单的重试逻辑:如果捕获到异常就重试整个请求。但问题是,之前的流已经吐了 20 个 token,这些 token 已经被前端通过回调渲染出去了。重试后,模型重新生成一遍,新的 token 会再次触发 on_llm_new_token 回调,前端把新 token 追加到了旧 token 后面,于是出现了两个回复拼接的“鬼打墙”现象。
这背后的核心错误是:没有把“流式响应”和“会话状态”当作一个整体来管理。流式响应本身是分段的,但如果你把每次重试都当作新的对话开始,而不是从中断点续传,那么客户端状态就会错乱。
正确的做法是:重试之前,必须把本次请求已经输出的 token 全部回滚。也就是说,前端需要把之前追加到界面上的内容撤销掉,然后再重新发起请求。而更优雅的方案是,客户端内部维护一个“会话生成缓冲区”,只有在流式连接完整结束(收到 [DONE])后,才把缓冲区内容提交给业务层。如果连接中断,缓冲区要清空,并触发重试。
我们来看一个带缓冲区和回滚能力的重试实现:
# robust_llm_client.py
import time
from langchain_community.chat_models import ChatOllama
from langchain_core.callbacks import BaseCallbackHandler
from langchain_core.messages import HumanMessage
class BufferCallbackHandler(BaseCallbackHandler):
"""
在内存中维护一个缓冲区,只有完整生成结束时才将内容写入最终结果。
如果中途发生异常,清空缓冲区。
"""
def __init__(self):
self.buffer = [] # 临时存放 token
self.final_text = None # 完整内容,仅当结束时赋值
self.error = None # 保存异常
def on_llm_new_token(self, token, **kwargs):
self.buffer.append(token)
# 这里不要往UI里输出,等待完整end
def on_llm_end(self, response, **kwargs):
# 只有当真正到达这里,且没有异常,才认为完整
self.final_text = "".join(self.buffer)
self.buffer = [] # 清空
print("[缓冲区] 收到完整结束标记,final_text长度:", len(self.final_text))
def on_llm_error(self, error, **kwargs):
self.error = error
# 清空缓冲区,准备重试
print("[缓冲区] 出错,清空缓冲区。错误内容:", error)
self.buffer = []
def invoke_with_retry(prompt, max_retries=3):
"""
带缓冲区和重试的调用函数。
策略:每次重试都重新创建handler,丢弃上一次所有未提交内容。
"""
for attempt in range(max_retries):
handler = BufferCallbackHandler()
llm = ChatOllama(
base_url="http://localhost:8000",
model="fake-model",
streaming=True,
temperature=0,
callbacks=[handler],
)
print(f"===== 第 {attempt+1} 次请求开始 =====")
try:
result = llm.invoke([HumanMessage(content=prompt)])
# invoke正常返回说明流式过程没有触发 on_llm_error
# 但是注意,如果连接中断没有被识别为错误,on_llm_end可能还会被触发
# 所以我们要额外校验 final_text 是否非空
if handler.final_text is None:
raise RuntimeError("final_text为空,说明流未完整结束")
return handler.final_text
except Exception as e:
print(f"!!! 第{attempt+1}次尝试失败: {type(e).__name__}: {e}")
time.sleep(1) # 简单退避
continue
raise RuntimeError("所有重试均失败")
看到问题了吗?on_llm_end 还是会触发。如果我们在 on_llm_end 里只判断 final_text 是否为空,是拦不住“半截话”的,因为 mock 服务没有发 [DONE],但是 FastAPI 正常关闭了流,LangChain 可能认为流正常结束并调用 on_llm_end。所以我们需要更底层的检查:如何知道流是否收到了完整的结束标记?
这里就有必要看一下 LangChain 内部的 Ollama 类是如何处理 SSE 的。在 langchain_community/llms/ollama.py 中,它读取 SSE 流,直到遇到 data: [DONE] 才终止循环。如果连接被服务端关闭,httpx 会在迭代下一个 chunk 时抛出 httpx.RemoteProtocolError 或返回空。LangChain 的代码里对异常的处理方式是:如果出现异常,会往错误队列里抛一个事件,从而触发 on_llm_error。但是,当连接优雅关闭(EOF)时,httpx 会把流结束当作正常结束,没有异常,所以 LangChain 的循环也会正常结束,并且不会发送 [DONE] 检查。这就导致它认为“流已结束”,进而触发 on_llm_end。
为了验证,我们可以在 mock 服务里加上乱发一个不完整的数据段,模拟半截 SSE 消息。但这里更简单的做法是:在回调里记录每次 on_llm_new_token 的 token 数,如果最终没有遇到 [DONE] 标记,那么 token 数一定不等于我们预期值。但我们不可能预先知道 token 数。所以根本解决之道是,重写或包装底层的 HTTP 流,自己来解析 SSE 事件,而不是全交给 LangChain。
四、自己掌控流式请求的生命周期
既然 LangChain 的黑盒行为不够透明,我们可以在客户端创建一个“流式请求包装器”,直接使用 httpx 发起请求,解析 SSE,然后手动调用 LangChain 的回调。这样连接何时断开、是否收到 [DONE],都清清楚楚。
技术栈:Python 3.11 + httpx + LangChain 回调。
下面是一个自制的流式调用器,它不依赖 ChatOllama 的 invoke,而是直接调用 Ollama API,但是通过 BaseCallbackHandler 与 LangChain 集成:
# manual_stream.py
import json
import httpx
from langchain_core.callbacks import BaseCallbackHandler, CallbackManager
from langchain_core.messages import HumanMessage
from langchain_core.outputs import LLMResult, Generation
class StreamingHandler(BaseCallbackHandler):
"""接收手动解析出来的令牌,并触发回调"""
def __init__(self):
self.token_count = 0
self.received_done = False # 是否收到[DONE]
self.buffer = []
def on_llm_new_token(self, token, **kwargs):
self.token_count += 1
self.buffer.append(token)
# 在这里可以做实时输出,但注意如果要回滚,需要让UI层支持撤销
# 为简化演示,只打印
print(f"token {self.token_count}: {token!r}")
def on_llm_end(self, response, **kwargs):
# 不在这里判断;我们会在外部判断 received_done
print("on_llm_end被触发,但我们会看received_done:", self.received_done)
def stream_ollama_with_manual_control(prompt, base_url="http://localhost:8000"):
"""
手动控制流式请求:
- 用httpx发起POST请求
- 一行一行解析SSE
- 只有收到[DONE]才认为完整
- 如果连接中断,抛出异常
"""
handler = StreamingHandler()
callback_manager = CallbackManager([handler])
# 回调开始事件
callback_manager.on_llm_start(
serialized={},
prompts=[prompt],
verbose=True
)
payload = {
"model": "fake-model",
"messages": [{"role": "user", "content": prompt}],
"stream": True,
}
# 设置超时,注意读写超时都要给足
timeout = httpx.Timeout(connect=5.0, read=30.0, write=5.0, pool=5.0)
# 使用with语句确保连接会被释放
with httpx.stream(
"POST",
f"{base_url}/api/chat",
json=payload,
timeout=timeout,
) as response:
if response.status_code != 200:
callback_manager.on_llm_error(Exception(f"HTTP {response.status_code}"))
raise RuntimeError(f"请求失败,状态码: {response.status_code}")
# 逐行读取SSE
try:
for line in response.iter_lines():
if not line.startswith("data:"):
continue # 忽略空行
data = line[5:].strip() # 去掉"data:"前缀
if data == "[DONE]":
handler.received_done = True
break # 正常结束
# 解析JSON
try:
json_data = json.loads(data)
token = json_data.get("message", {}).get("content", "")
if token:
# 手动触发新token回调
callback_manager.on_llm_new_token(token=token, verbose=True)
except json.JSONDecodeError as e:
# 如果是半截JSON,说明流断在数据中间,视为异常
callback_manager.on_llm_error(e)
raise RuntimeError(f"SSE数据解析失败,连接可能中断: {e}")
except httpx.RemoteProtocolError as e:
# 连接被重置,显式触发错误回调
callback_manager.on_llm_error(e)
raise RuntimeError(f"连接中断: {e}")
except httpx.ReadTimeout as e:
callback_manager.on_llm_error(e)
raise RuntimeError(f"读取超时: {e}")
# 根据是否收到[DONE]决定最终行为
if not handler.received_done:
error = RuntimeError("流结束但未收到[DONE],视为连接异常")
callback_manager.on_llm_error(error)
raise error
# 完整生成,生成LLMResult并触发on_llm_end
full_text = "".join(handler.buffer)
# 构造一个ChatResult/LLMResult对象,这里简单一点
generation = Generation(text=full_text)
result = LLMResult(generations=[[generation]])
callback_manager.on_llm_end(response=result, verbose=True)
return full_text
# 测试这个手动实现
if __name__ == "__main__":
try:
text = stream_ollama_with_manual_control("你好,请开始流式输出")
print("最终完整内容:", text)
except Exception as e:
print("捕获异常:", e)
这个手动实现有几点关键设计:
- 使用
httpx.stream而不是普通post,这样能逐行读取流。 - 用
response.iter_lines()逐行取得 SSE 数据。 - 显式检查
[DONE]。收到才标记完成。 - 如果连接中途断掉,
iter_lines可能抛出异常或正常结束。如果正常结束但没收到[DONE],我们抛出RuntimeError。 - 在异常路径中,我们显式调用
callback_manager.on_llm_error,这样上层可以知道失败原因。 - 完整结束时,我们手动构造
LLMResult并调用on_llm_end,这样其他依赖 LangChain 回调的组件(如 tracer)依然能正常工作。
在真机上测试,mock 服务在第 20 个 token 后断开,这个手动实现会立刻抛出 RuntimeError: 流结束但未收到[DONE],视为连接异常。而且 on_llm_error 也会被触发。这样重试机制就可以基于异常来进行,而不是盲目重试。
五、重试机制的正确姿势
有了上述能明确感知断连的调用器,下面设计重试机制。注意以下原则:
- 重试只发生在“尚未提交任何内容给用户”的场合。
- 如果已经通过
on_llm_new_token向 UI 实时输出,那么重试前必须通知 UI 回滚到本次请求开始前的状态。 - 重试次数要有上限,并且使用指数退避,避免雪崩。
- 重试次数用完后,要抛出明确的业务异常,而不是静默失败。
下面是一个集成了手动调用器的重试示例:
# retry_with_rollback.py
import time
import asyncio
# 假想一个前端界面对象的抽象,用于支持回滚操作
class UISession:
"""模拟前端会话界面,具有添加内容和回滚的功能"""
def __init__(self):
self.content = ""
self.checkpoints = [] # 栈式检查点
def snapshot(self):
"""在发起请求前记录一个快照"""
self.checkpoints.append(self.content)
def rollback(self):
"""回滚到最近一次快照"""
if self.checkpoints:
self.content = self.checkpoints.pop()
def append(self, text):
self.content += text
def invoke_with_retry_manual(ui_session, prompt, max_retries=3):
"""使用手动流控+UI快照回滚的重试"""
for attempt in range(max_retries):
# 每次重试前,先让UI回滚到本次请求开始前的位置
# 注意:第一次尝试时也需要有快照,所以在外部调用时先snapshot
# 这里简化处理,假设外部已经调用过 snapshot
# 如果之前的中断导致UI上已有部分token,则回滚
if attempt > 0:
ui_session.rollback()
print(f"[重试] 已回滚UI,当前UI内容: {ui_session.content!r}")
try:
text = stream_ollama_with_manual_control(prompt)
# 成功后,把完整内容追加到UI
ui_session.append(text)
return text
except Exception as e:
print(f"[重试] 第{attempt+1}次失败: {e}")
# 再次快照?不需要,因为下次重试时rollback会回到旧快照
time.sleep(2 ** attempt) # 指数退避:0,1,2,4秒(注意第一次0秒?这里用2**attempt)
continue
raise RuntimeError(f"重试{max_retries}次后仍然失败")
# 使用示例
if __name__ == "__main__":
ui = UISession()
ui.snapshot() # 初始快照
try:
result = invoke_with_retry_manual(ui, "讲个笑话")
print("重试成功后UI内容:", ui.content)
except RuntimeError as e:
print("最终失败,UI已回滚,内容为:", ui.content)
# 可以做降级处理
这个重试逻辑虽然简单,却解决了两个核心痛点:
- 每次重试都会回滚 UI,不会出现拼接问题。
- 只有收到完整
[DONE]才提交内容,否则线程不会退出,所以重试是有意义的。
六、从客户端请求生命周期入手的排查步骤
如果你也遇到了类似问题,我建议你按下面几步来排查,避免瞎猜:
6.1 抓取原始HTTP流
先用 curl 直接调 Ollama 接口,看是否稳定返回 [DONE]:
# 直接测试Ollama流式接口,看是否完整
curl -N -X POST http://localhost:11434/api/chat \
-H "Content-Type: application/json" \
-d '{
"model": "qwen2.5:7b",
"messages": [{"role": "user", "content": "你好"}],
"stream": true
}'
如果不设置超时,观察输出是否最后有 data: [DONE]。如果中间就断了,那么问题可能出在 Ollama 或网络层。如果完整,再继续下一步。
6.2 检查反向代理的 idle 超时
你的服务如果经过 Nginx,查看配置:
# nginx.conf 片段
server {
listen 80;
proxy_read_timeout 60s; # 如果模型生成超过60秒无数据,会被断开
proxy_send_timeout 60s;
proxy_buffering off; # 流式响应必须关闭缓冲
}
proxy_read_timeout 只要两次读操作之间的间隔超过阈值,Nginx 就会断开连接。一个长思考模型可能会卡在 “think” 阶段十几秒不发任何 token,这就会触发超时。
6.3 在LangChain侧配置超时和重试
如果你不想自己写 HTTP 层,也可以给 ChatOllama 设置 request_timeout,以及使用 LangChain 内置的重试逻辑。但注意,request_timeout 在流式场景下是 read 超时,不是总超时。设置太短会高频中断,设置太长又对死等不友好。建议设置成 120 秒。
# langchain_timeout_setting.py
from langchain_community.chat_models import ChatOllama
llm = ChatOllama(
base_url="http://localhost:11434",
model="qwen2.5:7b",
stream=True,
timeout=60, # 连接超时
num_predict=2048, # 限制生成长度
num_ctx=4096, # 上下文长度
)
6.4 监控回调事件队列
如果你使用 CallbackHandler,一定要在 on_llm_error 里把异常记录下来并上报监控,绝不能写成空实现。吞掉异常等于把问题埋进地雷。
七、应用场景、技术优缺点与注意事项
7.1 应用场景
这个排查思路适用于所有基于 LangChain 流式调用本地大模型的场景,比如:
- 企业内部知识库问答机器人,需要流式打字效果。
- 代码自动补全插件,后端调用 Ollama 流式生成。
- 智能客服系统,长回答中途可能被中断。
- 任何使用 SSE 作为传输格式的 AI 服务。
7.2 技术的优缺点
优点:
- 自己控制流式连接,可以获得完全透明的生命周期。
- 能够准确区分正常结束(
[DONE])和异常断开,不再静默吞错误。 - 结合回滚和重试,给用户一致的体验,不会出现半句或拼接。
缺点:
- 需要手动解析 SSE,工作量大。
- 绕开了 LangChain 内部的封装,可能会失去一些高级特性(如自动重试、token 统计等),都需要自己实现。
- 对开发者的网络编程能力要求更高。
7.3 注意事项
- 重试要幂等。大模型生成不是幂等的,重试会导致内容不同,所以回滚是必须的。
- 不要无限重试。建议最多 3 次,且每次等待时间翻倍。
- 流式连接要设置 read 超时,但不能太短。建议 30 到 120 秒之间。
- 处理半截 JSON。如果连接在
data: {"message": {...的中间断了,解析会失败,要能识别并触发错误回调。 - 前端要配合后端的回滚协议。如果你用后端服务代理 Ollama,后端重试时应该向前端发送一个
reset事件,让前端清掉已生成的半句话。 - 关注 Ollama 的 keep_alive 参数。如果模型不常驻内存,首次请求会加载模型,耗时增加,可能超出超时阈值。可以设置为
keep_alive: "5m"或更大。
八、文章总结
流式连接中断这个问题,表面上看起来像网络不稳定,但根子往往在客户端没有正确处理请求生命周期。我们需要把一次完整的流式请求分成:发起、连接、流式、结束/异常几个阶段,针对每个阶段设置明确的处理策略。尤其是“流式结束”和“异常断开”的区别,不能靠 LangChain 内部帮你区分,必须自己查看 [DONE] 标记或者捕获底层异常。
重试机制不能盲目地“再来一次”,要先回滚已经暴露给用户的中间结果,再用指数退避的方式重试。最好是把重试策略封装成独立的函数,配合回调处理器,让业务层感知不到底层波动。
虽然手动控制流式请求会多一些代码,但换来的是可观测性和可控性。在 AI 应用越来越复杂的今天,把命运的缰绳握在自己手里,总比被第三方库的隐式逻辑牵着走强。希望这篇文章能让你在下一次遇到“话说到一半就没了”的时候,不再抓狂,而是胸有成竹地打开代码,定位到那一个被遗忘的生命周期节点。
评论
围绕“基于LangChain调用Ollama流式输出时连接中断,重试机制与回调事件处理不当造成会话异常,从客户端请求生命周期入手的排查”参与讨论