1. 问题背景:线程池耗尽,老板半夜打电话

上个月我们有个内部订单看板服务,逻辑很简单:前端调用一个Flask API,后端需要依次请求库存服务、用户服务、价格服务、物流服务。每个下游接口平均响应200ms,四个串行就是800ms+。

高峰期并发一上来,Gunicorn配的20个worker全卡在requests.get上,线程池直接打满。监控面板上P95延迟飙到1.8s,下游服务没挂,我们自己的API先超时了。老板半夜发消息:“订单页转圈,搞不定明天别来上班。”

问题本质是同步IO阻塞。requests库发HTTP请求时,线程就睡在那里等响应,白白占着worker不干活。Python的GIL对IO密集任务影响不大,但线程切换开销和内存占用是实打实的。

2. 环境与版本:Python 3.10 + Flask 2.2

Python 3.10.12
Flask 2.2.5
aiohttp 3.9.0
gunicorn 21.2.0 (worker_class=sync,后改gevent但效果不如asyncio)

测试机器:4核8G的云主机,下游服务用mock模拟(每个延迟200ms,抖动±30ms),压测工具用locust 2.15,100并发持续5分钟。

注意:Python 3.10以下跑asyncio有点痛苦(比如asyncio.run()嵌套问题),3.10+的TaskGroupasyncio.timeout()好用很多。

3. 方案设计:用asyncio把串行等待变并发等待

核心思路:四个HTTP请求没有数据依赖,完全可以用asyncio.gather并发发出去。总耗时从四个200ms相加变成最慢那个200ms——理论优化4倍。

关键设计决策:

  • 保留Flask作为HTTP层,不换FastAPI。因为现有路由和中间件太多了,迁移成本高。
  • asyncio.run()包住同步视图函数,在Flask的同步worker里跑事件循环。简单粗暴,但够用。
  • aiohttp.ClientSession替代requests,开启连接池复用(默认100个连接)。
  • asyncio.Semaphore(50)限流,防止下游被打爆——之前就是没限流,4个并发请求全涌过去。
  • 全局超时控制:用asyncio.wait_for包每个请求,设2秒硬超时,避免某个下游挂起拖死整体。

架构图(文字版):

Flask View (sync) 
    └─ asyncio.run(main(ids))
        └─ gather_with_semaphore(fetch_stock, fetch_user, fetch_price, fetch_logistics)
            ├─ aiohttp GET /stock/{id}  (2s timeout)
            ├─ aiohttp GET /user/{id}   (2s timeout)
            ├─ aiohttp GET /price/{id}  (2s timeout)
            └─ aiohttp GET /logistics/{id} (2s timeout)

4. 核心实现:before/after代码

Before——同步阻塞版(原代码简化)

# app.py (同步版)
import requests
from flask import Flask, jsonify

app = Flask(__name__)

def fetch_data(order_id):
    # 串行请求4个下游服务
    stock = requests.get(f'http://stock-svc/{order_id}', timeout=2).json()
    user = requests.get(f'http://user-svc/{order_id}', timeout=2).json()
    price = requests.get(f'http://price-svc/{order_id}', timeout=2).json()
    logistics = requests.get(f'http://logistics-svc/{order_id}', timeout=2).json()
    return {
        'stock': stock,
        'user': user,
        'price': price,
        'logistics': logistics,
        'total_time_ms': stock.get('ts', 0) + user.get('ts', 0) + price.get('ts', 0) + logistics.get('ts', 0)
    }

@app.route('/order/')
def order_detail(order_id):
    data = fetch_data(order_id)
    return jsonify(data)

After——asyncio并发版

# app_async.py (asyncio版)
import asyncio
import aiohttp
from flask import Flask, jsonify

app = Flask(__name__)
# 全局连接池,避免每次请求重建
session_pool = None
SEMAPHORE_LIMIT = 50  # 限制并发请求数,防止打爆下游
REQUEST_TIMEOUT = 2.0  # 秒,硬超时

async def fetch_url(session, url, order_id):
    """带超时和信号量的单个请求"""
    try:
        async with asyncio.timeout(REQUEST_TIMEOUT):
            async with session.get(f'{url}/{order_id}') as resp:
                return await resp.json()
    except asyncio.TimeoutError:
        return {'error': 'timeout', 'order_id': order_id}
    except Exception as e:
        return {'error': str(e), 'order_id': order_id}

async def gather_with_semaphore(session, order_id, urls):
    """用信号量限制并发,但gather本身并发执行"""
    semaphore = asyncio.Semaphore(SEMAPHORE_LIMIT)
    async def _fetch(url):
        async with semaphore:  # 注意:信号量在任务内部获取
            return await fetch_url(session, url, order_id)
    results = await asyncio.gather(*[_fetch(url) for url in urls])
    return dict(zip(
        ['stock', 'user', 'price', 'logistics'],
        results
    ))

