Python 服务用 HTTPX 转发大文件、对象存储下载或上游 SSE 时,通常会启用流式响应,避免把整个响应体读入内存。上线后却可能出现一种很迷惑的故障:服务启动时一切正常,运行一段时间后请求开始批量报 httpx.PoolTimeout;重启进程能暂时恢复,单纯调大连接池只能延后复发。
这类问题往往不是上游连接慢,而是代码取得流式响应后,没有在成功、异常和下游断开等所有路径上关闭它。连接一直被响应对象占用,最终没有空闲连接可供新请求获取。
本文以异步 HTTP 代理为例,说明如何确认连接池耗尽的真正原因,如何正确使用自动和手动流式模式,以及怎样用并发测试与运行指标验证修复。
适用场景
本文适用于以下情况:
- 使用
httpx.AsyncClient调用 HTTP 或 HTTPS 上游; - 为控制内存占用,使用
client.stream()或client.send(..., stream=True); - 服务转发文件、SSE、分块响应或未知长度的响应体;
- 故障以
PoolTimeout为主,重启后暂时消失; - 连接池上限已经合理,但活跃请求数下降后连接仍不能恢复。
普通的 await client.get(url) 默认会读取完整响应体,HTTPX 随后可以回收连接。本文讨论的是显式开启流式模式后,响应体由业务代码负责消费和关闭的场景。
现象描述
典型日志如下:
获取上游文件失败 error=PoolTimeout operation=download_proxy duration_ms=1003
常见时间线是:
- 服务配置
max_connections=20; - 前 20 个流式请求成功取得响应头,但响应对象没有关闭;
- 即使客户端已经下载完成或中途断开,这 20 个连接仍未回到连接池;
- 第 21 个请求等待池中可用连接;
- 等待超过
pool超时后抛出httpx.PoolTimeout; - 重启进程关闭了整个客户端和底层连接,所以故障暂时消失。
下面是一个危险的代理实现:
import httpx
from fastapi import FastAPI
from starlette.responses import StreamingResponse
app = FastAPI()
http_client = httpx.AsyncClient()
@app.get("/files/{file_id}")
async def download_file(file_id: str) -> StreamingResponse:
request = http_client.build_request(
"GET",
f"https://storage.example.internal/files/{file_id}",
)
upstream = await http_client.send(request, stream=True)
upstream.raise_for_status()
# StreamingResponse 接管了迭代,却没有接管 upstream 的关闭责任。
return StreamingResponse(upstream.aiter_raw())
send(..., stream=True) 属于手动流式模式。函数返回后,创建响应的调用栈已经结束,但 upstream 仍然持有底层连接。代码既没有 async with,也没有调用 await upstream.aclose(),因此响应生命周期没有闭环。
先区分四种不同的超时
HTTPX 把超时拆成四类,排查时不能只看“请求超时”:
connect:建立 TCP 连接或完成相关连接步骤超时;read:等待下一段响应数据超时;write:发送一段请求数据超时;pool:等待连接池提供一个可用连接超时。
建议显式配置,避免所有阶段共用一个难以解释的数字:
import httpx
timeout = httpx.Timeout(
connect=3.0,
read=30.0,
write=10.0,
pool=1.0,
)
limits = httpx.Limits(
max_connections=100,
max_keepalive_connections=20,
keepalive_expiry=5.0,
)
http_client = httpx.AsyncClient(timeout=timeout, limits=limits)
关键字段的含义如下:
max_connections限制连接池允许同时占用的连接总数;max_keepalive_connections限制池中保留的空闲长连接数;keepalive_expiry控制空闲连接可以保留多久;pool只限制等待可用连接的时间,不限制上游读取完整响应体的总时长。
如果异常是 PoolTimeout,首先检查连接占用和并发队列,而不是立刻把 read 超时调大。
排查思路
1. 保留异常类型,不要只记录字符串
HTTPX 的超时异常有继承关系。应单独统计 PoolTimeout,并保留请求目标的受控信息:
import logging
import httpx
logger = logging.getLogger(__name__)
async def fetch_metadata(client: httpx.AsyncClient, url: str) -> dict[str, object]:
try:
response = await client.get(url)
response.raise_for_status()
return response.json()
except httpx.PoolTimeout as exc:
logger.warning(
"等待 HTTP 连接池超时",
extra={"operation": "fetch_metadata", "error_type": type(exc).__name__},
)
raise
except httpx.HTTPError as exc:
logger.error(
"调用上游 HTTP 接口失败",
extra={"operation": "fetch_metadata", "error_type": type(exc).__name__},
)
raise
不要记录 Authorization、Cookie、签名查询串或完整响应体。URL 如果可能携带临时凭据,应只记录主机名、固定路由模板和受控的上游名称。
2. 搜索所有手动流式入口
重点搜索以下调用:
rg -n "stream=True|\.send\(|\.stream\(" .
对每一个 stream=True,逐项确认:
- 正常读取完成后是否关闭;
raise_for_status()抛异常时是否关闭;- 处理响应体时抛异常是否关闭;
- 下游客户端断开时是否关闭;
- 任务取消时是否进入
finally; - 整个应用停机时是否关闭共享的
AsyncClient。
只检查成功路径不够。连接泄漏常发生在 404、上游读取中断或下游取消等低比例分支,所以通常要运行一段时间才暴露。
3. 记录业务层的在途流数量
HTTPX 没有把连接池内部结构承诺为稳定公共接口,不建议在线上依赖私有属性读取连接数。更稳妥的做法是在业务边界记录:
http_stream_active:当前正在转发的流;http_stream_open_total:已打开的流;http_stream_close_total:已关闭的流;http_pool_timeout_total:连接池等待超时次数;http_stream_duration_seconds:从取得响应头到流关闭的时长。
在流量回落后,active 应回到接近零,打开和关闭计数的增量应大体一致。如果活跃业务请求已经结束,但关闭计数持续落后,优先检查生命周期泄漏。
修复方案一:在同一作用域内读取
如果业务逻辑能在当前函数中消费响应体,优先使用 async with client.stream(...)。上下文退出时 HTTPX 会关闭响应,包括异常和取消路径:
from collections.abc import AsyncIterator
import httpx
async def iter_audit_lines(
client: httpx.AsyncClient,
url: str,
) -> AsyncIterator[str]:
async with client.stream("GET", url) as response:
response.raise_for_status()
async for line in response.aiter_lines():
if line:
yield line
这里的关键不是 aiter_lines(),而是响应的使用范围被 async with 明确包住。调用方停止迭代或任务被取消时,异步生成器必须被正常关闭;框架边界还应通过取消测试确认这一点。
不要在热点函数里反复创建 AsyncClient。客户端应在应用生命周期内复用,才能获得连接池收益,并在应用停机时统一执行 await client.aclose()。
修复方案二:把关闭责任交给代理响应
代理场景需要先取得上游响应头,再把响应体交给 Web 框架发送。此时无法让创建函数一直停留在普通的 async with 代码块中,可以使用手动流式模式,但必须把 aclose 注册到下游响应的后台清理阶段。
from contextlib import asynccontextmanager
from collections.abc import AsyncIterator
import httpx
from fastapi import FastAPI, HTTPException, Request
from starlette.background import BackgroundTask
from starlette.responses import StreamingResponse
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
app.state.http_client = httpx.AsyncClient(
timeout=httpx.Timeout(connect=3.0, read=30.0, write=10.0, pool=1.0),
limits=httpx.Limits(max_connections=100, max_keepalive_connections=20),
)
try:
yield
finally:
await app.state.http_client.aclose()
app = FastAPI(lifespan=lifespan)
@app.get("/files/{file_id}")
async def download_file(file_id: str, request: Request) -> StreamingResponse:
client: httpx.AsyncClient = request.app.state.http_client
upstream_request = client.build_request(
"GET",
f"https://storage.example.internal/files/{file_id}",
)
try:
upstream = await client.send(upstream_request, stream=True)
upstream.raise_for_status()
except httpx.HTTPStatusError as exc:
await exc.response.aclose()
raise HTTPException(status_code=502, detail="上游返回异常状态") from exc
except httpx.HTTPError as exc:
raise HTTPException(status_code=502, detail="上游连接失败") from exc
forwarded_headers = {}
content_type = upstream.headers.get("content-type")
if content_type is not None:
forwarded_headers["content-type"] = content_type
return StreamingResponse(
upstream.aiter_raw(),
status_code=upstream.status_code,
headers=forwarded_headers,
background=BackgroundTask(upstream.aclose),
)
这段实现有三个重要边界:
raise_for_status()可能抛出HTTPStatusError,此时尚未创建下游StreamingResponse,必须立即关闭上游响应;- 成功后由
BackgroundTask(upstream.aclose)在响应发送结束后回收连接; - 共享的
AsyncClient由 FastAPI lifespan 创建和关闭,而不是由单个请求随意销毁。
示例只转发了 Content-Type。实际代理不能无条件复制所有上游响应头;Connection、Transfer-Encoding 等逐跳头应由客户端和服务器协议栈管理,Content-Length 也必须与实际转发内容一致。
一个更容易审计的关闭包装器
如果所在框架不能保证后台清理在断开路径执行,可以把关闭动作放进异步迭代器的 finally,让资源所有权和数据迭代绑定在一起:
from collections.abc import AsyncIterator
import httpx
async def iter_and_close(response: httpx.Response) -> AsyncIterator[bytes]:
try:
async for chunk in response.aiter_raw():
yield chunk
finally:
await response.aclose()
然后将 iter_and_close(upstream) 交给框架。必须用集成测试确认:下游主动断开时,框架确实会关闭异步生成器并执行 finally。不同 ASGI 服务器和中间件组合的取消传播行为可能不同,不能只凭代码阅读下结论。
并发复现:小连接池能更快暴露泄漏
测试环境可以故意把连接池缩小,把“运行几小时后才出现”变成几次请求即可复现。下面的脚本要求 TARGET_URL 指向一个会持续输出数据的测试端点,不能对生产接口执行:
import asyncio
import os
import httpx
async def main() -> None:
target_url = os.environ["TARGET_URL"]
timeout = httpx.Timeout(connect=2.0, read=10.0, write=2.0, pool=0.2)
limits = httpx.Limits(max_connections=2, max_keepalive_connections=0)
async with httpx.AsyncClient(timeout=timeout, limits=limits) as client:
leaked_responses: list[httpx.Response] = []
try:
for request_number in range(3):
request = client.build_request("GET", target_url)
response = await client.send(request, stream=True)
leaked_responses.append(response)
print(f"已打开第 {request_number + 1} 个响应")
except httpx.PoolTimeout:
print("复现成功:第三个请求等待连接池超时")
finally:
for response in leaked_responses:
await response.aclose()
if __name__ == "__main__":
asyncio.run(main())
运行方式:
TARGET_URL=http://127.0.0.1:8001/slow-stream python reproduce_pool_timeout.py
Windows PowerShell 可以使用:
$env:TARGET_URL = "http://127.0.0.1:8001/slow-stream"
python .\reproduce_pool_timeout.py
如果测试端点在返回响应头后持续输出,前两个未关闭的响应会占满两个连接,第三个请求应在约 200 毫秒后触发 PoolTimeout。把循环改成 async with client.stream(...) 并在块内读取少量数据后退出,随后重复多轮,连接池应持续可用。
测试应覆盖哪些失败路径
只验证完整下载成功不足以证明资源会回收。至少覆盖以下场景:
- 上游返回 404 或 500,
raise_for_status()抛出异常; - 上游发送部分响应体后断开;
- 下游客户端读取一部分内容后主动断开;
- 转发任务在读取期间被取消;
- 并发数短时超过
max_connections,请求排队后能随连接释放继续执行; - 应用停机时仍有流式请求,客户端可以在有限时间内关闭;
- 连续执行数百轮后,打开与关闭计数没有持续扩大差值。
集成测试应使用真实的 ASGI 服务器和一个可控的本地上游,而不是只用 MockTransport。MockTransport 适合验证请求和响应逻辑,但不经过真实网络连接池,无法证明连接槽是否正确释放。
常见误区
只调大 max_connections
如果每个请求都会泄漏连接,把上限从 20 调到 200 只会让故障晚一点发生,还会增加文件描述符、上游连接和内存压力。先闭合生命周期,再根据峰值并发和上游承载能力调整池大小。
对 PoolTimeout 立即重试
连接槽没有释放时,立即重试只会增加等待者数量。只有在确认故障是短时并发尖峰而非泄漏后,才考虑带上限、退避和随机抖动的重试;流式下载还要确认是否已经向下游发送了数据,避免从头重放造成内容重复。
以为读到 EOF 等于所有路径都会回收
完整消费响应体通常可以让连接复用,但异常、取消和提前停止迭代仍可能绕开 EOF。显式的上下文管理或 finally 比依赖“正常情况下总会读完”更可靠。
在循环里创建客户端
每次请求都创建 AsyncClient 会丢失连接复用,并制造额外握手和套接字。客户端应按应用或明确的上游隔离维度复用,同时设置统一的生命周期、超时、连接池和 TLS 配置。
依赖 HTTPX 私有连接池属性做监控
私有属性可能随版本变化。线上告警应以稳定的业务指标、异常类型、在途任务和上游状态为基础;内部状态只适合临时诊断,并要锁定版本和做好兼容处理。
上线检查清单
上线前逐项确认:
- 所有
stream=True都能找到明确的aclose()所有者; - 优先使用
async with client.stream(...); - 手动流式响应在状态异常、取消和下游断开路径都会关闭;
AsyncClient在应用生命周期内复用,并在停机阶段关闭;connect、read、write、pool超时分别配置;- 连接池上限与应用并发、进程数和上游容量匹配;
- 集成测试使用真实连接池复现提前断开和异常路径;
- 指标能区分连接池等待超时与上游读取超时;
- 日志不包含令牌、Cookie、签名 URL 或原始敏感响应;
- 灰度期间观察打开/关闭流差值、
PoolTimeout和上游连接数。
参考资料
总结
PoolTimeout 表示请求没能在规定时间内从连接池取得连接,它不等于上游建立连接慢,也不等于响应读取超时。对于流式请求,最值得先查的是响应是否在所有路径上关闭。
可靠的修复不是无限调大连接池,而是明确资源所有权:能在当前作用域消费就使用 async with client.stream(...);必须把响应交给框架时,就把 aclose() 一并交给可靠的清理机制;共享客户端则由应用生命周期统一管理。最后用小连接池、真实网络、异常状态、下游断开和任务取消做集成验证,才能证明连接生命周期真正闭环。
Discussion
评论