1. 问题背景:线程池被打满,CPU却在看戏

上个月接手一个订单服务,核心接口/api/v1/orders需要聚合用户信息、商品快照、物流状态三个下游服务。最初实现是同步Flask视图:

@app.route('/api/v1/orders')
def get_orders():
    user = requests.get(f"http://user-svc/users/{uid}", timeout=2).json()
    products = requests.post("http://product-svc/batch", json=pid_list, timeout=2).json()
    logistics = requests.get(f"http://logistics-svc/tracks/{oid}", timeout=2).json()
    return jsonify(assemble(user, products, logistics))

压测结果惨不忍睹:单实例QPS 120,P95延迟380ms,最要命的是ThreadPoolExecutor的20个线程全部卡在requests.get的socket等待上,CPU占用不到30%。这显然是IO密集型瓶颈,同步IO把线程全占死了,机器资源完全浪费。

2. 环境与版本:Python 3.10 + aiohttp

  • Python 3.10.8(原生asyncio,没有额外装uvloop)
  • aiohttp 3.8.4(替代requests做HTTP调用)
  • Flask 2.2.3(保留框架,但视图改为异步)
  • gunicorn 20.1.0 + uvicorn workers(注意后面有坑)
  • wrk 4.2.0压测工具

为什么不用FastAPI?因为历史包袱,路由和中间件都是Flask的,迁移成本太高。好在Flask 2.2支持异步视图,只要装了asgiref兼容层就能跑。

3. 方案设计:全链路协程化 + 信号量限流

目标:把同步IO全部替换为异步IO,让单进程能同时处理上千个并发请求。

设计要点:
1. 用aiohttp.ClientSession代替requests,复用TCP连接
2. 用asyncio.gather并行发起三个下游调用
3. 用asyncio.Semaphore(100)控制并发量,防止打爆下游
4. 每个下游调用单独设置超时,防止整体超时

核心数据结构是asyncio.QueueSemaphore,配合asyncio.timeout(Python 3.11之前用asyncio.wait_for)。

4. 核心实现:改造后的异步视图

import asyncio
import aiohttp
from flask import Flask, jsonify, request

app = Flask(__name__)
# 全局session,复用连接池
session = None
# 信号量限制最大并发100个下游请求
semaphore = asyncio.Semaphore(100)
TIMEOUT = aiohttp.ClientTimeout(total=2.0)

async def fetch_json(client, url, payload=None):
    async with semaphore:
        try:
            if payload:
                async with client.post(url, json=payload) as resp:
                    return await resp.json()
            else:
                async with client.get(url) as resp:
                    return await resp.json()
        except asyncio.TimeoutError:
            return {"error": "timeout"}
        except aiohttp.ClientError as e:
            return {"error": str(e)}

async def get_orders_async(uid, pid_list, oid):
    async with aiohttp.ClientSession(timeout=TIMEOUT) as client:
        user_task = fetch_json(client, f"http://user-svc/users/{uid}")
        product_task = fetch_json(client, "http://product-svc/batch", payload={"ids": pid_list})
        logistics_task = fetch_json(client, f"http://logistics-svc/tracks/{oid}")
        user, products, logistics = await asyncio.gather(user_task, product_task, logistics_task)
        return assemble(user, products, logistics)

@app.route('/api/v1/orders', methods=['GET'])
async def get_orders():
    uid = request.args.get('uid')
    pid_list = request.args.getlist('pid')
    oid = request.args.get('oid')
    result = await get_orders_async(uid, pid_list, oid)
    return jsonify(result)

# 应用启动时创建全局session
@app.before_first_request
def init_session():
    global session
    session = aiohttp.ClientSession(timeout=TIMEOUT)

注意before_first_request在Flask 2.3已被移除,生产环境我用的是app.before_request加标志位。另外ClientSession必须复用,每次新建会浪费大量TCP握手时间。

5. 踩坑与优化:两个隐秘的性能杀手

坑1:Flask同步视图下的asyncio会失效

最初我直接在同步视图里写asyncio.run(get_orders_async(...)),结果QPS反而降到80。因为asyncio.run每次都会创建新的事件循环,线程还是被阻塞在事件循环的run_until_complete上。必须用async def定义视图,让gunicorn的worker直接驱动协程。

坑2:DNS解析阻塞EventLoop

压测时发现延迟抖动剧烈,用py-spy dump抓到栈,发现大量协程卡在getaddrinfo上。aiohttp默认使用线程池解析DNS,但线程池只有5个,一旦打满就阻塞整个loop。解决:

import socket
# 强制使用异步DNS解析
resolver = aiohttp.AsyncResolver(nameservers=["8.8.8.8", "1.1.1.1"])
connector = aiohttp.TCPConnector(resolver=resolver, use_dns_cache=True, ttl_dns_cache=300)
session = aiohttp.ClientSession(connector=connector, timeout=TIMEOUT)

坑3:gunicorn的worker类型

必须用uvicorn.workers.UvicornWorker,否则gunicorn的同步worker会把asyncio视图当普通函数跑,完全不起作用。启动命令:

gunicorn -w 4 -k uvicorn.workers.UvicornWorker -b 0.0.0.0:8000 app:app

6. 效果数据:QPS翻8倍,延迟降75%

压测环境:4核8G云主机,wrk -t8 -c500 -d30s。

指标 同步版本 异步版本
QPS 120 980
P50延迟 250ms 42ms
P95延迟 380ms 89ms
P99延迟 610ms 210ms
CPU占用 28% 75%
异常率 0.5% 0.1%

异步版本把并发能力从20个线程扩展到了上千协程,CPU终于跑满了。下游服务压力测试显示,商品服务的QPS从300涨到1200,说明之前是线程池瓶颈,不是下游瓶颈。

7. 总结与建议

  • 不要用requests做并发IO,线程池是有限资源,协程是廉价资源
  • asyncio.gather要配合return_exceptions=True使用,否则单个任务异常会取消所有任务
  • 生产环境必须监控EventLoop延迟,用loop.time()打点,超过100ms就要查慢任务
  • 如果下游API不支持批量接口,考虑用asyncio.Queue做流量整形,避免突发请求打崩下游

这次重构收益远大于成本,代码只改了200行,但性能提升了8倍。如果你的服务也是IO密集且CPU闲置,建议立刻尝试asyncio迁移。记住:线程是给CPU密集型任务准备的,IO密集型请交给协程