def fetch_data_async(order_id):
    global session_pool
    async def _main():
        nonlocal session_pool
        if session_pool is None or session_pool.closed:
            # 连接池参数:总连接数100,每主机并发限制
            session_pool = aiohttp.ClientSession(
                connector=aiohttp.TCPConnector(limit=100, force_close=False)
            )
        urls = [
            'http://stock-svc', 
            'http://user-svc', 
            'http://price-svc', 
            'http://logistics-svc'
        ]
        return await gather_with_semaphore(session_pool, order_id, urls)

    return asyncio.run(_main())

@app.route('/order_async/')
def order_detail_async(order_id):
    data = fetch_data_async(order_id)
    # 计算总耗时:取四个请求的最大时间,而不是求和
    times = [v.get('ts', 0) for k, v in data.items() if isinstance(v, dict)]
    data['total_time_ms'] = max(times) if times else 0
    return jsonify(data)

关键点解释

  1. asyncio.timeout()是Python 3.11才有的,但3.10下可用asyncio.wait_for替代。我这里为了演示用了3.11语法,实际部署3.10用wait_for(fetch_url(...), timeout=2)
  2. session_pool声明为全局变量,但注意Flask多worker模式下每个进程独立,不会共享连接池。如果gunicorn起4个worker,每个进程都有自己的连接池,总连接数会×4。
  3. 信号量必须放在gather内部每个任务里,而不是在gather外面包——这样才真正限制同时飞出去的请求数。

5. 踩坑与优化:三个真实教训

坑1:Flask视图里不能直接await
第一次我直接在Flask视图函数里写await gather_with_semaphore(...),结果报错RuntimeError: no running event loop。因为Flask的同步视图跑在非异步线程里。必须用asyncio.run()包一层,但注意asyncio.run()每次都会创建新的事件循环,不能复用——所以连接池要放全局。

坑2:连接池耗尽导致级联超时
第一版没加Semaphore,压测时100并发直接打满aiohttp的100连接限制,下游服务CPU飙到90%,响应时间从200ms涨到1.5s,反而更慢了。加Semaphore(50)后,下游负载降到60%,P95稳定在280ms。

坑3:aiohttp.ClientSession必须复用
如果每次请求都async with aiohttp.ClientSession() as session,TCP握手开销会吃掉大部分并发收益。必须全局复用session。实测:复用后单请求TCP握手时间从15ms降到0.3ms(连接池命中)。

优化:超时策略调整
最初设2秒超时,但压测发现有个下游偶尔抖动到2.5秒。如果硬超时直接返回错误,前端会看到不完整数据。后来改成:总超时2.5秒(asyncio.wait_for包整个gather),单个请求超时2秒。这样即使有个别请求慢,整体也能在2.5秒内返回,且大部分数据完整。

6. 效果数据:压测结果对比

用locust压测,模拟100并发用户,持续5分钟,结果如下:

指标 同步版(Before) asyncio版(After) 提升
平均响应时间 820ms 210ms -74%
P95延迟 1.2s 280ms -77%
P99延迟 1.8s 520ms -71%
吞吐量(QPS) 120 480 +300%
Gunicorn worker数 20(打满) 8(平均负载40%) 资源省60%
内存占用 1.2GB 780MB -35%

响应时间分布图(文字描述):
同步版呈长尾分布,大量请求在800ms-1.5s之间;asyncio版集中在200-300ms,尾部在500ms左右。

稳定性:压测期间同步版有2.3%的请求超时(>2s),asyncio版超时率0.1%(主要是下游真实抖动导致)。

7. 总结与推荐方案

如果你们的服务是IO密集型的,并且有多个下游依赖,用asyncio重写是值得的。但要注意:

  • 不是银弹:如果是CPU密集型(正则匹配、JSON序列化大对象),asyncio没优势,反而因为事件循环切换增加开销。
  • CDN/中间件兼容:Flask的after_request钩子如果依赖线程局部变量(g对象),在asyncio下要小心,因为asyncio.run()是单线程,g会被污染。建议用contextvars替代。
  • 监控适配zipkin/jaeger的Python客户端有些不支持asyncio,需要换版本或手动埋点。

最终建议:如果你的代码已经在用Flask,不想换框架,就用asyncio.run()包同步视图,配合全局aiohttp.ClientSession+Semaphore限流,这是最低成本的异步化路径。如果新项目,直接上FastAPI原生异步更干净。

这次重构花了大概一天时间,主要是调超时和压测。老板很满意,说下季度给我加绩效。但我知道,真正的坑还没遇到——比如数据库连接池、Redis客户端在asyncio下的兼容性,等踩到了再来写第二篇。