Agent生产环境部署架构与最佳实践

举报
柠檬🍋 发表于 2026/09/23 16:04:34 2026/09/23
【摘要】 Agent生产环境部署架构与最佳实践 引言Agent技术从实验室走向生产环境,是每一个技术团队必须面对的关键挑战。在实验室里,Agent可以跑在单机上,用Jupyter Notebook调试,失败了重启即可。但一旦进入生产环境,面对真实的用户流量、严格的可用性要求、复杂的运维场景,部署架构的设计就变得至关重要。本文将从服务化部署、容器化方案、负载均衡、高可用设计、灰度发布等多个维度,系统性...

Agent生产环境部署架构与最佳实践

引言

Agent技术从实验室走向生产环境,是每一个技术团队必须面对的关键挑战。在实验室里,Agent可以跑在单机上,用Jupyter Notebook调试,失败了重启即可。但一旦进入生产环境,面对真实的用户流量、严格的可用性要求、复杂的运维场景,部署架构的设计就变得至关重要。本文将从服务化部署、容器化方案、负载均衡、高可用设计、灰度发布等多个维度,系统性地讲解Agent生产环境部署的架构设计与最佳实践,并给出完整的代码实现。

一、Agent服务化部署的整体架构

1.1 从单机到分布式的演进

Agent应用的服务化部署,核心思路是将Agent能力封装为独立的服务单元,通过标准化的接口对外提供能力。一个典型的Agent生产环境架构包含以下层次:接入层负责接收外部请求,通常由API网关承担,处理认证、限流、路由分发。服务层是核心,包含Agent调度服务、模型推理服务、工具执行服务、记忆存储服务等。基础设施层提供底层的模型服务、向量数据库、对象存储等支撑能力。

在从单机向分布式演进的过程中,最关键的变化是Agent的状态管理。单机环境下,Agent的对话上下文、工具调用中间结果都存在内存中。分布式环境下,这些状态必须外置到独立的存储中,否则请求被路由到不同节点时,上下文就会丢失。

1.2 服务拆分原则

Agent服务的拆分需要遵循领域驱动的设计原则。按照业务能力进行拆分,而不是简单地按技术层拆分。一个合理的拆分方案如下:Agent编排服务负责整体流程调度,接收用户请求后协调各个子服务完成Agent的思考-行动-观察循环。LLM推理服务封装大语言模型的调用,处理prompt组装、参数调优、流式输出。工具服务负责执行外部工具调用,如搜索、代码执行、数据库查询等。记忆服务管理短期对话历史和长期知识检索。会话服务管理用户会话状态和上下文持久化。

每个服务独立部署、独立扩缩容、独立发布。服务之间通过gRPC或HTTP进行通信,通过消息队列进行异步解耦。

1.3 部署拓扑设计

生产环境的部署拓扑需要考虑容灾能力。典型的多可用区部署方案如下:在两个以上的可用区部署完整的服务栈,通过全局负载均衡器将流量分发到不同可用区。每个可用区内,服务以多副本形式运行,通过可用区内的负载均衡进行分发。数据库和存储层采用主从复制或分布式方案,确保跨可用区的数据一致性。这种拓扑设计能够应对可用区级别的故障,当某个可用区整体不可用时,全局负载均衡器会将流量自动切换到健康的可用区,用户几乎无感知。

二、容器化部署方案

2.1 容器镜像设计

容器化是Agent部署的基础。良好的镜像设计能够大幅提升部署效率和运行稳定性。基础镜像选择方面,推荐使用多阶段构建。第一阶段使用完整的构建环境镜像编译代码和依赖,第二阶段使用精简的运行时镜像只包含必要的运行时文件。依赖管理方面,Python生态的Agent应用通常依赖大量第三方库,包括transformers、langchain、openai等。这些依赖需要在构建阶段安装好,并生成依赖锁文件确保版本一致性。对于模型文件等大文件,不应该打包进镜像,而是通过对象存储在启动时动态加载。

# Dockerfile - Agent服务多阶段构建示例
FROM python:3.11-slim as builder
WORKDIR /build
COPY requirements.txt .
RUN pip install --user --no-cache-dir -r requirements.txt

