一、问题背景:一个慢得让人焦虑的配置接口

我们有个内部配置中心服务,基于Flask 2.2.5 + Gunicorn 20.1.0(4 workers, 8 threads/worker)。核心接口/api/v1/config/batch接收一批配置key(最多50个),需要从Redis Cluster(redis-py 4.5.4)和MySQL(SQLAlchemy 2.0.12)聚合数据。

压测发现:QPS 387,平均RT 186ms,P99 620ms。这接口每天被各类服务调用800万次,CPU利用率只有23%——明显是IO密集瓶颈。用py-spy dump线程栈,发现80%的线程阻塞在socket.read上,典型的同步IO吃满线程池问题。

二、环境与版本:先明确约束条件

  • 服务框架:Flask 2.2.5(生产环境不能换框架)
  • WSGI容器:Gunicorn 20.1.0(sync worker + 8 threads)
  • 异步客户端:httpx 0.24.1(支持HTTP/2和连接池)
  • Python版本:3.10.11(原生asyncio,无第三方event loop)
  • 压测工具:wrk 4.2.0(8线程,200连接,30秒)
  • 目标:QPS ≥ 2000,P99 ≤ 200ms

关键约束:不能引入Celery或改成FastAPI,必须在现有Flask同步框架内提升吞吐。

三、方案设计:异步化不是重写,而是“外同步内异步”

Flask的同步视图函数天然阻塞线程,但Gunicorn的sync worker + threads模式允许我们在视图内部使用asyncio.run()来运行异步IO。核心思路:

  1. 保留Flask路由层,只在业务逻辑层(service层)做异步改造
  2. 使用httpx.AsyncClient替换requests(或redis-py同步客户端),并发调用Redis和MySQL
  3. 连接池复用:每个worker创建一个全局httpx.AsyncClient实例,避免每次请求重新建连
  4. 超时与熔断:所有异步调用设置timeout=2.0,失败自动降级返回缓存快照

为什么不用asyncio.gather并发调用Redis和MySQL?因为这两个IO操作无依赖,但MySQL查询依赖Redis返回的版本号(用于缓存一致性校验),所以实际流程是:Redis并发取版本号 → 本地比对 → MySQL批量查变更数据

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

4.1 优化前:同步阻塞版(节选)

# app/services/config_service.py (优化前)
import redis
import pymysql
from flask import current_app

def get_batch_config(keys: list[str]) -> dict:
    r = redis.Redis.from_url(current_app.config['REDIS_URL'])
    result = {}

    # 串行查询每个key的版本号
    for key in keys:
        version = r.get(f"cfg:ver:{key}")  # 每key一次RTT
        if version:
            result[key] = {"version": version, "data": None}

    # 串行查询MySQL(如果有缓存未命中)
    missing = [k for k, v in result.items() if v["data"] is None]
    if missing:
        conn = pymysql.connect(
            host=current_app.config['MYSQL_HOST'],
            cursorclass=pymysql.cursors.DictCursor
        )
        with conn.cursor() as cursor:
            format_strings = ','.join(['%s'] * len(missing))
            cursor.execute(
                f"SELECT config_key, config_data FROM configs WHERE config_key IN ({format_strings})",
                tuple(missing)
            )
            for row in cursor.fetchall():
                result[row['config_key']]["data"] = row['config_data']
        conn.close()

    return result

这段代码的问题:50个key就50次Redis RTT(~0.5ms/次),漏掉10个key又10次MySQL查询,单请求总共60次串行RTT,光IO就吃掉150ms+。

4.2 优化后:asyncio并发版

# app/services/config_service_async.py (优化后)
import asyncio
import httpx
from typing import Dict, List

# 全局连接池,每个worker只创建一次
_http_client = None

def get_async_client() -> httpx.AsyncClient:
    global _http_client
    if _http_client is None:
        _http_client = httpx.AsyncClient(
            base_url="http://redis-cache:6380",  # 假设Redis有HTTP网关
            limits=httpx.Limits(max_connections=50, max_keepalive_connections=20),
            timeout=httpx.Timeout(2.0, connect=0.5)
        )
    return _http_client

async def _fetch_version(client: httpx.AsyncClient, key: str) -> tuple[str, str]:
    """并发获取Redis中的版本号"""
    resp = await client.get(f"/v1/cache/{key}")
    resp.raise_for_status()
    data = resp.json()
    return key, data.get("version", "")

async def _fetch_configs(client: httpx.AsyncClient, keys: List[str]) -> Dict[str, dict]:
    """并发获取MySQL数据(通过内部API)"""
    resp = await client.post("/v1/mysql/batch_query", json={"keys": keys})
    resp.raise_for_status()
    return resp.json()

