串行调用多个外部服务时,接口总耗时通常接近各调用耗时之和:只要其中一个依赖变慢,整体响应就会被拖长;如果某个请求卡住,后续调用甚至无法开始。直接把调用改成并行,能够减少等待叠加,但也可能同时放大连接数、线程数、下游请求量和错误数量。可靠的改造目标不是“尽可能并发”,而是在明确的并发上限、超时边界和部分成功语义下缩短整体等待时间。
先区分三种方案
假设接口需要调用服务 A、B、C,三者彼此没有数据依赖,耗时分别为 300、500、700 毫秒。
| 方案 | 理论等待时间 | 主要收益 | 主要风险 | 适用条件 |
|---|---|---|---|---|
| 保持串行 | 约 1500 毫秒 | 逻辑简单,资源压力可控 | 任一慢调用都会叠加到总耗时 | 调用存在严格依赖,或请求量很低 |
| 直接并行 | 约 700 毫秒 | 改造简单,延迟下降明显 | 瞬时请求量、连接数和失败量同时放大 | 调用数量固定且很少,下游容量明确 |
| 受控并行 | 接近最长调用耗时加调度开销 | 兼顾性能和资源控制 | 需要处理排队、超时、取消和结果状态 | 生产环境中的常规选择 |
这里的“并行”通常是异步 I/O 并发,而不是为每个请求无限创建线程。对网络等待占主导的调用,异步客户端可以让一个服务进程在等待响应时处理其他任务;但异步并不意味着下游可以承受无限请求。
改造前先确认调用关系
任务拆分应先回答三个问题:
- 调用是否真的互相独立:如果 B 必须使用 A 的结果,就不能简单地与 A 同时启动。
- 失败后是否仍有业务价值:例如推荐信息失败不影响主数据返回,而权限校验失败可能必须终止整个请求。
- 结果是否需要完整:有些接口要求全部依赖成功,有些接口允许返回部分结果并标记缺失项。
可以把调用分成三类:
- 关键任务:失败后无法形成合法响应,应快速失败。
- 可降级任务:失败后使用默认值、缓存或空结果。
- 附加任务:不影响主响应,可以取消、延后或转为异步处理。
只有第一类任务存在强依赖时,才应保持串行等待。其余独立任务可以并行执行,并在聚合阶段统一判断结果。
为什么不应直接使用无限制并行
最直接的实现往往是为所有外部调用创建任务并一次性等待:
results = await asyncio.gather(
call_service_a(),
call_service_b(),
call_service_c(),
)
这种写法的问题不在于语法错误,而在于缺少运行边界:
- 调用数量增加时,会同时建立大量连接或占用连接池。
- 下游变慢时,更多请求会处于等待状态,进一步消耗本地资源。
- 一个未设置超时的任务可能让整个聚合请求长期不结束。
gather的异常处理方式可能使调用方只看到一个错误,难以知道哪些任务已经成功、哪些任务尚未完成。- 如果上游请求已经断开,后台任务仍可能继续访问下游,形成无效流量。
因此,并行化应至少配套四个控制点:并发上限、单项超时、整体截止时间和任务取消。
最小的受控并行实现
下面的示例使用 Python asyncio 模拟服务调用,包含并发上限、单项超时、整体超时、部分失败、取消未完成任务和结果验证。真实项目中只需要将 call_external_service 替换成实际异步 HTTP、RPC 或 SDK 调用。
import asyncio
import random
import time
from typing import Any
async def call_external_service(name: str) -> dict[str, Any]:
"""模拟外部服务调用。真实实现应使用支持异步的客户端。"""
await asyncio.sleep(random.uniform(0.1, 1.2))
if name == "profile":
return {"user_id": 42, "name": "demo"}
if name == "quota":
return {"remaining": 100}
if name == "recommendation":
return {"items": ["a", "b"]}
raise RuntimeError(f"unknown service: {name}")
def validate_result(name: str, value: Any) -> bool:
"""只验证聚合逻辑真正依赖的字段。"""
if not isinstance(value, dict):
return False
required_fields = {
"profile": {"user_id", "name"},
"quota": {"remaining"},
"recommendation": {"items"},
}
return required_fields[name].issubset(value)
async def run_one(
name: str,
semaphore: asyncio.Semaphore,
per_call_timeout: float,
) -> dict[str, Any]:
start = time.monotonic()
try:
# 获取信号量的等待时间也受单项超时约束。
async with semaphore:
value = await asyncio.wait_for(
call_external_service(name),
timeout=per_call_timeout,
)
if not validate_result(name, value):
return {
"name": name,
"status": "invalid",
"value": None,
"elapsed_ms": elapsed_ms(start),
}
return {
"name": name,
"status": "success",
"value": value,
"elapsed_ms": elapsed_ms(start),
}
except asyncio.TimeoutError:
return {
"name": name,
"status": "timeout",
"value": None,
"elapsed_ms": elapsed_ms(start),
}
except asyncio.CancelledError:
# 取消异常应继续向上传播,不能被普通异常兜底吞掉。
raise
except Exception as exc:
return {
"name": name,
"status": "failed",
"error": type(exc).__name__,
"value": None,
"elapsed_ms": elapsed_ms(start),
}
def elapsed_ms(start: float) -> int:
return round((time.monotonic() - start) * 1000)
async def aggregate(
names: list[str],
max_concurrency: int = 2,
per_call_timeout: float = 0.8,
overall_timeout: float = 1.0,
) -> dict[str, Any]:
semaphore = asyncio.Semaphore(max_concurrency)
tasks = {
name: asyncio.create_task(
run_one(name, semaphore, per_call_timeout)
)
for name in names
}
done, pending = await asyncio.wait(
tasks.values(),
timeout=overall_timeout,
)
# 整体截止时间到达后,主动取消仍未完成的任务。
for task in pending:
task.cancel()
if pending:
await asyncio.gather(*pending, return_exceptions=True)
results: dict[str, Any] = {}
for name, task in tasks.items():
if task in pending:
results[name] = {
"name": name,
"status": "cancelled",
"value": None,
}
else:
results[name] = task.result()
# 关键任务失败时,返回整体失败;附加任务可以继续按降级策略处理。
critical_names = {"profile"}
critical_failed = any(
results[name]["status"] != "success"
for name in critical_names
if name in results
)
return {
"status": "failed" if critical_failed else "partial_or_success",
"results": results,
}
async def main() -> None:
response = await aggregate(
["profile", "quota", "recommendation"],
max_concurrency=2,
per_call_timeout=0.8,
overall_timeout=1.0,
)
print(response)
if __name__ == "__main__":
asyncio.run(main())
这个实现有几个边界需要明确:
max_concurrency=2限制的是当前聚合请求内的并发数,不等同于整个服务实例的全局并发控制。per_call_timeout限制单项任务,包括等待并发槽位和执行外部调用的时间。overall_timeout是聚合层的最终截止时间,避免多个单项超时叠加后仍长时间占用请求。validate_result应验证业务必需字段,而不是只判断 HTTP 状态是否成功。CancelledError不应被普通Exception逻辑转换成“失败”,否则上游取消后,后台任务可能继续运行。- 取消协程只表示本地任务不再等待;是否能真正停止网络请求,还取决于底层客户端是否支持取消,以及连接是否已将请求发送到下游。对不可取消的调用,仍需依靠客户端超时和连接回收止损。
部分失败不能只返回一个布尔值
聚合结果至少应区分以下状态:
success:调用完成,且结果通过校验。timeout:在单项时间预算内没有完成。failed:连接错误、协议错误或业务异常。invalid:调用返回了数据,但缺少必需字段或格式不符合预期。cancelled:由于整体截止时间或上游请求取消而未完成。
不要把超时、失败和空结果都统一转换成默认值。这样虽然能让接口“返回成功”,但会掩盖下游故障,导致调用方无法区分“确实没有数据”和“数据获取失败”。
更稳妥的响应结构是将业务数据和状态分开:
{
"status": "partial_or_success",
"data": {
"profile": {
"user_id": 42,
"name": "demo"
},
"quota": null,
"recommendation": {
"items": ["a", "b"]
}
},
"degraded": ["quota"],
"errors": {
"quota": {
"status": "timeout"
}
}
}
如果对外协议暂时不能增加状态字段,也应在服务内部保留完整状态,并在日志和指标中记录。否则只能看到一个最终响应,无法判断部分失败的来源。
并发上限如何设定
并发上限不是越大越好,可以从以下约束中取较小值:
- 下游允许的并发量:包括对方服务的配额、连接数和处理能力。
- 本地连接池容量:并发任务超过连接池后,任务可能只是排队等待。
- 单实例请求并发量:一个接口请求的并发上限乘以实例同时处理的请求数,才是实际下游压力。
- 超时预算和请求规模:如果单项任务耗时较长,过小的并发上限会导致排队时间吞噬收益。
- 失败时的放大效应:下游故障期间,过大的并发会制造更多超时和重试压力。
例如,一个实例同时处理 50 个聚合请求,每个请求最多启动 4 个下游调用,理论上可能产生 200 个下游请求。即使单个聚合请求的并发上限看起来不大,也需要结合实例级容量评估。
初始值可以通过小范围压测获得,而不是凭经验直接设置。先选择保守并发,再观察延迟、连接池等待和下游错误率,逐步增加,直到收益趋于平缓或资源指标出现拐点。
超时应形成层次,而不是只设一个数字
一个完整的时间预算通常至少包括:
- 连接超时:建立连接允许等待多久。
- 读取超时:连接建立后等待响应数据多久。
- 单项调用超时:一个外部任务最多占用多久。
- 整体聚合超时:本次接口最多等待多久。
- 上游请求截止时间:不能超过调用方或网关剩余时间。
整体超时必须小于上游真正能接受的总时间,并为序列化、日志、网络返回和重试预留余量。例如网关给接口 2 秒,服务不应把聚合整体超时设置成 2 秒,通常还需要留出返回响应的时间。
超时值也不应只看平均延迟。应结合压测中的高分位延迟、业务可接受等待时间和下游故障时的表现制定。单项超时过长,会让资源长时间被占用;过短,则可能把正常的长尾请求全部误判为失败。
取消策略要避免留下后台任务
当整体截止时间到达后,应取消未完成任务并等待它们结束清理。仅调用 task.cancel() 而不进一步 gather,可能留下未处理的取消异常或资源清理过程。
还要区分两种场景:
- 上游请求主动断开:聚合任务应监听请求取消,停止不再需要的下游调用。
- 单项任务超时:只取消该任务,不应无条件取消已经成功的其他任务。
- 关键任务失败:如果业务不允许继续,可以取消所有附加任务,避免继续产生无效流量。
- 附加任务失败:记录状态并降级,不应影响关键数据返回。
如果底层 SDK 是同步阻塞调用,把它直接放进异步函数并不能获得真正的异步并发。此时应使用有界线程池,并设置线程池大小、任务超时和关闭策略;如果 SDK 支持原生异步客户端,优先使用原生异步实现。
日志和指标要能还原一次请求
仅记录“接口耗时 1200 毫秒”无法定位问题。建议为每次聚合请求生成关联标识,并为每个子任务记录:
- 聚合请求标识、子任务名称和目标服务。
- 任务开始时间、结束时间和等待并发槽位的时间。
- 实际耗时、状态和异常类型。
- 是否发生超时、取消或结果校验失败。
- 当前并发上限、排队数量和连接池等待情况。
- 聚合结果是完整成功、部分成功还是整体失败。
指标可以按服务名称和状态拆分:
- 聚合接口请求总数、成功率和整体延迟。
- 各下游调用的调用次数、成功率、超时率和失败率。
- 各下游调用的延迟分布,而不只是平均值。
- 并发槽位等待时间和活动任务数。
- 部分成功比例、结果校验失败比例和取消任务数。
- 连接池耗尽、线程池队列长度等资源指标。
日志中应避免记录完整响应体和敏感参数。对于相同异常,可以使用异常类型、下游服务名和错误码聚合,防止高并发故障时日志量反过来拖慢服务。
压测要比较的不只是平均耗时
改造前后至少准备三组场景:
- 正常场景:所有依赖按典型延迟返回。
- 单项变慢:让一个依赖延迟明显升高,观察其他任务是否仍能完成。
- 单项失败或不返回:验证单项超时、部分成功和整体截止时间。
- 高并发场景:增加同时到达的聚合请求,观察下游请求量和本地资源。
- 取消场景:客户端提前断开,确认下游任务不会长期残留。
每组场景都应记录:
- 总响应时间的中位数和高分位值。
- 每个下游服务的耗时与状态分布。
- CPU、内存、连接数、连接池使用率和事件循环延迟。
- 下游服务收到的请求量、错误率和限流情况。
- 超时后仍在执行的任务数量。
压测结果应同时比较串行、直接并行和受控并行。例如直接并行可能在低负载下延迟最低,但在高并发下连接池迅速耗尽;受控并行的单请求延迟可能略高,却能保持更稳定的高分位表现。最终选择应以稳定性和资源边界为准,而不是只看一次测试中的最快响应。
什么时候应该保持串行
以下情况不适合为了降低延迟强行并行:
- 后一个调用必须使用前一个调用的结果。
- 下游明确要求严格顺序。
- 调用具有不可逆副作用,重复或提前执行会产生业务风险。
- 下游容量不足,并行会显著增加拒绝率。
- 返回协议要求所有结果同时有效,且没有合理的部分失败语义。
- 经过压测后,并行节省的等待时间不足以抵消资源和维护成本。
串行并不等于低质量实现。只要设置合理的单项超时、整体截止时间和清晰的失败处理,串行方案仍可能是更容易验证和维护的选择。只有在调用确实独立、响应时间受等待叠加影响,并且能够定义部分结果语义时,受控异步并行才值得引入。
发布者:jacky,转转请注明出处:https://kubiyun.com/archives/4581
评论列表(8条)
并发上限最容易被忽略,单请求看着不大,实例一叠加就完全是另一回事。
取消任务后还得等清理,这个细节挺容易漏掉。
HTTP 200不等于结果真的能用,字段校验很关键
@信号迷宫:没错,HTTP 200只能说明请求到达并返回了,字段完整性和业务状态还得单独校验,否则很容易把“成功响应”当成可用结果。
整体超时别顶到网关截止线,得留点返回余量
@电子蜃楼:对,整体超时最好早于网关截止线,给结果聚合、日志和响应返回留出余量,避免刚拿到结果就被网关切断。
把同步SDK塞进异步里,阻塞可不会消失
@高傲孔雀君:说得对,异步只对真正的异步 I/O 有效;同步 SDK 需要放到线程池,或换成异步客户端,否则事件循环还是会被卡住。