一、问题背景:一个慢得离谱的聚合接口

上个月接手了一个内部订单系统的重构任务。其中有个/api/v1/product/detail接口,前端需要展示商品的基础信息、实时库存、价格区间和用户评价摘要。原实现是典型的同步串行调用:

# before.py - 同步版本(Flask 2.2.5, Gunicorn 20.1.0, 4 workers)
import requests
from flask import Flask, jsonify

app = Flask(__name__)

def fetch_stock(product_id):
    resp = requests.get(f'http://stock-svc:8001/stock/{product_id}', timeout=2)
    return resp.json()

def fetch_price(product_id):
    resp = requests.get(f'http://price-svc:8002/price/{product_id}', timeout=2)
    return resp.json()

def fetch_reviews(product_id):
    resp = requests.get(f'http://review-svc:8003/reviews/{product_id}', timeout=2)
    return resp.json()

@app.route('/api/v1/product/detail/')
def product_detail(product_id):
    stock = fetch_stock(product_id)
    price = fetch_price(product_id)
    reviews = fetch_reviews(product_id)
    return jsonify({"product_id": product_id, "stock": stock, "price": price, "reviews": reviews})

用wrk压测(8线程,200连接,30秒):QPS≈120,P99延迟≈450ms。生产环境每天调用量约200万次,平均耗时210ms,光这一个接口就占了整个网关30%的响应时间。

二、环境与版本:别用老掉牙的Python

  • 操作系统:Ubuntu 22.04 LTS
  • Python:3.10.12(必须3.8+,asyncio的run_in_executorSemaphore在3.8之后才稳定)
  • Web框架:Flask 2.2.5(同步WSGI,但我们可以用asyncio.run()在请求内部跑协程)
  • HTTP客户端:aiohttp 3.8.4(比httpx更轻量,连接池表现更好)
  • 部署:Gunicorn 20.1.0 + 4 workers(注意:worker数不能再多了,后面讲原因
  • 压测工具:wrk 8.1.2(参数:-t8 -c200 -d30s

核心思路:Flask的同步视图函数内部,用asyncio.run()启动一个事件循环,把所有阻塞IO换成协程并发执行。 因为Flask本身是同步的,每个worker一次只处理一个请求,内部并发不会破坏WSGI协议。

三、方案设计:信号量限流 + 协程并发

三个上游服务平均延迟约为:

  • 库存服务:70ms(偶发300ms尖峰)
  • 价格服务:85ms(最稳定)
  • 评价服务:120ms(数据量大,最慢)

串行总耗时≈275ms。如果三者并发,理论上最大耗时≈max(70, 85, 120)=120ms。但实际还要考虑连接建立、内核调度,预期能压到150ms以内。

方案细节:

  1. 使用asyncio.gather()并发发起3个请求。
  2. asyncio.Semaphore(10)控制并发度——防止上游服务被我们打爆,同时限制连接池大小。
  3. 为每个子任务单独设置超时(30秒总体超时,子任务超时2.5秒),避免单个服务拖垮整个接口。
  4. 兜底策略:如果某个服务失败,返回缓存或默认值,不要让整个接口500

四、核心实现:before/after代码对比

4.1 初始化aiohttp连接池和事件循环

注意:不要在每个请求里创建新的ClientSession,必须复用! 否则连接池失效,性能反而不如requests。

# after.py - 异步版本(核心代码)
import asyncio
import aiohttp
from flask import Flask, jsonify

app = Flask(__name__)

# 全局连接池:每worker最多100个连接,并发请求上限50
CONNECTOR = aiohttp.TCPConnector(
    limit=50,          # 每worker最大并发连接数
    limit_per_host=20, # 每个host最多20连接
    ttl_dns_cache=300, # DNS缓存5分钟
    enable_cleanup_closed=True  # 自动清理废弃连接
)
SESSION = aiohttp.ClientSession(
    connector=CONNECTOR,
    timeout=aiohttp.ClientTimeout(total=30, connect=2)
)

# 并发信号量:限制同时打向上游的请求数
SEMAPHORE = asyncio.Semaphore(10)

async def fetch_json(session, url, semaphore, timeout=2.5):
    """带信号量限流 + 独立超时的异步GET"""
    async with semaphore:
        try:
            async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout)) as resp:
                if resp.status == 200:
                    return await resp.json()
                else:
                    # 非200返回空dict,由调用方兜底
                    return {"error": f"HTTP {resp.status}"}
        except asyncio.TimeoutError:
            return {"error": "timeout"}
        except aiohttp.ClientError as e:
            return {"error": str(e)}