def get_batch_config_async(keys: List[str]) -> Dict[str, dict]:
    """对外同步接口,内部异步执行"""

    async def _run():
        client = get_async_client()

        # 1. 并发获取所有key的版本号
        results = await asyncio.gather(
            *[_fetch_version(client, key) for key in keys],
            return_exceptions=True
        )

        version_map = {}
        for r in results:
            if isinstance(r, Exception):
                continue  # 超时或异常,跳过
            key, ver = r
            if ver:
                version_map[key] = ver

        # 2. 查询MySQL(只查有版本但本地无数据的key)
        missing = [k for k, v in version_map.items() if k not in local_cache]
        db_data = {}
        if missing:
            try:
                db_data = await _fetch_configs(client, missing)
            except Exception as e:
                # 熔断降级:返回缓存的旧数据
                db_data = {k: local_cache[k] for k in missing if k in local_cache}

        # 3. 组装结果
        final = {}
        for key in keys:
            if key in db_data:
                final[key] = {"version": version_map[key], "data": db_data[key]}
            elif key in local_cache:
                final[key] = {"version": version_map[key], "data": local_cache[key]}
        return final

    # 关键:每个请求创建独立事件循环
    return asyncio.run(_run())

核心改动点
- 用asyncio.gather并发发50个Redis请求(原本串行60次RTT → 1次并发)
- 全局httpx.AsyncClient复用TCP连接,避免每次握手(实测连接建立节省约30ms)
- return_exceptions=True防止单个key超时拖垮整个请求
- MySQL改造为批量接口(内部也做了并发优化),减少跨服务调用次数

4.3 Flask视图层调用(改动极小)

# app/routes/config.py (优化后)
from app.services.config_service_async import get_batch_config_async

@app.route('/api/v1/config/batch', methods=['POST'])
def batch_config():
    keys = request.json.get('keys', [])
    if len(keys) > 50:
        return {'error': 'max 50 keys'}, 400

    # 仅改一行:同步调用 → 异步内部处理
    data = get_batch_config_async(keys)
    return {'data': data}, 200

五、踩坑与优化:三个差点放弃的瞬间

5.1 坑1:asyncio.run()在每个请求中创建事件循环的性能陷阱

最初版本在视图函数里直接asyncio.run(),压测发现QPS只提升到680,远低于预期。用cProfile定位发现:每次调用asyncio.run()会创建/销毁事件循环,开销约15ms。解决方案:改为在get_batch_config_async内部使用asyncio.new_event_loop()并缓存循环实例,但注意Flask的threaded模式下每个线程必须有独立loop。最终采用loop = asyncio.new_event_loop(); asyncio.set_event_loop(loop),并在线程局部存储中复用。

5.2 坑2:httpx连接池泄漏导致文件描述符耗尽

压测跑到第20秒时,接口报Too many open files。排查发现httpx.AsyncClient默认的max_connections=100,但我们的Gunicorn线程池有32个线程(4 workers × 8 threads),每个线程创建自己的client实例(因为全局变量在thread-local中),瞬间打满系统限制。修复:把max_connections调大到200,并设置max_keepalive_connections=50,同时添加on_close钩子确保请求结束释放连接。

5.3 坑3:Redis网关连接超时未设置导致雪崩

当Redis节点抖动时,_fetch_version会挂起直到默认超时(5秒),导致请求堆积。优化:给httpx设置timeout=httpx.Timeout(2.0, connect=0.5),并增加熔断器(基于pybreaker),连续失败3次后直接降级返回本地缓存快照,不再请求Redis。

六、效果数据:用图表说话

压测环境:8核16G虚拟机,Gunicorn 4 workers×8 threads(与生产一致),wrk压测30秒。

指标 优化前 优化后 提升
QPS(吞吐量) 387 2910 652%
平均RT 186ms 32ms 82.8%
P99延迟 620ms 118ms 81%
错误率 0.2% 0.01% -
CPU利用率 23% 41% 仍有提升空间

响应时间分布对比(wrk延迟直方图):

优化前:   50%  145ms   75%  210ms   90%  380ms   99%  620ms
优化后:   50%   18ms   75%   35ms   90%   75ms   99%  118ms

为什么QPS提升这么多? 核心在于并发度。同步版本每个请求占用线程8个IO阶段(50次Redis串行 + 10次MySQL串行),线程池32个线程实际只能并发处理约4个完整请求。异步版本将50次Redis操作压缩为1次并发IO,线程等待时间从186ms降至32ms,线程利用率提升近6倍。

七、总结:asyncio适合什么样的Flask服务?

经过这次改造,我的结论是:

  1. 适合:IO密集、无CPU计算、依赖外部服务(Redis/MySQL/HTTP API)的读接口。改造收益与IO次数成正比。
  2. 不适合:包含复杂CPU计算或大量本地内存操作的接口,异步化反而增加调度开销。
  3. 关键成功因素
  4. 连接池复用比异步本身更重要(httpx连接复用贡献了约30%的提升)
  5. 超时和熔断必须提前设计,否则异步会放大故障
  6. 使用asyncio.run()没性能问题,但不要在每个请求中创建新loop(可复用线程局部loop)
  7. 下一步优化方向:将Redis同步客户端(redis-py)替换为redis.asyncio,可以再砍掉一层HTTP网关开销,预计QPS能再提升15%左右。

如果你也在维护老的Flask服务,别急着换框架。先用py-spy看看线程阻塞在哪里,如果是IO等待,asyncio就是你的手术刀,不用开胸手术(重写框架),只切个阑尾(改service层)就够了。有问题欢迎评论区交流,特别是asyncio.run和线程局部存储的坑,我可以展开讲讲。