FROM python:3.11-slim
WORKDIR /app
COPY --from=builder /root/.local /root/.local
COPY . .
ENV PATH=/root/.local/bin:$PATH
ENV PYTHONUNBUFFERED=1
ENV AGENT_ENV=production
EXPOSE 8080
HEALTHCHECK --interval=30s --timeout=5s --retries=3 \
    CMD curl -f http://localhost:8080/health || exit 1
CMD ["python", "-m", "agent.server", "--host", "0.0.0.0", "--port", "8080"]

2.2 Kubernetes部署编排

在Kubernetes上部署Agent服务,需要编写Deployment、Service、ConfigMap、Secret等资源清单。以下是一个完整的部署配置示例,包含了滚动更新策略、Pod反亲和性、资源配额、健康探针和优雅停机等关键设计:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: agent-orchestrator
  namespace: agent-prod
  labels:
    app: agent-orchestrator
spec:
  replicas: 6
  strategy:
    type: RollingUpdate
    rollingUpdate:
      maxSurge: 2
      maxUnavailable: 1
  selector:
    matchLabels:
      app: agent-orchestrator
  template:
    metadata:
      labels:
        app: agent-orchestrator
    spec:
      affinity:
        podAntiAffinity:
          preferredDuringSchedulingIgnoredDuringExecution:
          - weight: 100
            podAffinityTerm:
              labelSelector:
                matchLabels:
                  app: agent-orchestrator
              topologyKey: kubernetes.io/hostname
      containers:
      - name: agent-orchestrator
        image: registry.cn-beijing.aliyuncs.com/agent/orchestrator:v2.1.0
        ports:
        - containerPort: 8080
          name: http
        env:
        - name: AGENT_ENV
          value: "production"
        - name: REDIS_URL
          valueFrom:
            secretKeyRef:
              name: agent-secrets
              key: redis-url
        resources:
          requests:
            cpu: "2000m"
            memory: "4Gi"
          limits:
            cpu: "4000m"
            memory: "8Gi"
        livenessProbe:
          httpGet:
            path: /health
            port: 8080
          initialDelaySeconds: 30
          periodSeconds: 20
        readinessProbe:
          httpGet:
            path: /ready
            port: 8080
          initialDelaySeconds: 10
          periodSeconds: 10
        lifecycle:
          preStop:
            exec:
              command: ["/bin/sh", "-c", "sleep 15 && curl -X POST http://localhost:8080/shutdown"]

这个部署配置包含了几个关键设计点。滚动更新策略设置了maxSurge和maxUnavailable,确保更新过程中始终有足够的可用实例。Pod反亲和性配置确保副本分散在不同节点上,避免单节点故障导致服务不可用。优雅停机通过preStop钩子实现,先从负载均衡中摘除,等待15秒处理完存量请求后再关闭。

2.3 模型服务的特殊考量

Agent应用依赖的大语言模型服务,在部署上有特殊考量。如果使用自建模型服务,需要GPU资源调度,Kubernetes通过NVIDIA Device Plugin管理GPU资源。对于使用云端模型API的场景,需要考虑API调用的网络延迟和限流。建议在Agent服务和模型API之间增加一层本地代理,负责连接池管理、请求重试、限流控制。这样即使模型API出现抖动,Agent服务也能通过本地代理的缓冲机制保持稳定。

三、负载均衡设计

3.1 多层负载均衡架构

Agent生产环境的负载均衡通常采用多层架构。最外层是全局负载均衡器,负责跨可用区的流量分发。中间层是入口负载均衡器,负责跨实例的流量分发。内层是服务网格的Sidecar代理,负责服务间通信的负载均衡。每一层负载均衡的策略需要根据特点进行设计:全局层主要基于地理位置和可用区健康状态进行路由,入口层可以采用加权轮询或最少连接数策略,服务间通信层可以采用更精细的延迟感知策略。

3.2 会话亲和性处理

