Agent生产环境部署架构与最佳实践
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生产系统。
- 点赞
- 收藏
- 关注作者
评论(0)