Python 异步编程与 连接池 实践
前段时间把手头一个同步的采集服务改成 asyncio 版本,改之前想得挺美,觉得换成异步以后并发随便拉,吞吐量翻几倍不是问题。结果一压测就露馅了,并发开到几百的时候请求开始大量超时,日志里全是 TimeoutError 和 ServerDisconnectedError,反而比同步版本还不稳定。折腾了两天才把问题理清楚,发现锅基本都在连接池上。网上讲 asyncio 的教程大多只教你怎么 async with 发请求,很少讲连接池到底该怎么配,这里把踩过的坑记录一下。
先泼盆冷水:异步不等于快
asyncio 解决的是“等待的时候别干等着”的问题,不是把你的带宽和对方的处理能力变强了。我最开始的想法是把任务列表直接丢给 asyncio.gather,几百个协程同时跑,美滋滋。但 aiohttp 底层的 TCPConnector 默认 limit=100,也就是说不管你开了多少协程,真正同时在飞的连接就 100 条,剩下的全在池子里排队。更坑的是这个排队过程是静默的,不报错也不打日志,表现就是请求莫名变慢、然后超时,你根本不知道瓶颈在哪。这么关键的默认值,官方文档把它埋在参数表格里,可能是我瞎没看到吧...
所以第一件事是把两个概念分开:协程数量是你想发多少请求,连接池限制是同一时刻真正能建立多少条 TCP 连接,后者才是实际瓶颈。
aiohttp 这边的实际配置
connector = aiohttp.TCPConnector(
limit=128, # 总连接数上限
limit_per_host=32, # 单个主机的连接数,抓同一家站点的时候一定要设
ttl_dns_cache=300, # DNS 缓存,默认才 10 秒,高并发下解析也是开销
keepalive_timeout=30,
enable_cleanup_closed=True,
)
limit_per_host 这个是我踩坑之后才加上去的,之前没设,同一个目标站点瞬间被打出限流,直接开始 429。ttl_dns_cache 默认只有 10 秒,压测的时候等于每 10 秒全体协程重新解析一遍域名,如果你的服务跑在内网、解析特别快,没准儿感知不到,但公网环境下这个差异挺明显的。
连接要及时还回去
还有一个隐蔽的问题:响应体没读完就把响应扔了,这条连接是没法复用的,只能断开重建,TLS 握手的开销全浪费了。正确的姿势是老老实实在上下文里把数据读出来:
# 错误示范
async def fetch(session, url):
await session.get(url) # 响应没读,连接状态很尴尬
# 正确姿势
async def fetch(session, url):
async with session.get(url) as resp:
return await resp.text()
反过来也一样,别拿到响应之后在上下文里做耗时的处理,连接一直被占着,后面的请求排队等着,池子等于白配。
Semaphore 和连接池要一起看
落到代码上,并发和连接数就是两个东西各管各的:
sem = asyncio.Semaphore(64)
async def fetch(session, url):
async with sem:
async with session.get(url) as resp:
return await resp.text()
我的习惯是 Semaphore 稍微给大一点,让请求在应用层排队,连接池保持一个和下游能承受的量级匹配的值。要是 Semaphore 远大于连接池,请求就堆在 aiohttp 内部排队,超时了都死在里面,特别不好排查。
数据库这边:asyncpg 的坑
HTTP 这边稳定了,结果数据库又炸了。最开始是每个请求里直接 asyncpg.connect() 裸连,几百个并发下去直接把 PostgreSQL 打满。老老实实换成池子:
pool = await asyncpg.create_pool(
dsn="postgresql://user:password@127.0.0.1:5432/mydb",
min_size=5,
max_size=10,
command_timeout=30,
max_inactive_connection_lifetime=300,
)
async with pool.acquire() as conn:
rows = await conn.fetch("SELECT id, url FROM task WHERE status = $1", "pending")
这里有两个点。一个是 max_size 不是越大越好,我一开始手一抖给了 20,然后服务起了 4 个 worker,每个 worker 各自一个池子,4 × 20 = 80 条连接,加上别的服务也在连这个库,PostgreSQL 默认 max_connections=100,直接 too many clients already 了。所以池子大小要按 worker 数量乘一下再定。
另一个是 max_inactive_connection_lifetime,数据库或者中间的代理一般会主动断开空闲连接,池子里的老连接失效了你还拿去用就会报错,设个几百秒让它自己回收重建,能省不少心。程序退出的时候记得 await pool.close(),不然收尾的时候 warning 刷屏。
断连和重试兜底
长连接绕不开一个问题:你这边觉得连接还活着,服务端已经把它关了,表现就是偶发的 ServerDisconnectedError,复现还看概率。除了把客户端的 keepalive_timeout 调到比服务端空闲超时短,反正是幂等的 GET,重试也没什么心理负担:
for attempt in range(3):
try:
async with session.get(url) as resp:
return await resp.text()
except (aiohttp.ServerDisconnectedError, asyncio.TimeoutError):
if attempt == 2:
raise
await asyncio.sleep(2 ** attempt) # 1s, 2s 退避
最后的完整配置
import asyncio
import aiohttp
sem = asyncio.Semaphore(64)
timeout = aiohttp.ClientTimeout(total=15, connect=5)
connector = aiohttp.TCPConnector(
limit=128,
limit_per_host=32,
ttl_dns_cache=300,
keepalive_timeout=30,
enable_cleanup_closed=True,
)
async def fetch(session: aiohttp.ClientSession, url: str) -> str:
async with sem:
for attempt in range(3):
try:
async with session.get(url) as resp:
return await resp.text()
except (aiohttp.ServerDisconnectedError, asyncio.TimeoutError):
if attempt == 2:
raise
await asyncio.sleep(2 ** attempt)
async def main():
async with aiohttp.ClientSession(connector=connector, timeout=timeout) as session:
urls = [...] # 你的任务列表
results = await asyncio.gather(*(fetch(session, u) for u in urls), return_exceptions=True)
print(results)
asyncio.run(main())
把这套配上去之后,同样的并发量,超时率从百分之几十降到基本没有,数据库那边也再没报过连接数的问题。总结下来就三句话:协程数、连接池、下游承受能力是三个东西,要分开看;连接用完及时还,断连和超时靠重试兜底;池子大小记得乘上 worker 数量。要是你也在做异步改造,建议先把连接池这块理清楚,能少熬不少夜...
评论
还没有评论。
发表评论
提交后评论将经过自动审核,审核通过后公开展示。