Agent应用的一个特点是会话上下文的连续性。同一个用户的连续请求需要访问相同的上下文数据。如果上下文完全外置到共享存储中,任何节点都可以处理任何请求。但如果为了性能将部分上下文缓存在本地内存中,就需要会话亲和性。会话亲和性的实现可以通过一致性哈希,负载均衡器根据会话ID计算哈希值,将请求路由到固定的节点。当节点扩缩容时,一致性哈希能保证大部分会话不受影响,只有少量会话需要迁移。但会话亲和性也带来了风险:如果某个节点故障,该节点上的所有会话都需要重新路由,因此即使使用会话亲和性,也必须确保上下文数据在外部存储中有完整副本。

3.3 负载均衡器实现

以下是一个基于Python的智能负载均衡器实现,支持加权轮询、健康检查和会话亲和性:

import hashlib
import time
import threading
from dataclasses import dataclass
from typing import Optional
import requests

@dataclass
class BackendNode:
    address: str
    port: int
    weight: int = 1
    healthy: bool = True
    active_connections: int = 0
    avg_response_time: float = 0.0
    total_requests: int = 0
    failed_requests: int = 0

    @property
    def endpoint(self):
        return f"http://{self.address}:{self.port}"

    @property
    def error_rate(self):
        if self.total_requests == 0:
            return 0.0
        return self.failed_requests / self.total_requests


class LoadBalancer:
    """智能负载均衡器,支持加权轮询、最少连接、延迟感知、一致性哈希"""

    def __init__(self, strategy="weighted_round_robin"):
        self.strategy = strategy
        self.nodes = []
        self._lock = threading.RLock()
        self._rr_index = 0
        self._hash_ring = {}
        self._virtual_nodes = 150
        self._running = False

    def add_node(self, address, port, weight=1):
        with self._lock:
            node = BackendNode(address=address, port=port, weight=weight)
            self.nodes.append(node)
            if self.strategy == "consistent_hash":
                self._add_to_hash_ring(node)

    def remove_node(self, address, port):
        with self._lock:
            self.nodes = [n for n in self.nodes if not (n.address == address and n.port == port)]
            if self.strategy == "consistent_hash":
                self._rebuild_hash_ring()

    def select(self, session_id=None):
        with self._lock:
            healthy = [n for n in self.nodes if n.healthy]
            if not healthy:
                return None
            if self.strategy == "consistent_hash" and session_id:
                return self._select_by_hash(session_id, healthy)
            elif self.strategy == "least_connections":
                return self._select_least_connections(healthy)
            elif self.strategy == "latency_aware":
                return self._select_latency_aware(healthy)
            else:
                return self._select_weighted_rr(healthy)

    def _select_weighted_rr(self, nodes):
        total_weight = sum(n.weight for n in nodes)
        self._rr_index = (self._rr_index + 1) % total_weight
        current = 0
        for node in nodes:
            current += node.weight
            if self._rr_index <= current:
                node.active_connections += 1
                node.total_requests += 1
                return node
        return nodes[-1]

    def _select_least_connections(self, nodes):
        selected = min(nodes, key=lambda n: n.active_connections)
        selected.active_connections += 1
        selected.total_requests += 1
        return selected

    def _select_latency_aware(self, nodes):
        def score(n):
            latency = n.avg_response_time * 10
            error = n.error_rate * 1000
            return n.active_connections + latency + error
        selected = min(nodes, key=score)
        selected.active_connections += 1
        selected.total_requests += 1
        return selected

    def _select_by_hash(self, session_id, nodes):
        h = int(hashlib.md5(session_id.encode()).hexdigest(), 16)
        sorted_keys = sorted(self._hash_ring.keys())
        for key in sorted_keys:
            if key >= h:
                node = self._hash_ring[key]
                if node in nodes:
                    node.total_requests += 1
                    return node
        node = self._hash_ring[sorted_keys[0]]
        return node if node in nodes else nodes[0]

    def _add_to_hash_ring(self, node):
        for i in range(self._virtual_nodes):
            vk = f"{node.address}:{node.port}#{i}"
            self._hash_ring[int(hashlib.md5(vk.encode()).hexdigest(), 16)] = node

    def _rebuild_hash_ring(self):
        self._hash_ring.clear()
        for node in self.nodes:
            self._add_to_hash_ring(node)

    def release(self, node, success, response_time=0):
        with self._lock:
            node.active_connections = max(0, node.active_connections - 1)
            if not success:
                node.failed_requests += 1
            if response_time > 0:
                alpha = 0.3
                node.avg_response_time = alpha * response_time + (1 - alpha) * node.avg_response_time

    def start_health_check(self):
        self._running = True
        t = threading.Thread(target=self._health_loop, daemon=True)
        t.start()

    def _health_loop(self):
        while self._running:
            for node in self.nodes:
                try:
                    resp = requests.get(f"{node.endpoint}/health", timeout=3)
                    node.healthy = resp.status_code == 200
                except Exception:
                    node.healthy = False
            time.sleep(10)

    def stop(self):
        self._running = False

