当业务规模增长,单次API调用已不能满足吞吐量需求时,并发控制与资源调度就成为系统架构的核心问题。本节将深入探讨多请求并发处理、速率限制应对、队列管理、负载均衡和资源池化等关键技术,帮助你构建高可用、低成本的大模型API调用系统。
大模型API的典型响应时间在1-10秒之间。如果采用串行调用,系统吞吐量极低。通过合理的并发控制,可以在不超出API速率限制的前提下,将吞吐量提升数倍乃至数十倍。
串行 vs 并发调用对比 串行模式(5个请求): 请求1 ████████████░░░░░░░░░░░░░░░░░░░░ 2s 请求2 ░░░░░░░░████████████░░░░░░░░░░░░ 4s 请求3 ░░░░░░░░░░░░░░░░████████████░░░ 6s 请求4 ░░░░░░░░░░░░░░░░░░░░░██████████ 8s 请求5 ░░░░░░░░░░░░░░░░░░░░░░░░░████████ 10s 总耗时:10s 并发模式(5个请求,并发度=3): 请求1 ████████████░░░░░░░░░░░░░░░░░░░░ 2s 请求2 █████████████░░░░░░░░░░░░░░░░░░ 2.2s 请求3 ██████████████░░░░░░░░░░░░░░░░░ 2.5s 请求4 ░░░░░░░░░░░░░░████████████░░░░ 4.7s 请求5 ░░░░░░░░░░░░░░░░████████████░░ 5.0s 总耗时:5.0s(吞吐提升2倍)
大多数大模型API供应商都实施了速率限制,通常包括:
经典的速率限制算法是令牌桶(Token Bucket),它允许一定程度的突发流量,同时保证长期平均速率不超过限制。
在大量API请求场景中,直接并发往往会导致速率超限触发429错误。通过消息队列进行流量整形,可以将不均匀的请求流量平滑化,确保API调用始终在限制范围内。
常见队列方案:
当单一API供应商的速率限制无法满足需求时,可以通过多供应商负载均衡来扩展容量。核心策略包括:
import time import threading from collections import deque class TokenBucketRateLimiter: """令牌桶速率限制器,支持TPM和RPM双重限制""" def __init__(self, rpm: int, tpm: int): self.rpm = rpm self.tpm = tpm self.rpm_tokens = rpm self.tpm_tokens = tpm self.last_refill = time.time() self.lock = threading.Lock() self.request_times = deque() # 记录请求时间用于精确RPM def _refill(self): """补充令牌""" now = time.time() elapsed = now - self.last_refill # RPM补充 self.rpm_tokens = min(self.rpm, self.rpm_tokens + elapsed * self.rpm / 60) # TPM补充 self.tpm_tokens = min(self.tpm, self.tpm_tokens + elapsed * self.tpm / 60) self.last_refill = now def acquire(self, estimated_tokens: int = 0) -> float: """ 获取令牌,返回需要等待的时间(秒) 返回0表示无需等待 """ with self.lock: self._refill() wait_time = 0 # 检查RPM if self.rpm_tokens < 1: wait_time = max(wait_time, 60 / self.rpm) # 检查TPM if estimated_tokens > 0 and self.tpm_tokens < estimated_tokens: wait_time = max( wait_time, estimated_tokens * 60 / self.tpm ) if wait_time > 0: time.sleep(wait_time) self._refill() self.rpm_tokens -= 1 self.tpm_tokens -= estimated_tokens self.request_times.append(time.time()) return wait_time class AdaptiveRateLimiter: """自适应限速器:根据429错误动态调整速率""" def __init__(self, base_rpm: int, base_tpm: int): self.base_rpm = base_rpm self.base_tpm = base_tpm self.current_rpm = base_rpm self.current_tpm = base_tpm self.consecutive_429 = 0 self.last_success_time = time.time() self.limiter = TokenBucketRateLimiter(base_rpm, base_tpm) def on_success(self): """成功调用后恢复速率""" self.consecutive_429 = 0 self.last_success_time = time.time() # 逐步恢复到基准速率(每次恢复20%) if self.current_rpm < self.base_rpm: recovery = int((self.base_rpm - self.current_rpm) * 0.2) self.current_rpm = min(self.base_rpm, self.current_rpm + max(1, recovery)) self.current_tpm = min(self.base_tpm, int(self.current_tpm * 1.2)) self.limiter = TokenBucketRateLimiter(self.current_rpm, self.current_tpm) def on_rate_limit(self, retry_after: int = None): """遇到429后降低速率""" self.consecutive_429 += 1 # 指数退避式降速 factor = 0.5 ** self.consecutive_429 self.current_rpm = max(1, int(self.base_rpm * factor)) self.current_tpm = max(100, int(self.base_tpm * factor)) self.limiter = TokenBucketRateLimiter(self.current_rpm, self.current_tpm)
import asyncio import aiohttp from typing import Dict, Any, Optional import json class QueueBasedDispatcher: """基于队列的API调度器""" def __init__( self, api_config: Dict[str, str], rate_limiter: AdaptiveRateLimiter, max_concurrent: int = 10, max_retries: int = 3 ): self.api_config = api_config self.rate_limiter = rate_limiter self.max_concurrent = max_concurrent self.max_retries = max_retries self.semaphore = asyncio.Semaphore(max_concurrent) self.stats = { "total_sent": 0, "total_completed": 0, "total_failed": 0, "total_retried": 0, "total_saved_by_limit": 0 } async def _call_api(self, messages: list, model: str) -> Dict: """单次API调用(带限速与重试)""" estimated_tokens = sum(len(m["content"]) for m in messages) * 2 self.rate_limiter.limiter.acquire(estimated_tokens) headers = { "Authorization": f"Bearer {self.api_config['api_key']}", "Content-Type": "application/json" } payload = { "model": model, "messages": messages, "temperature": 0.7 } async with aiohttp.ClientSession() as session: for attempt in range(self.max_retries): try: async with session.post( self.api_config["endpoint"], headers=headers, json=payload, timeout=aiohttp.ClientTimeout(total=30) ) as resp: if resp.status == 200: data = await resp.json() self.rate_limiter.on_success() return { "success": True, "content": data["choices"][0]["message"]["content"], "usage": data.get("usage", {}) } elif resp.status == 429: retry_after = int(resp.headers.get("Retry-After", 5)) self.rate_limiter.on_rate_limit(retry_after) self.stats["total_retried"] += 1 await asyncio.sleep(retry_after) else: error_text = await resp.text() return { "success": False, "error": f"HTTP {resp.status}: {error_text}" } except asyncio.TimeoutError: if attempt == self.max_retries - 1: return {"success": False, "error": "请求超时"} await asyncio.sleep(2 ** attempt) async def submit(self, messages: list, model: str = "default") -> Dict: """提交请求到调度队列""" async with self.semaphore: self.stats["total_sent"] += 1 result = await self._call_api(messages, model) if result["success"]: self.stats["total_completed"] += 1 else: self.stats["total_failed"] += 1 return result async def submit_batch( self, batch: list, model: str = "default" ) -> list: """批量提交请求""" tasks = [self.submit(item["messages"], model) for item in batch] return await asyncio.gather(*tasks)
import random from dataclasses import dataclass from typing import List @dataclass class ProviderConfig: name: str endpoint: str api_key: str rpm_limit: int tpm_limit: int cost_per_1k_tokens: float avg_latency_ms: float quality_score: float # 0-1 class MultiProviderLoadBalancer: """多供应商负载均衡器""" def __init__(self, providers: List[ProviderConfig]): self.providers = providers self.rate_limiters = { p.name: AdaptiveRateLimiter(p.rpm_limit, p.tpm_limit) for p in providers } self.provider_stats = {p.name: { "success": 0, "failure": 0, "total_latency": 0 } for p in providers} def select_provider( self, strategy: str = "cost", priority: str = "quality" ) -> ProviderConfig: """根据策略选择供应商""" if strategy == "cost": # 成本优先:选最便宜的可用供应商 available = [ p for p in self.providers if self.rate_limiters[p.name].current_rpm > 0 ] return min(available, key=lambda p: p.cost_per_1k_tokens) elif strategy == "speed": # 速度优先:选延迟最低的可用供应商 available = [ p for p in self.providers if self.rate_limiters[p.name].current_rpm > 0 ] return min(available, key=lambda p: p.avg_latency_ms) elif strategy == "quality": # 质量优先:选质量评分最高的 return max(self.providers, key=lambda p: p.quality_score) elif strategy == "weighted": # 加权随机:综合成本、速度、质量 weights = [] for p in self.providers: score = ( 0.3 * (1 - p.cost_per_1k_tokens / 0.1) + 0.3 * (1 - p.avg_latency_ms / 10000) + 0.4 * p.quality_score ) weights.append(max(0.1, score)) total = sum(weights) weights = [w / total for w in weights] return random.choices(self.providers, weights=weights, k=1)[0] return self.providers[0] def report_stats(self) -> dict: """报告各供应商统计数据""" return { name: { **stats, "success_rate": ( stats["success"] / (stats["success"] + stats["failure"]) * 100 if (stats["success"] + stats["failure"]) > 0 else 0 ), "avg_latency": ( stats["total_latency"] / stats["success"] if stats["success"] > 0 else 0 ) } for name, stats in self.provider_stats.items() }
Q1:如何确定最佳的并发度?
A1:最佳并发度取决于多个因素:API供应商的RPM限制、网络延迟、业务对吞吐量的需求。一个实用的经验公式是:并发度 = RPM限制 / 平均响应时间(秒)。例如RPM为60、平均响应2秒,则理论最大并发度为2-3。建议从保守值开始,通过压测逐步调整。同时要预留20%-30%的余量应对突发流量。
Q2:遇到频繁429错误怎么处理?
A2:首先检查速率限制器配置是否正确。然后分析请求模式:如果是均匀流量超限,需要降低并发度;如果是突发流量导致,需要增加队列缓冲。确保实现了指数退避重试机制,并正确解析Retry-After响应头。对于多供应商场景,可以在当前供应商限速时自动切换到备用供应商。
Q3:消息队列会不会增加延迟?
A3:会,但可控。对于实时性要求高的场景,可使用内存队列(微秒级延迟);对于可接受秒级延迟的批处理任务,Redis队列的延迟通常在1-5毫秒。队列的核心价值不是减少延迟,而是平滑流量、避免429错误、提升整体吞吐量。权衡延迟和吞吐量时,建议优先保证系统稳定性。
最佳实践:
常见避坑:
并发控制系统架构图 ┌────────────┐ ┌──────────────┐ ┌──────────────┐ │ 业务请求 │────→│ 队列管理器 │────→│ 限速控制器 │ │ (Web/API) │ │ (流量整形) │ │ (Token Bucket│ └────────────┘ └──────────────┘ │ + 自适应) │ └──────┬───────┘ │ ┌──────▼───────┐ │ 负载均衡器 │ │ (多供应商) │ └──┬──┬──┬─────┘ │ │ │ ┌──────┘ │ └──────┐ ▼ ▼ ▼ ┌────────┐┌────────┐┌────────┐ │供应商A ││供应商B ││供应商C │ │(主力) ││(备用) ││(备用) │ └────────┘└────────┘└────────┘
本节从并发请求的设计原理出发,系统讲解了速率限制的令牌桶算法、自适应限速策略、队列化调度架构和多供应商负载均衡方案。通过令牌桶与自适应限速的组合,系统可以智能应对429错误;通过队列化调度,可以将不均匀的请求流量平滑化;通过多供应商负载均衡,可以在不超出单一供应商限制的前提下扩展整体吞吐量。这些技术的组合使用,是构建大规模大模型API调用系统的基石。
并发优化效果总结 无优化 ██ 吞吐量基准 高成本(频繁429) 固定限速 ████ 吞吐提升2x 成本稳定 自适应限速 ██████ 吞吐提升3x 成本降低 队列+均衡 ████████ 吞吐提升4x 成本最优