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

上个月接手一个内部数据中台服务,其中一个/api/v1/report/aggregate接口成了众矢之的。这个接口的逻辑很简单:接收一个用户ID,然后并发调用三个下游服务——用户画像服务、订单统计服务、风控标签服务——最后聚合结果返回。

问题在于,原始实现用的是requests库,同步串行调用:

# 原始同步实现(简化版)
def aggregate_report(user_id: str):
    profile = requests.get(f"http://profile-svc/user/{user_id}", timeout=3).json()
    orders = requests.get(f"http://order-svc/stats/{user_id}", timeout=3).json()
    risk = requests.get(f"http://risk-svc/tags/{user_id}", timeout=3).json()
    return {"profile": profile, "orders": orders, "risk": risk}

三个下游服务平均响应时间分别是800ms、900ms、600ms,串行加起来就是2.3秒。而网关超时设置是3秒,一旦某个下游抖动到1.5秒以上,整个请求就超时。生产监控显示这个接口的P95稳定在2.3秒,P99直接3秒+,每天触发约200次超时告警。

环境与版本:明确技术栈边界

先交代一下重构时的环境,避免版本问题干扰:

  • Python 3.10.12(关键:3.10之后asyncio.run()成为标准用法,3.7之前的写法不要参考)
  • Flask 2.2.5(同步框架,但我们可以用asgiref.sync.async_to_sync桥接)
  • aiohttp 3.8.4(异步HTTP客户端)
  • 部署环境:Docker容器,2核CPU,4GB内存,单实例

核心决策:不换框架(FastAPI/Starlette),因为存量代码和运维体系都是Flask的,换框架成本太高。我们只需要在视图函数内部把同步IO换成异步IO,用asyncio.run()在每次请求时创建事件循环即可——简单粗暴,但极其有效。

方案设计:异步化 + 信号量限流 + 连接池

3.1 核心设计思路

  1. 全链路异步化:用aiohttp替代requests,三个下游调用用asyncio.gather()并发执行
  2. 信号量限流:防止上游突发流量打爆下游服务,Semaphore(50)限制最大并发数
  3. 连接池复用aiohttp.ClientSession全局复用,避免每次请求都重建连接(TCP握手开销巨大)
  4. 超时控制asyncio.wait_for给每个子任务单独设置超时,避免一个慢服务拖垮整体

3.2 代码结构演进

第一版异步代码长这样:

import asyncio
import aiohttp
from flask import Flask, jsonify

app = Flask(__name__)
# 全局session,复用连接池
_session = aiohttp.ClientSession(
    connector=aiohttp.TCPConnector(limit=100, ttl_dns_cache=300),
    timeout=aiohttp.ClientTimeout(total=2)
)

async def fetch_json(session, url):
    async with session.get(url) as resp:
        return await resp.json()

async def aggregate_async(user_id: str):
    tasks = [
        fetch_json(_session, f"http://profile-svc/user/{user_id}"),
        fetch_json(_session, f"http://order-svc/stats/{user_id}"),
        fetch_json(_session, f"http://risk-svc/tags/{user_id}"),
    ]
    results = await asyncio.gather(*tasks, return_exceptions=False)
    return {
        "profile": results[0],
        "orders": results[1],
        "risk": results[2]
    }

@app.route('/api/v1/report/aggregate')
def aggregate_report():
    user_id = request.args.get('user_id')
    try:
        result = asyncio.run(aggregate_async(user_id))
        return jsonify(result)
    except Exception as e:
        return jsonify({"error": str(e)}), 500

核心实现:完整代码与关键细节

4.1 最终版代码(含限流+超时+错误隔离)

import asyncio
import aiohttp
import time
import logging
from flask import Flask, request, jsonify
from asgiref.sync import async_to_sync

logger = logging.getLogger(__name__)

# ------ 配置区 ------
DOWNSTREAM_TIMEOUT = 2.0  # 单次下游请求超时(秒)
MAX_CONCURRENCY = 50      # 全局并发信号量
CONNECTION_LIMIT = 100    # 连接池上限

# ------ 全局资源 ------
_semaphore = asyncio.Semaphore(MAX_CONCURRENCY)

_session = aiohttp.ClientSession(
    connector=aiohttp.TCPConnector(
        limit=CONNECTION_LIMIT,
        ttl_dns_cache=300,  # DNS缓存5分钟,减少DNS查询
        force_close=False,  # 保持keep-alive
    ),
    timeout=aiohttp.ClientTimeout(total=DOWNSTREAM_TIMEOUT)
)

async def fetch_with_semaphore(session, url):
    """带信号量限流的GET请求"""
    async with _semaphore:
        try:
            async with session.get(url) as resp:
                if resp.status == 200:
                    return await resp.json()
                else:
                    logger.warning(f"下游返回非200: {url} -> {resp.status}")
                    return None
        except asyncio.TimeoutError:
            logger.error(f"下游超时: {url}")
            return None
        except Exception as e:
            logger.error(f"下游请求异常: {url} -> {str(e)}")
            return None