这个负载均衡器实现了四种调度策略。加权轮询按照节点权重分配请求,适合节点配置均匀的场景。最少连接数选择当前连接数最少的节点,适合请求处理时间差异较大的场景。延迟感知策略综合考虑响应时间、错误率和连接数,是最智能的策略。一致性哈希策略保证相同会话ID的请求路由到同一节点,适合需要会话亲和性的场景。健康检查线程定期检测每个节点的健康状态,自动剔除故障节点。

四、高可用设计

4.1 高可用的层次设计

Agent系统的高可用设计需要从多个层次考虑。基础设施层的高可用通过多可用区部署、冗余电源、冗余网络实现。平台层的高可用通过Kubernetes的自动调度、健康检查、自动重启实现。应用层的高可用通过无状态设计、优雅降级、熔断限流实现。数据层的高可用通过主从复制、多副本、定期备份实现。每一层都需要有对应的故障检测和恢复机制,确保单点故障不会扩散为系统级故障。

4.2 优雅降级策略

当系统部分组件出现故障时,优雅降级能够保证核心功能可用。Agent系统的降级策略可以设计为多级:一级降级,当模型推理服务延迟升高时,切换到更轻量的模型或缓存的历史回复。二级降级,当工具服务不可用时,Agent跳过工具调用,仅基于自身知识回答。三级降级,当记忆服务不可用时,Agent使用当前请求的上下文进行回答,不检索历史记忆。四级降级,当整个Agent编排服务过载时,返回预设的友好提示,引导用户稍后重试。每一级降级都需要在代码中预先定义好触发条件和降级行为,并通过配置中心动态控制阈值。

4.3 故障恢复机制

故障恢复的核心是快速检测和自动恢复。Kubernetes的liveness探针可以检测应用死锁,readiness探针可以检测应用就绪状态。当liveness探针失败时,Kubernetes会自动重启Pod。当readiness探针失败时,Kubernetes会将Pod从Service的端点列表中移除。对于数据库故障,需要实现自动故障转移,主数据库故障时自动提升从数据库为主数据库。对于缓存故障,需要有降级到直接访问数据库的机制,缓存恢复后逐步重建缓存数据,避免缓存击穿导致数据库过载。

五、灰度发布实践

5.1 灰度发布策略

灰度发布是降低发布风险的关键手段。Agent系统的灰度发布可以采用以下策略:基于权重的灰度发布,通过调整Service的流量权重,逐步将流量从旧版本切换到新版本,典型流程是先切1%流量观察,然后逐步增加到5%、10%、25%、50%、100%,每一步都需要观察关键指标,如有异常立即回滚。基于特征的灰度发布,通过请求特征筛选灰度用户,可以将内部用户的请求路由到新版本。基于地域的灰度发布,先在一个地域部署新版本,验证稳定后推广到其他地域。

5.2 灰度发布实现

以下是一个基于Kubernetes和Istio的灰度发布管理器实现,包含自动指标监控和回滚机制:

import subprocess
import json
import time
import logging
from dataclasses import dataclass

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

@dataclass
class CanaryConfig:
    service_name: str
    namespace: str
    new_version: str
    old_version: str
    weight_steps: list
    observation_period: int = 120
    rollback_error_rate: float = 0.01
    rollback_latency_p99: float = 2000.0


