Agent弹性伸缩与资源调度策略
Agent弹性伸缩与资源调度策略
引言
Agent应用在生产环境中面临的流量模式往往具有高度不确定性。一次产品推广可能导致请求量在几分钟内增长数十倍,而深夜时段又可能只有极低的基线流量。传统的固定容量规划方式要么导致资源浪费,要么在流量高峰时服务不可用。弹性伸缩和资源调度是解决这一矛盾的核心技术手段。本文将深入探讨Agent系统的自动伸缩策略、资源监控体系、队列积压处理、成本优化方案和多租户隔离机制,并给出完整的代码实现。
一、弹性伸缩的核心挑战
1.1 Agent应用的伸缩特点
Agent应用与普通Web应用在伸缩方面有显著差异。首先是资源消耗的不均匀性,Agent的请求处理涉及大语言模型推理,CPU和内存消耗远高于普通Web请求。单个Agent请求可能需要数秒甚至数十秒的处理时间,期间持续占用计算资源。其次是流量模式的突发性,Agent应用常被集成到各类产品中,上游产品的流量波动会直接传导到Agent服务。第三是状态依赖性,Agent的对话上下文需要跨请求保持,伸缩时需要考虑状态迁移的成本。这些特点决定了Agent的伸缩策略不能简单照搬传统Web应用的HPA方案,需要定制化的伸缩指标、更精细的资源调度策略,以及专门针对长时请求的处理机制。
1.2 伸缩维度分析
Agent系统的弹性伸缩可以从多个维度进行。水平伸缩通过增加或减少服务实例数量来应对流量变化,是最常用的伸缩方式。垂直伸缩通过调整单个实例的资源配额来应对负载变化,适合处理短期的资源压力。区域伸缩通过在不同地理区域动态增减实例来应对地域性流量变化。模型伸缩通过切换不同规格的模型来平衡质量和成本,在低峰期使用轻量模型,高峰期切换到高质量模型。实际生产中,这几种维度通常组合使用,形成多维度的弹性伸缩体系。
1.3 伸缩决策的关键指标
伸缩决策不能仅依赖CPU和内存使用率。对于Agent应用,更重要的指标包括:并发请求队列深度,反映当前积压的未处理请求数量。请求等待时间,反映用户等待响应的时间。模型推理延迟,反映底层模型服务的处理速度。错误率,反映系统健康状态。活跃会话数,反映当前在线用户数量。这些指标的综合分析才能做出准确的伸缩决策。单一指标容易产生误判,例如CPU使用率高可能是因为模型推理密集,也可能是发生了死循环,需要结合其他指标综合判断。
二、自动伸缩策略设计
2.1 基于指标的伸缩
基于指标的伸缩是最基础的策略。系统持续监控关键指标,当指标超过上限阈值时触发扩容,低于下限阈值时触发缩容。为了避免频繁伸缩导致的抖动,需要设置冷却时间,在每次伸缩操作后等待一段时间再进行下一次决策。对于Agent应用,推荐的指标组合是:并发请求队列深度作为主要伸缩指标,CPU使用率作为辅助指标,模型推理延迟作为质量保障指标。当队列深度超过阈值时立即扩容,当CPU使用率持续低于阈值时考虑缩容,当模型推理延迟超过阈值时触发紧急扩容。
2.2 基于预测的伸缩
基于预测的伸缩通过分析历史流量模式,预测未来的流量变化,提前进行扩容。这种策略特别适合有规律性流量模式的Agent应用。例如,工作日的流量高峰通常出现在上午9-11点和下午14-16点,可以提前30分钟开始扩容。预测模型可以采用时间序列分析方法,如移动平均、指数平滑、ARIMA等,也可以采用机器学习方法如LSTM神经网络。预测的准确性直接影响伸缩效果,需要持续优化预测模型。
2.3 混合伸缩策略
实际生产中,单一伸缩策略往往不够。混合策略结合了基于指标和基于预测两种方式的优点。预测伸缩负责应对可预见的流量变化,提前扩容避免高峰期响应慢。指标伸缩负责应对突发流量,实时监控并在必要时紧急扩容。两者协同工作,预测伸缩作为基线保障,指标伸缩作为动态调整。
2.4 伸缩策略实现
以下是一个完整的弹性伸缩控制器实现,支持基于指标和预测的混合伸缩策略:
import time
import threading
import logging
import statistics
from dataclasses import dataclass
from collections import deque
from typing import Optional
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
logger = logging.getLogger("autoscaler")
@dataclass
class ScalingPolicy:
"""伸缩策略配置"""
min_replicas: int = 3
max_replicas: int = 50
scale_up_threshold: float = 0.7
scale_down_threshold: float = 0.3
scale_up_step: int = 2
scale_down_step: int = 1
cooldown_scale_up: int = 60
cooldown_scale_down: int = 300
emergency_latency_ms: float = 8000.0
@dataclass
class MetricsSnapshot:
"""指标快照"""
timestamp: float
queue_depth: int = 0
active_requests: int = 0
cpu_usage: float = 0.0
memory_usage: float = 0.0
avg_latency_ms: float = 0.0
error_rate: float = 0.0
current_replicas: int = 0
class MetricsCollector:
"""指标采集器,从Prometheus等监控系统获取指标"""
def __init__(self):
self._history = deque(maxlen=1440)
def collect(self) -> MetricsSnapshot:
snapshot = MetricsSnapshot(
timestamp=time.time(),
queue_depth=self._query_int("agent_queue_depth"),
active_requests=self._query_int("agent_active_requests"),
cpu_usage=self._query_float("rate(container_cpu_usage_seconds_total[5m])"),
memory_usage=self._query_float("container_memory_usage_bytes / container_memory_working_set_bytes"),
avg_latency_ms=self._query_float("histogram_quantile(0.99, agent_request_duration_seconds_bucket)") * 1000,
error_rate=self._query_float("rate(agent_errors_total[5m])"),
current_replicas=self._query_int("kube_deployment_status_replicas")
)
self._history.append(snapshot)
return snapshot
def _query_int(self, query):
return 0 # 实际调用Prometheus API
def _query_float(self, query):
return 0.0
def get_history(self, minutes: int) -> list:
count = min(minutes, len(self._history))
return list(self._history)[-count:]
class TrafficPredictor:
"""流量预测器,基于加权移动平均和趋势分析"""
def __init__(self, history_minutes: int = 60):
self.history_minutes = history_minutes
def predict(self, history: list, predict_minutes: int = 30) -> list:
if len(history) < 10:
return [0] * predict_minutes
recent = [h.active_requests for h in history[-self.history_minutes:]]
weights = [0.4, 0.3, 0.2, 0.1]
predictions = []
for i in range(predict_minutes):
if len(recent) >= len(weights):
weighted_sum = sum(w * recent[-(j+1)] for j, w in enumerate(weights))
else:
weighted_sum = statistics.mean(recent) if recent else 0
if len(recent) >= 5:
trend = (recent[-1] - recent[-5]) / 5
else:
trend = 0
predicted = max(0, weighted_sum + trend * 0.5)
predictions.append(int(predicted))
recent.append(predicted)
return predictions
class AutoScaler:
"""弹性伸缩控制器,支持指标驱动和预测驱动的混合伸缩"""
def __init__(self, policy: ScalingPolicy, collector: MetricsCollector,
predictor: Optional[TrafficPredictor] = None):
self.policy = policy
self.collector = collector
self.predictor = predictor
self._last_scale_up = 0
self._last_scale_down = 0
self._running = False
def evaluate(self) -> Optional[int]:
metrics = self.collector.collect()
current = metrics.current_replicas
logger.info(f"指标: replicas={current}, queue={metrics.queue_depth}, "
f"active={metrics.active_requests}, cpu={metrics.cpu_usage:.1%}, "
f"latency={metrics.avg_latency_ms:.0f}ms")
# 紧急扩容
if metrics.avg_latency_ms > self.policy.emergency_latency_ms:
target = min(current + self.policy.scale_up_step * 3, self.policy.max_replicas)
logger.warning(f"紧急扩容到 {target} 副本")
return target
# 基于负载比例判断
capacity = current * 10
load_ratio = metrics.active_requests / max(capacity, 1)
if load_ratio > self.policy.scale_up_threshold:
if not self._in_cooldown(True):
target = min(current + self.policy.scale_up_step, self.policy.max_replicas)
logger.info(f"负载 {load_ratio:.2f} 超阈值,扩容 {current} -> {target}")
self._last_scale_up = time.time()
return target
if load_ratio < self.policy.scale_down_threshold:
if not self._in_cooldown(False):
target = max(current - self.policy.scale_down_step, self.policy.min_replicas)
if target < current:
logger.info(f"负载 {load_ratio:.2f} 低于阈值,缩容 {current} -> {target}")
self._last_scale_down = time.time()
return target
# 预测驱动伸缩
if self.predictor:
history = self.collector.get_history(60)
predictions = self.predictor.predict(history, 30)
max_predicted = max(predictions) if predictions else 0
if max_predicted > capacity * self.policy.scale_up_threshold:
if not self._in_cooldown(True):
needed = int(max_predicted / 10) + self.policy.scale_up_step
target = min(needed, self.policy.max_replicas)
if target > current:
logger.info(f"预测驱动扩容: 预计峰值 {max_predicted},扩容到 {target}")
self._last_scale_up = time.time()
return target
return None
def _in_cooldown(self, is_scale_up: bool) -> bool:
now = time.time()
if is_scale_up:
return (now - self._last_scale_up) < self.policy.cooldown_scale_up
return (now - self._last_scale_down) < self.policy.cooldown_scale_down
def start(self, interval: int = 30):
self._running = True
thread = threading.Thread(target=self._loop, args=(interval,), daemon=True)
thread.start()
def _loop(self, interval):
while self._running:
target = self.evaluate()
if target is not None:
self._scale_to(target)
time.sleep(interval)
def _scale_to(self, target: int):
logger.info(f"执行伸缩到 {target} 副本")
# subprocess.run(["kubectl", "scale", "deployment/agent", f"--replicas={target}"])
def stop(self):
self._running = False
这个弹性伸缩控制器实现了混合伸缩策略。ScalingPolicy定义了伸缩的阈值、步长和冷却时间。MetricsCollector负责采集当前系统指标。TrafficPredictor基于加权移动平均和趋势分析预测未来流量。AutoScaler综合实时指标和预测结果做出伸缩决策,支持紧急扩容、基于负载比例的常规伸缩和预测驱动的提前扩容三种模式。
三、资源监控体系
3.1 监控指标分层
Agent系统的监控指标体系应该分为多个层次。基础设施层监控包括节点CPU、内存、磁盘IO、网络带宽等。容器层监控包括Pod的CPU请求与限制、内存使用、重启次数等。应用层监控包括请求QPS、延迟分布、错误率、队列深度等。业务层监控包括活跃用户数、对话轮次、工具调用成功率等。模型层监控包括模型推理延迟、token消耗量、模型调用错误率等。每一层指标都有其特定的告警阈值和响应策略。
3.2 监控数据采集
监控数据的采集需要覆盖push和pull两种模式。Prometheus采用pull模式,定期从应用的/metrics端点拉取指标。对于短生命周期任务和批处理作业,采用push模式通过Pushgateway推送指标。日志数据通过Fluentd或Filebeat采集,发送到Elasticsearch或Loki进行存储和分析。链路追踪数据通过OpenTelemetry SDK采集,发送到Jaeger或Tempo进行存储和查询。Agent系统特别需要关注模型调用的追踪,包括prompt长度、生成token数、推理耗时等维度。
3.3 监控告警规则
告警规则的设计需要平衡灵敏度和误报率。对于Agent系统,关键的告警规则包括:队列深度持续超过阈值5分钟触发扩容告警,错误率超过1%持续2分钟触发服务异常告警,P99延迟超过目标SLO持续5分钟触发性能告警,Pod重启次数超过3次在10分钟内触发稳定性告警,模型推理延迟超过阈值持续3分钟触发模型服务告警,节点资源使用率超过90%持续5分钟触发容量告警。每条告警规则都应配置明确的处理流程和升级机制。
四、队列积压处理
4.1 请求队列设计
Agent系统处理请求的时间较长,同步处理容易导致请求积压。引入请求队列可以将突发流量平滑化,避免后端服务被压垮。队列类型选择上,对于要求严格顺序的场景使用FIFO队列,对于可以并行处理的场景使用优先级队列。对于Agent对话场景,不同用户的请求之间没有顺序依赖,可以使用多分区队列,每个分区内部保持顺序,分区之间并行处理。队列容量规划上,队列长度需要根据预期的峰值流量和处理能力来设置,一个合理的策略是设置队列长度为平均处理能力的3-5倍,同时设置最大等待时间,超时请求直接返回友好提示。
4.2 背压机制
当队列积压超过阈值时,需要启动背压机制,防止系统进一步恶化。背压机制的实现包括几个层次。第一层是请求限流,当队列深度超过80%时,降低API网关的限流阈值,减少新请求进入。第二层是降级响应,当队列深度超过90%时,对低优先级请求返回降级响应,只处理高优先级请求。第三层是熔断保护,当队列满时,直接拒绝新请求,返回503状态码和重试建议。背压机制的核心思想是宁可拒绝部分请求,也不能让整个系统崩溃。
4.3 队列积压处理器实现
以下是一个完整的队列积压处理器实现,包含优先级队列、背压机制和动态消费者管理:
import time
import threading
import heapq
import logging
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional, Callable
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
logger = logging.getLogger("queue-manager")
class Priority(Enum):
HIGH = 1
NORMAL = 2
LOW = 3
@dataclass(order=True)
class QueueItem:
priority: int
timestamp: float
request_id: str = field(compare=False)
payload: dict = field(compare=False)
callback: Callable = field(compare=False)
class RequestQueue:
"""优先级请求队列,支持背压和超时丢弃"""
def __init__(self, max_size: int = 1000, max_wait_ms: float = 30000):
self.max_size = max_size
self.max_wait_ms = max_wait_ms
self._queue = []
self._lock = threading.Lock()
self._not_empty = threading.Condition(self._lock)
self._stats = {"enqueued": 0, "dequeued": 0, "dropped": 0, "timeout": 0}
def enqueue(self, request_id: str, payload: dict,
priority: Priority = Priority.NORMAL, callback: Callable = None) -> bool:
with self._lock:
if len(self._queue) >= self.max_size:
logger.warning(f"队列已满({len(self._queue)}/{self.max_size}),拒绝 {request_id}")
self._stats["dropped"] += 1
return False
item = QueueItem(priority.value, time.time(), request_id, payload,
callback or (lambda r: None))
heapq.heappush(self._queue, item)
self._stats["enqueued"] += 1
self._not_empty.notify()
return True
def dequeue(self, timeout: float = 1.0) -> Optional[QueueItem]:
with self._not_empty:
while not self._queue:
if not self._not_empty.wait(timeout=timeout):
return None
item = heapq.heappop(self._queue)
self._stats["dequeued"] += 1
wait_ms = (time.time() - item.timestamp) * 1000
if wait_ms > self.max_wait_ms:
logger.warning(f"请求 {item.request_id} 排队 {wait_ms:.0f}ms 超时")
self._stats["timeout"] += 1
item.callback({"status": "timeout"})
return self.dequeue(timeout=0.1)
return item
def depth_ratio(self) -> float:
with self._lock:
return len(self._queue) / self.max_size
def get_stats(self) -> dict:
with self._lock:
return {"size": len(self._queue), **self._stats,
"depth_ratio": f"{self.depth_ratio():.2%}"}
class QueueWorker:
"""队列消费者"""
def __init__(self, worker_id: str, queue: RequestQueue, handler: Callable):
self.worker_id = worker_id
self.queue = queue
self.handler = handler
self._running = False
self._processed = 0
self._errors = 0
def start(self):
self._running = True
threading.Thread(target=self._run, daemon=True).start()
def _run(self):
while self._running:
item = self.queue.dequeue(timeout=1.0)
if item is None:
continue
try:
result = self.handler(item.payload)
item.callback({"status": "success", "result": result})
self._processed += 1
except Exception as e:
logger.error(f"Worker {self.worker_id} 处理 {item.request_id} 失败: {e}")
item.callback({"status": "error", "message": str(e)})
self._errors += 1
def stop(self):
self._running = False
class BackpressureController:
"""背压控制器,根据队列深度动态调整流量"""
def __init__(self, queue: RequestQueue):
self.queue = queue
self._rate_limit = 1.0 # 1.0=不限流
def get_rate_limit(self) -> float:
"""获取当前限流比例 (0.0-1.0)"""
ratio = self.queue.depth_ratio()
if ratio > 0.9:
self._rate_limit = 0.1 # 只放行10%
elif ratio > 0.8:
self._rate_limit = 0.3
elif ratio > 0.6:
self._rate_limit = 0.7
else:
self._rate_limit = 1.0
return self._rate_limit
def should_accept(self) -> bool:
"""判断是否应该接受新请求"""
import random
return random.random() < self.get_rate_limit()
class DynamicConsumerPool:
"""动态消费者池,根据队列深度自动调整消费者数量"""
def __init__(self, queue: RequestQueue, handler: Callable,
min_workers: int = 2, max_workers: int = 20):
self.queue = queue
self.handler = handler
self.min_workers = min_workers
self.max_workers = max_workers
self.workers = []
self._running = False
def start(self):
self._running = True
for i in range(self.min_workers):
self._add_worker()
threading.Thread(target=self._adjust_loop, daemon=True).start()
def _add_worker(self):
wid = f"worker-{len(self.workers)}"
worker = QueueWorker(wid, self.queue, self.handler)
worker.start()
self.workers.append(worker)
logger.info(f"添加消费者 {wid},当前共 {len(self.workers)} 个")
def _remove_worker(self):
if len(self.workers) > self.min_workers:
worker = self.workers.pop()
worker.stop()
logger.info(f"移除消费者 {worker.worker_id},当前共 {len(self.workers)} 个")
def _adjust_loop(self):
while self._running:
ratio = self.queue.depth_ratio()
if ratio > 0.5 and len(self.workers) < self.max_workers:
self._add_worker()
elif ratio < 0.2 and len(self.workers) > self.min_workers:
self._remove_worker()
time.sleep(10)
def stop(self):
self._running = False
for w in self.workers:
w.stop()
这个队列积压处理器实现了完整的请求管理链路。RequestQueue使用堆结构实现优先级队列,支持超时自动丢弃。QueueWorker是消费者线程,从队列中取出请求并处理。BackpressureController根据队列深度动态调整流量准入比例,实现背压。DynamicConsumerPool根据队列积压情况自动增减消费者数量,实现弹性处理能力。
五、成本优化方案
5.1 资源利用率优化
Agent系统的成本主要由计算资源、模型调用和存储三部分构成。计算资源的成本优化核心是提高资源利用率。通过弹性伸缩,在低峰期自动缩减实例数量,可以大幅降低计算成本。通过请求队列的削峰填谷,可以将突发流量平滑化,减少为峰值预留的资源。通过模型推理的批处理,可以将多个请求合并处理,提高GPU利用率。
5.2 模型调用成本优化
模型调用通常是Agent系统最大的成本项。优化策略包括:模型分级使用,简单任务使用轻量模型,复杂任务使用高质量模型,通过路由策略自动选择。响应缓存,对于相同或相似的请求,直接返回缓存结果。Token优化,通过prompt压缩、上下文裁剪等技术减少token消耗。模型蒸馏,将大模型的能力蒸馏到小模型中,降低推理成本。
5.3 成本监控与预算控制
成本监控需要细化到每个维度。按租户统计资源消耗和模型调用成本,用于多租户计费。按功能模块统计成本,识别成本热点。设置预算阈值,当消耗接近预算时自动告警或触发降级。以下是一个成本监控器的实现:
import time
import threading
import logging
from dataclasses import dataclass, field
from collections import defaultdict
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("cost-monitor")
@dataclass
class CostRecord:
timestamp: float
tenant_id: str
component: str # model_inference, tool_execution, storage
cost: float
tokens_used: int = 0
requests: int = 1
class CostMonitor:
"""成本监控与预算控制"""
def __init__(self):
self._records = []
self._lock = threading.Lock()
self._budgets = {} # tenant_id -> monthly_budget
self._alerts = []
def record(self, tenant_id: str, component: str, cost: float,
tokens: int = 0, requests: int = 1):
with self._lock:
record = CostRecord(time.time(), tenant_id, component, cost, tokens, requests)
self._records.append(record)
self._check_budget(tenant_id)
def set_budget(self, tenant_id: str, monthly_budget: float):
self._budgets[tenant_id] = monthly_budget
def _check_budget(self, tenant_id: str):
budget = self._budgets.get(tenant_id)
if not budget:
return
spent = self.get_tenant_cost(tenant_id, days=30)
ratio = spent / budget
if ratio > 0.9:
alert = {"tenant": tenant_id, "type": "budget_exceeded",
"spent": spent, "budget": budget, "ratio": ratio}
self._alerts.append(alert)
logger.warning(f"租户 {tenant_id} 预算告警: 已用 {spent:.2f}/{budget:.2f} ({ratio:.0%})")
def get_tenant_cost(self, tenant_id: str, days: int = 1) -> float:
cutoff = time.time() - days * 86400
with self._lock:
return sum(r.cost for r in self._records
if r.tenant_id == tenant_id and r.timestamp >= cutoff)
def get_cost_breakdown(self, tenant_id: str = None, days: int = 1) -> dict:
cutoff = time.time() - days * 86400
with self._lock:
records = [r for r in self._records if r.timestamp >= cutoff]
if tenant_id:
records = [r for r in records if r.tenant_id == tenant_id]
breakdown = defaultdict(float)
for r in records:
breakdown[r.component] += r.cost
return dict(breakdown)
def get_daily_trend(self, tenant_id: str = None, days: int = 7) -> list:
"""获取每日成本趋势"""
cutoff = time.time() - days * 86400
with self._lock:
records = [r for r in self._records if r.timestamp >= cutoff]
if tenant_id:
records = [r for r in records if r.tenant_id == tenant_id]
daily = defaultdict(float)
for r in records:
day = time.strftime("%Y-%m-%d", time.localtime(r.timestamp))
daily[day] += r.cost
return sorted(daily.items())
六、多租户隔离机制
6.1 隔离层次设计
Agent系统服务于多个租户时,必须实现资源隔离,防止一个租户的异常影响其他租户。隔离可以分为几个层次。网络层隔离通过命名空间和网络策略实现,不同租户的Pod之间网络隔离。资源层隔离通过ResourceQuota和LimitRange实现,限制每个租户的CPU、内存和存储使用量。应用层隔离通过租户上下文和线程池隔离实现,不同租户的请求使用独立的处理线程池。数据层隔离通过数据库Schema或行级安全策略实现,确保租户数据互不可见。
6.2 资源配额管理
每个租户应该有明确的资源配额。CPU和内存配额限制租户可以使用的最大资源量。请求QPS配额限制租户的请求速率。模型调用配额限制租户的token消耗量。存储配额限制租户的数据存储量。当租户超过配额时,系统应该进行限流或降级,而不是直接拒绝服务。配额的设置应该根据租户的套餐级别动态调整,高级别租户享有更高的配额。
6.3 多租户资源调度器实现
以下是一个多租户资源调度器的实现,支持资源配额管理和优先级调度:
import time
import threading
import logging
from dataclasses import dataclass, field
from collections import defaultdict
from typing import Optional
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("tenant-scheduler")
@dataclass
class TenantQuota:
"""租户配额"""
tenant_id: str
max_qps: int = 100
max_concurrent: int = 50
max_tokens_per_day: int = 1000000
max_storage_gb: int = 10
priority: int = 0 # 0=标准, 1=高级, 2=企业
class TokenBucket:
"""令牌桶限流器"""
def __init__(self, rate: float, capacity: float):
self.rate = rate
self.capacity = capacity
self.tokens = capacity
self._last_time = time.time()
self._lock = threading.Lock()
def acquire(self, tokens: float = 1) -> bool:
with self._lock:
now = time.time()
elapsed = now - self._last_time
self.tokens = min(self.capacity, self.tokens + elapsed * self.rate)
self._last_time = now
if self.tokens >= tokens:
self.tokens -= tokens
return True
return False
class TenantScheduler:
"""多租户资源调度器"""
def __init__(self):
self._quotas = {}
self._rate_limiters = {}
self._concurrent_counts = defaultdict(int)
self._token_usage = defaultdict(int)
self._token_usage_date = {}
self._lock = threading.Lock()
def register_tenant(self, quota: TenantQuota):
self._quotas[quota.tenant_id] = quota
self._rate_limiters[quota.tenant_id] = TokenBucket(
rate=quota.max_qps, capacity=quota.max_qps * 2
)
logger.info(f"注册租户 {quota.tenant_id}: QPS={quota.max_qps}, "
f"并发={quota.max_concurrent}, 优先级={quota.priority}")
def acquire(self, tenant_id: str, tokens: int = 0) -> bool:
"""获取资源许可"""
with self._lock:
quota = self._quotas.get(tenant_id)
if not quota:
logger.warning(f"未注册租户: {tenant_id}")
return False
# 检查QPS限流
limiter = self._rate_limiters[tenant_id]
if not limiter.acquire():
logger.debug(f"租户 {tenant_id} QPS限流")
return False
# 检查并发限制
if self._concurrent_counts[tenant_id] >= quota.max_concurrent:
logger.debug(f"租户 {tenant_id} 并发数超限")
return False
# 检查token用量
self._reset_daily_usage_if_needed(tenant_id)
if tokens > 0 and self._token_usage[tenant_id] + tokens > quota.max_tokens_per_day:
logger.warning(f"租户 {tenant_id} token用量超限: "
f"{self._token_usage[tenant_id]}/{quota.max_tokens_per_day}")
return False
# 分配资源
self._concurrent_counts[tenant_id] += 1
self._token_usage[tenant_id] += tokens
return True
def release(self, tenant_id: str):
"""释放资源"""
with self._lock:
if self._concurrent_counts[tenant_id] > 0:
self._concurrent_counts[tenant_id] -= 1
def _reset_daily_usage_if_needed(self, tenant_id: str):
today = time.strftime("%Y-%m-%d")
last_date = self._token_usage_date.get(tenant_id)
if last_date != today:
self._token_usage[tenant_id] = 0
self._token_usage_date[tenant_id] = today
def get_usage(self, tenant_id: str) -> dict:
with self._lock:
quota = self._quotas.get(tenant_id, TenantQuota(tenant_id))
return {
"tenant_id": tenant_id,
"current_concurrent": self._concurrent_counts[tenant_id],
"max_concurrent": quota.max_concurrent,
"tokens_used_today": self._token_usage[tenant_id],
"max_tokens_per_day": quota.max_tokens_per_day,
"priority": quota.priority
}
if __name__ == "__main__":
scheduler = TenantScheduler()
scheduler.register_tenant(TenantQuota("tenant-a", max_qps=50, max_concurrent=20,
max_tokens_per_day=500000, priority=1))
scheduler.register_tenant(TenantQuota("tenant-b", max_qps=100, max_concurrent=50,
max_tokens_per_day=1000000, priority=2))
for i in range(10):
if scheduler.acquire("tenant-a", tokens=100):
print(f"租户A请求{i}获准")
scheduler.release("tenant-a")
if scheduler.acquire("tenant-b", tokens=100):
print(f"租户B请求{i}获准")
scheduler.release("tenant-b")
time.sleep(0.1)
print("租户A用量:", scheduler.get_usage("tenant-a"))
print("租户B用量:", scheduler.get_usage("tenant-b"))
这个多租户资源调度器实现了完整的租户隔离机制。TenantQuota定义了每个租户的资源配额。TokenBucket实现了令牌桶限流算法,控制租户的请求速率。TenantScheduler综合管理QPS限流、并发数控制和token用量限制,确保每个租户在配额范围内使用资源,不会影响其他租户。
七、总结
Agent系统的弹性伸缩与资源调度是一个需要持续优化的领域。本文从自动伸缩策略、资源监控、队列积压处理、成本优化和多租户隔离五个方面进行了系统阐述,并给出了完整的代码实现。核心要点是:伸缩决策需要基于多维指标而非单一指标,队列机制是应对突发流量的关键手段,成本优化需要贯穿从模型选择到资源调度的全链路,多租户隔离是SaaS化Agent服务的基石。在实际落地中,需要根据业务特点持续调整策略参数,找到性能、成本和可靠性之间的最佳平衡点。
- 点赞
- 收藏
- 关注作者
评论(0)