async def fetch_all(product_id):
    """并发获取三个上游数据"""
    base = "http://"
    tasks = [
        fetch_json(SESSION, f"{base}stock-svc:8001/stock/{product_id}"),
        fetch_json(SESSION, f"{base}price-svc:8002/price/{product_id}"),
        fetch_json(SESSION, f"{base}review-svc:8003/reviews/{product_id}")
    ]
    # gather返回顺序与tasks顺序一致,不是完成顺序
    return await asyncio.gather(*tasks, return_exceptions=False)

@app.route('/api/v1/product/detail/')
def product_detail(product_id):
    # Flask同步视图内部启动一次性事件循环
    stock, price, reviews = asyncio.run(fetch_all(product_id))
    return jsonify({
        "product_id": product_id,
        "stock": stock,
        "price": price,
        "reviews": reviews
    })

4.2 Gunicorn配置:worker数量与事件循环的恩怨

踩坑警告:如果你有4个Gunicorn worker,每个worker内部跑一个asyncio.run(),相当于每worker一个事件循环。如果把Gunicorn的worker数改成--worker-class=gthread --threads=8,每个线程里都会创建/销毁事件循环,导致aiohttp连接池内部状态错乱

我的最终配置:

# gunicorn.conf.py
workers = 4          # 进程数,不宜超过CPU核数
worker_class = 'sync' # 保持同步worker,内部用asyncio.run
timeout = 30         # 请求超时30秒
keepalive = 5        # 长连接5秒
max_requests = 1000  # 每worker处理1000个请求后重启(防内存泄漏)

五、踩坑与优化:三个真实的坑

坑1:事件循环与Signal冲突
第一次部署上线,发现重启Gunicorn时偶尔报RuntimeError: asyncio.run() cannot be called from a running event loop。排查了半天,发现是Gunicorn的--reload模式和asyncio.run()冲突。解决方案:生产环境禁用reload,本地调试用flask run --reload时不要跑异步代码。

坑2:连接池耗尽导致雪崩
压测时把SEMAPHORE设为50,CONNECTOR.limit设为100,结果上游服务出现大量ConnectionResetError。分析后发现aiohttp的limit_per_host默认是0(无限),我忘了设置,导致同一host连接数失控。最终调成:limit=50, limit_per_host=20, SEMAPHORE=10

坑3:超时设置分层
最初只设了ClientTimeout(total=30),结果某个上游服务卡死时,整个接口要等30秒才返回。后来改成全局total=30秒,子任务total=2.5秒,一旦某个服务超时,立刻返回错误数据,总耗时永远不超过3秒。

六、效果数据:用数字说话

压测环境:4核8G云主机,wrk 8线程200连接,30秒。

指标 同步版(before) 异步版(after) 提升倍数
平均响应时间 210ms 32ms 6.5x
P99响应时间 450ms 78ms 5.8x
QPS 120 850 7.1x
CPU使用率 68% 56% -12%(受益于IO等待减少)
内存占用 450MB 610MB(aiohttp连接池开销) +35%

注意内存涨了35%,这是aiohttp连接池和Semaphore队列的代价。如果你的服务器内存紧张,可以调低limit到20,QPS会降到600左右,但内存只涨15%。

七、总结:异步不是银弹,但这次值了

适合异步的场景:IO密集型,且上游服务延迟差异大(比如本例中评价服务最慢,串行时白白等待120ms)。
不适合的场景:CPU密集型计算(加密、压缩),或者上游服务本地延迟低于1ms(异步收益忽略不计)。

三个经验:
1. 不要相信“异步一定会更快”——你得有压测数据。这次优化后QPS提升7倍,但如果你的代码里没有IO等待,异步反而增加开销。
2. 连接池和信号量必须精心调参——limit_per_host和SEMAPHORE不设置,生产环境必炸。
3. 兜底策略比并发更重要——异步代码里每个子任务都要有超时和默认值,否则一个上游挂了,你的接口就跟着挂。

最后提醒:asyncio.run()在Flask同步视图里用没问题,但如果你用FastAPI或Sanic,直接定义async def视图函数,永远不要调用asyncio.run()——那样会创建第二个事件循环,必出bug。

如果这篇对你启发,欢迎评论区交流你的异步改造踩坑经历。