class CanaryDeployer:
    """灰度发布管理器,支持渐进式流量切换和自动回滚"""

    def __init__(self, config):
        self.config = config
        self.current_weight = 0

    def deploy(self):
        logger.info(f"开始灰度: {self.config.service_name} v{self.config.old_version} -> v{self.config.new_version}")
        if not self._deploy_new_version():
            return False
        for weight in self.config.weight_steps:
            logger.info(f"调整灰度权重: {self.current_weight}% -> {weight}%")
            self._set_canary_weight(weight)
            self.current_weight = weight
            time.sleep(self.config.observation_period)
            metrics = self._collect_metrics()
            logger.info(f"当前指标: {json.dumps(metrics, indent=2)}")
            if self._should_rollback(metrics):
                logger.warning("指标异常,触发自动回滚!")
                self.rollback()
                return False
            logger.info(f"权重 {weight}% 阶段指标正常")
        logger.info("灰度发布完成,全量切换")
        self._set_canary_weight(100)
        self._cleanup_old_version()
        return True

    def _deploy_new_version(self):
        manifest = f"""
apiVersion: apps/v1
kind: Deployment
metadata:
  name: {self.config.service_name}-v{self.config.new_version}
  namespace: {self.config.namespace}
  labels:
    app: {self.config.service_name}
    version: v{self.config.new_version}
spec:
  replicas: 3
  selector:
    matchLabels:
      app: {self.config.service_name}
      version: v{self.config.new_version}
  template:
    metadata:
      labels:
        app: {self.config.service_name}
        version: v{self.config.new_version}
    spec:
      containers:
      - name: {self.config.service_name}
        image: registry.cn-beijing.aliyuncs.com/agent/{self.config.service_name}:v{self.config.new_version}
        ports:
        - containerPort: 8080
        resources:
          requests:
            cpu: "1000m"
            memory: "2Gi"
          limits:
            cpu: "2000m"
            memory: "4Gi"
"""
        result = subprocess.run(
            ["kubectl", "apply", "-f", "-"],
            input=manifest, capture_output=True, text=True, timeout=60
        )
        if result.returncode != 0:
            logger.error(f"部署失败: {result.stderr}")
            return False
        subprocess.run(
            ["kubectl", "rollout", "status",
             f"deployment/{self.config.service_name}-v{self.config.new_version}",
             "-n", self.config.namespace, "--timeout=300s"],
            timeout=300, check=True
        )
        return True

    def _set_canary_weight(self, weight):
        old_weight = 100 - weight
        vs = f"""
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
  name: {self.config.service_name}
  namespace: {self.config.namespace}
spec:
  hosts:
  - {self.config.service_name}
  http:
  - route:
    - destination:
        host: {self.config.service_name}
        subset: v{self.config.old_version}
      weight: {old_weight}
    - destination:
        host: {self.config.service_name}
        subset: v{self.config.new_version}
      weight: {weight}
"""
        subprocess.run(["kubectl", "apply", "-f", "-"], input=vs,
                       capture_output=True, text=True, timeout=30)

    def _collect_metrics(self):
        try:
            queries = {
                "error_rate": f'sum(rate(http_requests_total{{service="{self.config.service_name}",version="v{self.config.new_version}",code=~"5.."}}[2m])) / sum(rate(http_requests_total{{service="{self.config.service_name}",version="v{self.config.new_version}"}}[2m]))',
                "latency_p99": f'histogram_quantile(0.99, sum(rate(http_request_duration_seconds_bucket{{service="{self.config.service_name}",version="v{self.config.new_version}"}}[2m])) by (le))',
                "qps": f'sum(rate(http_requests_total{{service="{self.config.service_name}",version="v{self.config.new_version}"}}[2m]))'
            }
            metrics = {}
            for name, query in queries.items():
                result = subprocess.run(
                    ["kubectl", "exec", "-n", "monitoring", "prometheus-0", "--",
                     "wget", "-qO-", f"http://localhost:9090/api/v1/query?query={query}"],
                    capture_output=True, text=True, timeout=10
                )
                if result.returncode == 0:
                    data = json.loads(result.stdout)
                    if data.get("status") == "success" and data["data"]["result"]:
                        metrics[name] = float(data["data"]["result"][0]["value"][1])
                    else:
                        metrics[name] = 0.0
            return metrics
        except Exception as e:
            logger.error(f"收集指标异常: {e}")
            return {"error_rate": 0.0, "latency_p99": 0.0, "qps": 0.0}

    def _should_rollback(self, metrics):
        error_rate = metrics.get("error_rate", 0.0)
        latency_ms = metrics.get("latency_p99", 0.0) * 1000
        if error_rate > self.config.rollback_error_rate:
            logger.warning(f"错误率 {error_rate:.2%} 超过阈值")
            return True
        if latency_ms > self.config.rollback_latency_p99:
            logger.warning(f"P99延迟 {latency_ms:.0f}ms 超过阈值")
            return True
        return False

    def rollback(self):
        logger.info("开始回滚...")
        self._set_canary_weight(0)
        time.sleep(10)
        subprocess.run(["kubectl", "delete", "deployment",
                        f"{self.config.service_name}-v{self.config.new_version}",
                        "-n", self.config.namespace],
                       capture_output=True, timeout=30)
        logger.info("回滚完成")

    def _cleanup_old_version(self):
        subprocess.run(["kubectl", "delete", "deployment",
                        f"{self.config.service_name}-v{self.config.old_version}",
                        "-n", self.config.namespace],
                       capture_output=True, timeout=30)
        logger.info("旧版本已清理")


