Agent弹性伸缩与资源调度策略

举报
柠檬🍋 发表于 2026/09/07 10:58:39 2026/09/07
【摘要】 Agent弹性伸缩与资源调度策略 引言Agent应用在生产环境中面临的流量模式往往具有高度不确定性。一次产品推广可能导致请求量在几分钟内增长数十倍,而深夜时段又可能只有极低的基线流量。传统的固定容量规划方式要么导致资源浪费,要么在流量高峰时服务不可用。弹性伸缩和资源调度是解决这一矛盾的核心技术手段。本文将深入探讨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服务的基石。在实际落地中,需要根据业务特点持续调整策略参数,找到性能、成本和可靠性之间的最佳平衡点。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。