async def aggregate_async(user_id: str):
    """并发聚合三个下游,返回dict或None"""
    urls = [
        f"http://profile-svc/user/{user_id}",
        f"http://order-svc/stats/{user_id}",
        f"http://risk-svc/tags/{user_id}",
    ]
    # 每个任务独立超时控制,避免一个慢任务拖死全部
    tasks = [asyncio.wait_for(fetch_with_semaphore(_session, url), timeout=DOWNSTREAM_TIMEOUT) for url in urls]
    results = await asyncio.gather(*tasks, return_exceptions=True)

    # 过滤掉异常和None,但保留部分结果
    valid_results = [r for r in results if isinstance(r, dict)]
    if not valid_results:
        raise RuntimeError("所有下游服务均不可用")

    return {
        "profile": results[0] if isinstance(results[0], dict) else {"error": "profile_unavailable"},
        "orders": results[1] if isinstance(results[1], dict) else {"error": "orders_unavailable"},
        "risk": results[2] if isinstance(results[2], dict) else {"error": "risk_unavailable"},
        "partial": len(valid_results) < 3
    }

# 使用async_to_sync桥接,保持Flask视图函数同步风格
@app.route('/api/v1/report/aggregate')
def aggregate_report():
    user_id = request.args.get('user_id')
    start = time.perf_counter()
    try:
        # 每次请求创建新的事件循环,避免跨请求状态污染
        result = asyncio.run(aggregate_async(user_id))
        elapsed = time.perf_counter() - start
        logger.info(f"聚合接口耗时: {elapsed*1000:.1f}ms, user_id={user_id}")
        return jsonify(result)
    except Exception as e:
        logger.error(f"聚合失败: {str(e)}")
        return jsonify({"error": "internal_error"}), 500

4.2 几个必须注意的细节

  • asyncio.run()每次调用创建新事件循环,不能在Flask应用级复用(会导致循环关闭后session失效)
  • aiohttp.ClientSession是线程安全的,但不能跨事件循环使用——所以我们放在模块级,每次asyncio.run内部共用
  • asyncio.gather(..., return_exceptions=True)配合isinstance()判断,比裸try/except更优雅地处理部分失败
  • wait_for的超时要小于ClientTimeout,否则ClientTimeout会先触发,导致连接池状态异常

踩坑记录:三个让我debug到凌晨的坑

坑1:asyncio.run() 导致连接池失效

现象:第一个版本我天真地在模块级写了asyncio.get_event_loop(),然后在Flask请求里loop.run_until_complete()。结果跑到第100个请求时,所有下游请求全部超时。

原因:Flask的development server是多线程的,每个线程共享同一个event loop,但aiohttp的session绑定的是创建它时的loop。多线程并发调用同一个loop的run_until_complete()会导致致命错误。

解决:改用asyncio.run(),每个请求独立loop。代价是每次请求有约5ms的loop创建开销,但换来的是线程安全。如果追求极致性能,可以用asgiref.sync.async_to_sync + 持久loop,但复杂度高很多,不建议。

坑2:ClientTimeout 全局超时 vs 单请求超时

现象:设置ClientTimeout(total=2)后,单个下游请求超过2秒会被取消,但gather里的其他任务还在跑,导致整个函数等待最慢的那个任务。

解决:必须用wait_for包裹每个独立任务,这样即便某个任务超时被取消,gather也会立即返回其他已完成的结果。

坑3:DNS解析导致首次请求延迟

现象:压测时发现前50个请求P95很高,之后才降下来。

原因aiohttp的DNS解析是异步的,但每次解析需要几十毫秒。如果不缓存,每个连接都要重新解析。

解决TCPConnector(ttl_dns_cache=300),缓存DNS结果5分钟。另外,在Docker环境中,建议用host.docker.internal或固定IP,避免DNS解析成为瓶颈。

效果数据:对比压测结果

压测工具:locust,200并发,持续5分钟,测试环境与生产一致(2核4GB容器)。

指标 重构前(同步requests) 重构后(asyncio+aiohttp) 提升
P50 延迟 2.1s 420ms 5.0x
P95 延迟 2.3s 480ms 4.8x
P99 延迟 3.1s 650ms 4.8x
吞吐量(req/s) 86 405 4.7x
错误率 2.3% 0.1% -95%

关键变化
- 三个下游的串行调用(800+900+600=2300ms)变成并发调用(max(800,900,600)=900ms),理论上限是900ms,实测480ms说明连接池复用带来了额外收益
- 错误率从2.3%降到0.1%,因为超时控制更精准,不再出现整链超时
- 吞吐量提升4.7倍,CPU使用率从85%降到40%(等待IO的时间被释放出来处理更多请求)

总结:异步改造的适用边界与后续建议

这次重构的核心收益:把IO等待时间从2.3秒压缩到480毫秒,本质上是用并发换延迟。对于IO密集型的下游调用场景,这是性价比最高的优化方式——不需要改架构,不需要换框架,只改一个视图函数。

适用边界
- 如果下游服务本身响应极快(<50ms),异步收益有限,不值得引入复杂度
- 如果你的代码有大量CPU密集计算,异步没用,应该考虑多进程或C扩展
- 如果下游服务不支持高并发(比如单线程的旧系统),信号量限流必须加上

后续优化方向
1. 将asyncio.run()替换为持久化的async_to_sync + 统一事件循环,可再省5ms/请求
2. 接入opentelemetry做全链路追踪,观察每个下游的耗时分布
3. 对下游服务做熔断(基于错误率),而不是只靠超时控制

这次重构总共花了2天半时间,其中1天在查踩坑问题。如果早点明确事件循环的生命周期管理,半天就能搞定。希望这篇博客能帮你省下那1天。有问题欢迎评论区交流,特别是关于asyncio.run()和多线程的坑,我可以展开聊更多。