if __name__ == "__main__":
    config = CanaryConfig(
        service_name="agent-orchestrator",
        namespace="agent-prod",
        new_version="2.2.0",
        old_version="2.1.0",
        weight_steps=[1, 5, 10, 25, 50],
        observation_period=120,
        rollback_error_rate=0.01,
        rollback_latency_p99=2000
    )
    deployer = CanaryDeployer(config)
    success = deployer.deploy()
    logger.info("灰度发布成功!" if success else "灰度发布失败,已回滚")

这个灰度发布管理器实现了完整的渐进式发布流程。它首先部署新版本的Deployment,然后通过Istio VirtualService逐步调整流量权重。每个权重阶段都有观察期,期间通过Prometheus收集错误率、P99延迟和QPS等关键指标。如果指标超过预设阈值,自动触发回滚,将流量切回旧版本并清理新版本资源。

六、部署架构的代码实现

6.1 Agent服务框架

以下是一个完整的Agent服务框架实现,集成了健康检查、优雅停机、指标暴露等生产环境必需的功能:

import os
import time
import asyncio
import logging
from contextlib import asynccontextmanager
from dataclasses import dataclass
import json

from fastapi import FastAPI, Request, Response
import uvicorn
from prometheus_client import Counter, Histogram, Gauge, generate_latest

logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
logger = logging.getLogger("agent-server")

REQUEST_COUNT = Counter('agent_requests_total', 'Total requests', ['endpoint', 'status'])
REQUEST_LATENCY = Histogram('agent_request_duration_seconds', 'Request latency', ['endpoint'])
ACTIVE_CONNECTIONS = Gauge('agent_active_connections', 'Active connections')
MODEL_CALL_COUNT = Counter('agent_model_calls_total', 'Model calls', ['model', 'status'])


@dataclass
class ServerConfig:
    host: str = "0.0.0.0"
    port: int = 8080
    max_concurrent_requests: int = 50
    model_endpoint: str = ""
    model_name: str = "gpt-4"
    redis_url: str = ""


