一、问题背景:一个慢接口拖垮了整个服务

上个月接手一个内部工单系统,用户反馈“导出报表经常超时”。查看监控发现/api/v1/report/export接口在高峰时段P99延迟高达1.2s,而该接口内部逻辑很简单:查数据库获取1000条工单ID,然后循环调用外部CRM系统的/batch/status接口查询状态,最后聚合返回。

问题出在循环调用上——用了requests.post同步请求,每次等待200-500ms,1000次串行就是200-500秒。后来加了线程池(ThreadPoolExecutor)勉强能跑,但线程切换开销大,且连接复用差。压测数据触目惊心:50并发下QPS仅120,线程数飙到300+,CPU闲置但IO等待严重。

二、环境与版本:别用最新,用最稳

  • Python 3.10.12(pyenv管理,3.11+的asyncio有性能提升但团队统一3.10)
  • aiohttp 3.9.1(pip install aiohttp==3.9.1,3.10版本有连接池bug)
  • uvloop 0.17.0(可选,但强烈推荐,事件循环性能提升约30%)
  • 压测工具:wrk 4.2.0(--threads=4 --connections=50 --duration=30s

关键配置aiohttpTCPConnector需要显式设置limit=100(默认100,但容易耗尽),并启用keepalive_timeout=30。数据库连接用asyncpg 0.27.0(比aiomysql更稳定)。

三、方案设计:从同步阻塞到异步非阻塞

核心思路:将requests.post替换为aiohttp.ClientSession.post,配合asyncio.gather并发发起请求。但直接无限制并发会打死CRM系统,所以用asyncio.Semaphore(20)限制最大并发数(根据CRM压测,20并发是安全阈值)。

架构对比

Before:  Flask(线程池) → requests.post(同步阻塞) → CRM
After:   Flask(单线程) → asyncio.run_coroutine_threadsafe → aiohttp(异步IO) → CRM

注意:Flask是WSGI同步框架,不能直接跑async def视图函数。我采用asyncio.run_coroutine_threadsafe将协程提交到独立的事件循环线程中运行,主线程保持响应。

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

Before(同步阻塞版)

import requests
from concurrent.futures import ThreadPoolExecutor
from flask import Flask, jsonify

app = Flask(__name__)
executor = ThreadPoolExecutor(max_workers=20)

@app.route('/api/v1/report/export', methods=['GET'])
def export_report():
    # 模拟查询数据库获取工单ID
    ticket_ids = list(range(1, 1001))

    def fetch_status(tid):
        resp = requests.post(
            'https://crm.internal/api/batch/status',
            json={"ticket_id": tid},
            timeout=2
        )
        return resp.json()['status']

    results = list(executor.map(fetch_status, ticket_ids))
    return jsonify({"count": len(results), "statuses": results})

After(异步非阻塞版)

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

app = Flask(__name__)

# 全局事件循环(在独立线程中运行)
loop = uvloop.new_event_loop()
asyncio.set_event_loop(loop)

# 全局aiohttp会话,复用连接池
connector = aiohttp.TCPConnector(limit=100, keepalive_timeout=30)
session = aiohttp.ClientSession(connector=connector, timeout=aiohttp.ClientTimeout(total=3))

async def fetch_status(sem, tid):
    async with sem:
        async with session.post(
            'https://crm.internal/api/batch/status',
            json={"ticket_id": tid}
        ) as resp:
            data = await resp.json()
            return data['status']

async def run_async(ticket_ids):
    sem = asyncio.Semaphore(20)  # 限制并发20
    tasks = [asyncio.create_task(fetch_status(sem, tid)) for tid in ticket_ids]
    return await asyncio.gather(*tasks)

@app.route('/api/v1/report/export', methods=['GET'])
def export_report():
    ticket_ids = list(range(1, 1001))
    # 通过run_coroutine_threadsafe将协程提交到事件循环线程
    future = asyncio.run_coroutine_threadsafe(run_async(ticket_ids), loop)
    results = future.result(timeout=10)  # 主线程阻塞等待(但事件循环不阻塞)
    return jsonify({"count": len(results), "statuses": results})

# 应用退出时关闭资源
@app.teardown_appcontext
def close_session(exception=None):
    if hasattr(session, 'close'):
        asyncio.run_coroutine_threadsafe(session.close(), loop).result()

关键点
1. asyncio.Semaphore(20)控制并发,防止打爆CRM
2. aiohttp.ClientSession全局复用,避免每次请求创建新连接(这是性能提升的核心)
3. uvloop替换默认事件循环,减少系统调用开销

五、踩坑与优化:两个血泪教训

坑1:事件循环策略导致RuntimeError
第一次直接在视图函数里写asyncio.run(),结果报错“This event loop is already running”。因为Flask的请求处理线程和事件循环线程不是同一个,必须用run_coroutine_threadsafe。后来在应用启动时创建全局loop,并设置为uvloop

坑2:连接池泄漏导致内存暴涨
最初没有复用ClientSession,每个请求async with aiohttp.ClientSession() as session,压测30分钟后内存从200MB涨到2GB。原因是ClientSession内部维护连接池,创建/销毁频繁导致文件描述符泄漏。改为全局单例后,内存稳定在350MB。

性能优化
- 将json={"ticket_id": tid}改为json={"ticket_ids": [tid]}批量请求,但CRM接口不支持,只能作罢
- 使用asyncio.gatherreturn_exceptions=True,避免单个请求失败导致全部取消
- 设置aiohttp.ClientTimeout(total=3),防止慢请求拖垮整体

六、效果数据:QPS提升15倍

使用wrk -t4 -c50 -d30s http://localhost:5000/api/v1/report/export压测结果:

指标 Before(同步) After(异步) 提升
QPS 120 1850 15.4x
平均延迟 420ms 54ms 7.8x
P99延迟 800ms 210ms 3.8x
线程数 300+ 6(主线程+事件循环) -
CPU占用 15% 45%(IO等待变计算) 利用更充分

压测环境:4核8G虚拟机,CRM服务部署在另一台机器(千兆内网)。延迟主要来自CRM端点(单次请求约50ms),异步并发20个请求时,总耗时约1000/20 * 50ms = 2.5s,但QPS提升是因为Flask主线程不再阻塞,可同时处理多个请求。

注意:QPS提升15倍的前提是CRM服务不成为瓶颈。实际生产中将Semaphore调至30后,CRM响应时间从50ms涨到120ms,说明已到极限。

七、总结:异步不是银弹,但IO密集场景是神器

这次重构让我深刻理解:同步代码适合CPU密集,异步适合IO密集。本场景中80%时间在等待网络IO,异步非阻塞将等待时间让渡给其他请求,资源利用率大幅提升。

后续建议
1. 如果接口需要数据库操作,配合asyncpg + SQLAlchemy 1.4+的异步扩展
2. 监控事件循环的loop.slow_callback_duration,设置0.1s警告
3. 生产环境部署时,用gunicorn + uvicorn workers(但Flask需迁移为QuartFastAPI才能完全异步化,当前方案是妥协)

最后想说:asyncio的坑不少,但一旦跑通,性能回报非常可观。如果你们也有类似同步循环调外部API的场景,值得一试。有问题欢迎评论区交流。