class AgentServer:
    """Agent生产级服务器,支持优雅停机、健康检查、指标暴露、流式响应"""

    def __init__(self, config):
        self.config = config
        self.ready = False
        self.shutting_down = False
        self._semaphore = asyncio.Semaphore(config.max_concurrent_requests)
        self.app = self._create_app()

    def _create_app(self):
        @asynccontextmanager
        async def lifespan(app):
            logger.info("Agent服务器启动中...")
            await self._startup()
            self.ready = True
            logger.info("Agent服务器就绪")
            yield
            logger.info("Agent服务器关闭中...")
            self.ready = False
            self.shutting_down = True
            await self._shutdown()
            logger.info("Agent服务器已关闭")

        app = FastAPI(title="Agent Service", version="2.1.0", lifespan=lifespan)

        @app.get("/health")
        async def health():
            if self.shutting_down:
                return Response(status_code=503, content="shutting down")
            return {"status": "ok", "version": "2.1.0"}

        @app.get("/ready")
        async def ready():
            if not self.ready or self.shutting_down:
                return Response(status_code=503, content="not ready")
            return {"status": "ready"}

        @app.get("/metrics")
        async def metrics():
            return Response(content=generate_latest(), media_type="text/plain")

        @app.post("/v1/chat")
        async def chat(request: Request):
            if self.shutting_down:
                return Response(status_code=503, content="service shutting down")
            async with self._semaphore:
                ACTIVE_CONNECTIONS.inc()
                start = time.time()
                status = "200"
                try:
                    body = await request.json()
                    result = await self._process_chat(body)
                    REQUEST_COUNT.labels(endpoint="/v1/chat", status=status).inc()
                    REQUEST_LATENCY.labels(endpoint="/v1/chat").observe(time.time() - start)
                    return result
                except asyncio.TimeoutError:
                    REQUEST_COUNT.labels(endpoint="/v1/chat", status="504").inc()
                    return Response(status_code=504, content="timeout")
                except Exception as e:
                    REQUEST_COUNT.labels(endpoint="/v1/chat", status="500").inc()
                    logger.error(f"请求处理失败: {e}", exc_info=True)
                    return Response(status_code=500, content=str(e))
                finally:
                    ACTIVE_CONNECTIONS.dec()

        @app.post("/shutdown")
        async def shutdown():
            self.shutting_down = True
            return {"status": "shutting down"}

        return app

    async def _startup(self):
        """启动初始化:连接Redis、预热模型、加载配置"""
        logger.info("初始化Redis连接...")
        await asyncio.sleep(0.5)
        logger.info("预热模型连接...")
        await asyncio.sleep(0.5)
        logger.info("加载配置...")
        await asyncio.sleep(0.2)

    async def _shutdown(self):
        """优雅关闭:等待存量请求处理完成"""
        logger.info("等待存量请求处理完成...")
        await asyncio.sleep(1)

    async def _process_chat(self, body):
        """处理Agent对话请求"""
        messages = body.get("messages", [])
        session_id = body.get("session_id", "")
        model_start = time.time()
        try:
            response = await self._call_model(messages)
            MODEL_CALL_COUNT.labels(model=self.config.model_name, status="success").inc()
            return {"response": response, "session_id": session_id}
        except Exception as e:
            MODEL_CALL_COUNT.labels(model=self.config.model_name, status="error").inc()
            raise

    async def _call_model(self, messages):
        """调用LLM模型"""
        await asyncio.sleep(0.1)
        last_msg = messages[-1]["content"] if messages else ""
        return f"Agent回复: 已收到您的消息「{last_msg[:50]}」,正在处理中..."

    def run(self):
        uvicorn.run(
            self.app, host=self.config.host, port=self.config.port,
            workers=1, log_level="info"
        )


if __name__ == "__main__":
    config = ServerConfig(
        host="0.0.0.0", port=8080,
        model_endpoint=os.getenv("MODEL_ENDPOINT", ""),
        model_name=os.getenv("MODEL_NAME", "gpt-4"),
        redis_url=os.getenv("REDIS_URL", "")
    )
    server = AgentServer(config)
    server.run()

这个Agent服务框架实现了生产环境的核心功能。通过FastAPI的lifespan上下文管理器实现优雅启动和停机。健康检查端点/health供Kubernetes的liveness探针使用,就绪端点/ready供readiness探针使用。Prometheus指标端点/metrics暴露请求计数、延迟分布、活跃连接数等关键指标。信号量_semaphore控制最大并发请求数,防止过载。/shutdown端点配合preStop钩子实现优雅停机,先标记关闭状态拒绝新请求,等待存量请求处理完成后再退出。

6.2 部署流水线

完整的部署流水线应该包含代码编译、镜像构建、安全扫描、自动化测试、部署 staging、部署 production 等阶段。以下是一个基于Python的CI/CD流水线编排器:

import subprocess
import time
import logging
from dataclasses import dataclass, field
from typing import Callable

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)


@dataclass
class PipelineStage:
    name: str
    action: Callable
    timeout: int = 300
    retry: int = 0
    required: bool = True


class DeploymentPipeline:
    """部署流水线编排器"""

    def __init__(self):
        self.stages = []

    def add_stage(self, stage):
        self.stages.append(stage)
        return self

    def run(self):
        logger.info(f"启动部署流水线,共 {len(self.stages)} 个阶段")
        for i, stage in enumerate(self.stages, 1):
            logger.info(f"[{i}/{len(self.stages)}] 执行阶段: {stage.name}")
            success = self._run_stage(stage)
            if not success and stage.required:
                logger.error(f"阶段 {stage.name} 失败,流水线中止")
                return False
            elif not success:
                logger.warning(f"阶段 {stage.name} 失败但非必需,继续执行")
        logger.info("部署流水线全部完成")
        return True

    def _run_stage(self, stage):
        for attempt in range(stage.retry + 1):
            try:
                result = stage.action()
                if result:
                    logger.info(f"阶段 {stage.name} 完成")
                    return True
                if attempt < stage.retry:
                    logger.warning(f"阶段 {stage.name} 第{attempt+1}次失败,重试中...")
                    time.sleep(5)
            except Exception as e:
                if attempt < stage.retry:
                    logger.warning(f"阶段 {stage.name} 异常: {e},重试中...")
                    time.sleep(5)
                else:
                    logger.error(f"阶段 {stage.name} 失败: {e}")
                    return False
        return False


def build_image():
    result = subprocess.run(
        ["docker", "build", "-t", "agent-orchestrator:v2.1.0", "."],
        capture_output=True, text=True, timeout=600
    )
    return result.returncode == 0


def push_image():
    result = subprocess.run(
        ["docker", "push", "registry.cn-beijing.aliyuncs.com/agent/orchestrator:v2.1.0"],
        capture_output=True, text=True, timeout=600
    )
    return result.returncode == 0


def run_tests():
    result = subprocess.run(
        ["python", "-m", "pytest", "tests/", "--tb=short", "-q"],
        capture_output=True, text=True, timeout=300
    )
    return result.returncode == 0


def deploy_staging():
    result = subprocess.run(
        ["kubectl", "apply", "-f", "k8s/staging/", "--record"],
        capture_output=True, text=True, timeout=120
    )
    return result.returncode == 0


def deploy_production():
    result = subprocess.run(
        ["kubectl", "apply", "-f", "k8s/production/", "--record"],
        capture_output=True, text=True, timeout=120
    )
    return result.returncode == 0


if __name__ == "__main__":
    pipeline = DeploymentPipeline()
    pipeline.add_stage(PipelineStage("运行测试", run_tests, timeout=300, retry=1))
    pipeline.add_stage(PipelineStage("构建镜像", build_image, timeout=600, retry=1))
    pipeline.add_stage(PipelineStage("推送镜像", push_image, timeout=600, retry=2))
    pipeline.add_stage(PipelineStage("部署Staging", deploy_staging, timeout=120))
    pipeline.add_stage(PipelineStage("部署Production", deploy_production, timeout=120))
    success = pipeline.run()

七、总结

Agent生产环境部署是一个系统工程,需要从架构设计、容器化、负载均衡、高可用、灰度发布等多个维度综合考虑。本文给出的方案和代码实现,涵盖了从镜像构建到Kubernetes编排、从负载均衡到灰度发布的完整链路。在实际落地中,还需要根据具体的业务场景和团队能力进行调整。核心原则是:服务无状态化、状态外置存储、多层冗余、渐进式发布、自动故障恢复。只有将这些原则贯穿到部署架构的每一个环节,才能构建出真正可靠、可扩展、易运维的Agent生产系统。

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

评论(0)

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

全部回复

上滑加载中

设置昵称

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

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

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