LLMOps 生产化落地
本文档系统阐述 LLMOps 的完整生产化技术体系,涵盖模型服务部署、流量路由与降级、链路追踪、Token 与成本监控、延迟与性能优化、Prompt 版本管理、CI/CD 流水线、灰度发布与回滚、质量监控告警及完整生产系统集成的全链路知识
1. 模型服务部署
1.1 定义
模型服务部署是将 LLM 应用从开发环境推向生产环境的过程,核心组件包括模型路由器(ModelRouter)、API 网关(APIGateway)和流式输出处理器(StreamingHandler),提供统一的 API 接口、负载均衡、限流和健康检查能力。
对应 Demo: demos/08_LLMOps生产化/01_模型服务部署.py
1.2 部署架构
┌──────────────────────────────────────────────────────────────┐
│ LLM 生产服务架构 │
│ │
│ 客户端请求 │
│ │ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ APIGateway(API 网关) │ │
│ │ - 认证鉴权 - 限流 - 日志 - CORS │ │
│ └──────────────────────┬───────────────────────────────┘ │
│ ▼ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ ModelRouter(模型路由器) │ │
│ │ - 按任务选择模型 - 按成本优化 - 按延迟选择 │ │
│ └─────┬──────────┬──────────┬──────────┬───────────────┘ │
│ ▼ ▼ ▼ ▼ │
│ ┌─────────┐┌─────────┐┌─────────┐┌─────────┐ │
│ │GPT-4 ││Claude 3 ││Gemini ││本地模型 │ │
│ │Endpoint ││Endpoint ││Endpoint ││Endpoint │ │
│ └─────────┘└─────────┘└─────────┘└─────────┘ │
│ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ StreamingHandler(流式处理器) │ │
│ │ - SSE 流式输出 - 分块传输 - 实时中断 │ │
│ └──────────────────────────────────────────────────────┘ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ HealthChecker(健康检查器) │ │
│ └──────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘1.3 模型路由器
class ModelRouter:
"""模型路由器:根据策略选择最优模型"""
def __init__(self):
self.endpoints: dict[str, ModelEndpoint] = {}
self.routing_mode = RoutingMode.COST_OPTIMIZED
def route(self, request: APIRequest) -> ModelEndpoint:
"""路由请求到合适的模型端点"""
if self.routing_mode == RoutingMode.QUALITY_FIRST:
return self._best_quality(request)
elif self.routing_mode == RoutingMode.COST_OPTIMIZED:
return self._cheapest_adequate(request)
elif self.routing_mode == RoutingMode.LATENCY_FIRST:
return self._fastest(request)
else:
return self._match_task(request)1.4 API 网关
class APIGateway:
"""API 网关:统一入口,提供认证/限流/日志"""
def __init__(self):
self.rate_limiter = RateLimiter(requests_per_minute=100, tokens_per_minute=100000)
self.router = ModelRouter()
async def handle(self, request: APIRequest) -> APIResponse:
"""处理 API 请求"""
if not self._authenticate(request.api_key):
return APIResponse(status=401, error="认证失败")
if not self.rate_limiter.allow(request.user_id):
return APIResponse(status=429, error="请求频率超限")
endpoint = self.router.route(request)
try:
result = await endpoint.call(request)
return APIResponse(status=200, data=result)
except Exception as e:
return APIResponse(status=500, error=str(e))1.5 SSE 流式输出
class StreamingHandler:
"""SSE 流式输出处理器"""
async def stream_response(self, request, endpoint):
"""流式返回模型输出: data: {chunk}\n\n"""
async for chunk in endpoint.stream(request):
yield f"data: {json.dumps(chunk)}\n\n"
yield "data: [DONE]\n\n" # 结束标记SSE 流式输出时序:
客户端 服务端
│ ── POST /chat ────────────────>│
│ <── HTTP 200, text/event-stream│
│ <── data: {"token":"Hello"} ──│ (第1个 Token)
│ <── data: {"token":" world"} ─│ (第2个 Token)
│ <── data: [DONE] ─────────────│ (结束)
优势: 首个 Token 即可显示(TTFT 低)1.6 健康检查
class HealthChecker:
"""健康检查器"""
async def liveness(self) -> dict:
"""存活检查:服务是否运行"""
return {"status": "alive", "timestamp": time.time()}
async def readiness(self) -> dict:
"""就绪检查:是否准备好接收请求"""
checks = {
"model_endpoints": await self._check_endpoints(),
"database": await self._check_db(),
"cache": await self._check_cache(),
}
all_ready = all(checks.values())
return {"status": "ready" if all_ready else "not_ready", "checks": checks}1.7 与其他概念的关联
-> 流量路由与降级:部署后需要流量管理
-> 链路追踪:部署后需要可观测性
-> 完整生产系统:部署是生产系统的基础
2. 流量路由与降级
2.1 定义
流量路由与降级保障 LLM 服务高可用,通过负载均衡分发流量、降级链处理故障、级联处理器管理多模型调用、熔断器防止级联失败。
对应 Demo: demos/08_LLMOps生产化/02_流量路由与降级.py
2.2 核心组件
┌──────────────────────────────────────────────────────────────┐
│ 流量路由与降级体系 │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │LoadBalancer │ │FallbackChain │ │CascadeHandler│ │
│ │(负载均衡) │ │(降级链) │ │(级联处理) │ │
│ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │
│ └──────────┬──────┘ │ │
│ ▼ │ │
│ ┌──────────┐ │ │
│ │CircuitBkr│<──────────────────┘ │
│ │(熔断器) │ │
│ └──────────┘ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ DegradationManager(降级管理器) │ │
│ │ - 缓存降级 - 简化模型 - 静态回复 - 排队等待 │ │
│ └──────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘2.3 负载均衡
class LoadBalancer:
"""负载均衡器:将流量分发到多个端点"""
def __init__(self, strategy: RoutingStrategy = RoutingStrategy.ROUND_ROBIN):
self.strategy = strategy
self.endpoints: list[ModelEndpoint] = []
self._index = 0
def select(self) -> ModelEndpoint:
"""选择一个端点"""
available = [e for e in self.endpoints if e.healthy]
if not available:
raise NoAvailableEndpoint("所有端点不可用")
if self.strategy == RoutingStrategy.ROUND_ROBIN:
endpoint = available[self._index % len(available)]
self._index += 1
elif self.strategy == RoutingStrategy.LEAST_LATENCY:
endpoint = min(available, key=lambda e: e.avg_latency)
elif self.strategy == RoutingStrategy.WEIGHTED:
endpoint = self._weighted_select(available)
return endpoint2.4 降级链
class FallbackChain:
"""降级链:主模型失败时依次尝试备选模型"""
def __init__(self):
self.chain: list[ModelEndpoint] = []
def add(self, endpoint: ModelEndpoint, priority: int = 0):
self.chain.append((endpoint, priority))
self.chain.sort(key=lambda x: x[1])
async def call(self, request) -> ModelResponse:
errors = []
for endpoint, _ in self.chain:
try:
return await endpoint.call(request)
except Exception as e:
errors.append(f"{endpoint.name}: {e}")
continue
raise AllEndpointsFailed(f"所有端点失败: {errors}")降级链示例:
GPT-4 (高质量,贵) <- 失败(限流)
|
Claude 3 (高质量,备选) <- 失败(超时)
|
GPT-3.5 (中质量,便宜) <- 失败(网络错误)
|
缓存/静态回复 (兜底) <- 成功(质量低,但可用)
降级原则: 质量递减,可用性递增2.5 级联处理器
class CascadeHandler:
"""级联处理器:先便宜模型,质量不够再用强模型"""
async def handle(self, request) -> ModelResponse:
result = await self.cheap_model.call(request)
if self._is_good_enough(result):
return result # 质量足够
# 质量不够,用强模型重新处理(带上下文)
return await self.expensive_model.call(request, context=result)2.6 熔断器
class CircuitBreaker:
"""熔断器:防止级联失败 (CLOSED -> OPEN -> HALF_OPEN -> CLOSED)"""
def __init__(self, threshold: int = 5, timeout: float = 60.0):
self.threshold = threshold
self.timeout = timeout
self.state = CircuitState.CLOSED
self.failure_count = 0
self.last_failure_time = 0
async def call(self, func, *args):
if self.state == CircuitState.OPEN:
if time.time() - self.last_failure_time > self.timeout:
self.state = CircuitState.HALF_OPEN
else:
raise CircuitOpenError("熔断器开启,请求被拒绝")
try:
result = await func(*args)
self.failure_count = 0
self.state = CircuitState.CLOSED
return result
except Exception as e:
self.failure_count += 1
self.last_failure_time = time.time()
if self.failure_count >= self.threshold:
self.state = CircuitState.OPEN
raise2.7 降级管理器
class DegradationManager:
"""降级管理器:在故障时提供降级服务"""
async def degraded_response(self, request) -> ModelResponse:
# 1. 尝试缓存命中
cache_key = self._hash(request)
if cache_key in self.cache:
return ModelResponse(content=self.cache[cache_key], degraded=True, source="cache")
# 2. 尝试静态回复
category = self._classify(request)
if category in self.static_replies:
return ModelResponse(content=self.static_replies[category], degraded=True, source="static")
# 3. 返回排队提示
return ModelResponse(content="服务繁忙,请稍后重试", degraded=True, source="fallback")2.8 与其他概念的关联
<- 模型服务部署:路由降级是部署后的流量管理
-> 链路追踪:降级过程需要追踪
-> 完整生产系统:降级是生产系统的可用性保障
3. 链路追踪
3.1 定义
链路追踪(Tracing)记录 LLM 应用中每个请求的完整执行路径,包括 LLM 调用、工具调用、检索操作等,通过 Span/Trace 层级结构可视化请求的全生命周期。
对应 Demo: demos/08_LLMOps生产化/03_链路追踪.py
3.2 Span 与 Trace 结构
Trace(一条完整请求链路)
│
├── Span: API Gateway (2.5s)
│ ├── Span: Auth (0.01s)
│ ├── Span: Rate Limit (0.005s)
│ └── Span: Model Router (0.02s)
│
├── Span: Agent Execution (2.3s)
│ ├── Span: LLM Call #1 (0.8s)
│ │ ├── Span: Prompt Construction (0.05s)
│ │ └── Span: API Call (0.7s)
│ ├── Span: Tool Call: search (0.5s)
│ └── Span: LLM Call #2 (1.0s)
│
└── Span: Response Formatting (0.15s)3.3 Span 数据结构
class Span:
"""Span:链路追踪的最小单元"""
span_id: str
trace_id: str # 所属 Trace ID
parent_id: str | None # 父 Span ID
name: str # Span 名称
start_time: float
end_time: float
status: SpanStatus # OK / ERROR / TIMEOUT
attributes: dict # 自定义属性(模型名/Token数/工具名)
events: list[dict] # 事件日志
class Trace:
"""Trace:一条完整的请求链路"""
trace_id: str
spans: list[Span]
start_time: float
end_time: float
def duration(self) -> float:
return self.end_time - self.start_time3.4 Agent 追踪器
class AgentTracer:
"""Agent 追踪器:自动追踪 Agent 执行"""
@contextmanager
def trace_llm_call(self, model: str, prompt: str, parent_span=None):
"""追踪 LLM 调用"""
with self.tracer.start_span("llm_call", parent=parent_span) as span:
span.set_attribute("model", model)
span.set_attribute("prompt_length", len(prompt))
start = time.time()
try:
yield span
except Exception as e:
span.set_status(SpanStatus.ERROR)
span.add_event("error", {"message": str(e)})
raise
finally:
span.set_attribute("duration", time.time() - start)
@contextmanager
def trace_tool_call(self, tool_name: str, args: dict, parent_span=None):
"""追踪工具调用"""
with self.tracer.start_span(f"tool:{tool_name}", parent=parent_span) as span:
span.set_attribute("tool", tool_name)
span.set_attribute("args", json.dumps(args))
yield span3.5 追踪分析器
class TraceAnalyzer:
"""追踪分析器:从 Trace 中提取洞察"""
def analyze(self, trace: Trace) -> dict:
return {
"total_duration": trace.duration(),
"span_count": len(trace.spans),
"llm_calls": sum(1 for s in trace.spans if "llm" in s.name),
"tool_calls": sum(1 for s in trace.spans if "tool" in s.name),
"errors": [s for s in trace.spans if s.status == SpanStatus.ERROR],
"slowest_span": max(trace.spans, key=lambda s: s.end_time - s.start_time),
"token_usage": sum(s.attributes.get("tokens", 0) for s in trace.spans),
}
def find_bottleneck(self, trace: Trace) -> Span:
"""找到最慢的 Span(瓶颈)"""
return max(trace.spans, key=lambda s: s.end_time - s.start_time)3.6 与其他概念的关联
<- 流量路由与降级:降级过程通过追踪可观测
-> Token 与成本监控:追踪中包含 Token 使用信息
-> 延迟与性能优化:追踪是性能优化的依据
-> 完整生产系统:追踪是生产系统的可观测性基础
4. Token 与成本监控
4.1 定义
Token 与成本监控实时追踪 LLM 应用的 Token 消耗和成本支出,通过预算管理器设置阈值,在接近预算时发出三级告警,并支持成本优化策略。
对应 Demo: demos/08_LLMOps生产化/04_Token与成本监控.py
4.2 定价表
class PricingTable:
"""模型定价表"""
PRICING = {
"gpt-4": ModelPricing("gpt-4", 30.0, 60.0), # $30/$60 per 1M tokens
"gpt-4-turbo": ModelPricing("gpt-4-turbo", 10.0, 30.0),
"gpt-3.5-turbo": ModelPricing("gpt-3.5-turbo", 0.5, 1.5),
"claude-3-opus": ModelPricing("claude-3-opus", 15.0, 75.0),
"claude-3-sonnet": ModelPricing("claude-3-sonnet", 3.0, 15.0),
}
def calculate_cost(self, model: str, input_tokens: int, output_tokens: int) -> float:
pricing = self.PRICING.get(model)
if not pricing:
return 0.0
return input_tokens / 1_000_000 * pricing.input_price + output_tokens / 1_000_000 * pricing.output_price4.3 使用量追踪与预算管理
class UsageTracker:
"""使用量追踪器"""
def record(self, model, input_tokens, output_tokens, user_id=""):
cost = self.pricing.calculate_cost(model, input_tokens, output_tokens)
self.records.append(TokenUsage(model, input_tokens, output_tokens, cost, user_id))
class BudgetManager:
"""预算管理器:三级告警"""
def check(self, cost: float) -> list[BudgetAlert]:
self.current_usage += cost
ratio = self.current_usage / self.budget.limit
if ratio >= 1.0:
return [BudgetAlert(AlertLevel.CRITICAL, f"预算已用尽!")]
elif ratio >= 0.8:
return [BudgetAlert(AlertLevel.WARNING, f"预算使用 80%!")]
elif ratio >= 0.5:
return [BudgetAlert(AlertLevel.INFO, f"预算使用 50%")]
return []4.4 三级告警
┌──────────┬──────────────────────────────────────────────┐
│ INFO │ 预算使用 50% -> 记录日志,开发团队知悉 │
│ WARNING │ 预算使用 80% -> 通知运维,可能限流/降级 │
│ CRITICAL │ 预算使用 100% -> 立即告警,强制降级或拒绝请求 │
└──────────┴──────────────────────────────────────────────┘4.5 成本优化
class CostOptimizer:
"""成本优化器:建议更经济的模型"""
def optimize(self, request, current_model: str) -> str:
if self._is_simple_task(request):
return "gpt-3.5-turbo" # 简单任务用便宜模型
if self._is_complex_task(request):
return current_model # 复杂任务保持强模型
return "gpt-4-turbo" # 中等任务用中等模型4.6 与其他概念的关联
<- 链路追踪:追踪中包含 Token 使用数据
-> 延迟与性能优化:Token 优化同时降低延迟和成本
-> 完整生产系统:成本监控是生产系统的财务保障
5. 延迟与性能优化
5.1 定义
延迟与性能优化通过缓存(精确 + 语义)、批量处理和流式优化三大策略,降低 LLM 应用的响应延迟,提升用户体验。核心指标包括 TTFT(首 Token 延迟)和整体延迟。
对应 Demo: demos/08_LLMOps生产化/05_延迟与性能优化.py
5.2 三大优化策略
┌──────────┬───────────────────────────────────────────────────┐
│ Cache │ 精确缓存(完全匹配) + 语义缓存(语义相似) │
│ │ 命中时延迟接近 0,适用: 重复查询、事实性问题 │
├──────────┼───────────────────────────────────────────────────┤
│ Batch │ 合并多个请求为一次 API 调用 │
│ │ 减少请求次数,适用: 独立的多个小任务 │
├──────────┼───────────────────────────────────────────────────┤
│ Stream │ SSE 流式输出,降低 TTFT │
│ │ 用户无需等待完整响应,适用: 所有文本生成场景 │
└──────────┴───────────────────────────────────────────────────┘5.3 缓存层
class ExactCache:
"""精确缓存:完全匹配查询"""
def get(self, key: str) -> Any | None:
entry = self.cache.get(key)
if entry and time.time() - entry.timestamp < self.ttl:
return entry.value
return None
class SemanticCache:
"""语义缓存:语义相似查询命中"""
def get(self, query_embedding: list[float]) -> Any | None:
for entry in self.entries:
similarity = self._cosine_similarity(query_embedding, entry.embedding)
if similarity >= self.threshold:
return entry.value
return None
class CacheManager:
"""缓存管理器:精确缓存 + 语义缓存"""
async def get_or_compute(self, query: str, compute_func) -> Any:
# 1. 精确缓存
cached = self.exact.get(self._hash(query))
if cached is not None:
return cached # 命中!
# 2. 语义缓存
embedding = await self._embed(query)
cached = self.semantic.get(embedding)
if cached is not None:
return cached # 语义命中!
# 3. 计算
result = await compute_func(query)
# 4. 写入缓存
self.exact.set(self._hash(query), result)
return result5.4 批量处理器
class BatchProcessor:
"""批量处理器:合并多个请求"""
def __init__(self, max_batch_size: int = 10, max_wait: float = 0.1):
self.max_batch_size = max_batch_size
self.queue: list[BatchRequest] = []
async def submit(self, request) -> BatchResult:
future = asyncio.Future()
self.queue.append(BatchRequest(request=request, future=future))
if len(self.queue) >= self.max_batch_size:
asyncio.create_task(self._flush()) # 队列满,触发批处理
return await future5.5 延迟分解
LLM 请求延迟分解 (典型值):
┌──────────────┬──────────┬────────────────────────────┐
│ 网络延迟 │ 0.1-0.5s │ 客户端到服务器的网络传输 │
│ 排队延迟 │ 0-2.0s │ 限流队列等待时间 │
│ 模型推理 │ 0.3-1.0s │ Prefill 阶段(处理 Prompt) │
│ Token 生成 │ 0.5-3.0s │ Decode 阶段(逐 Token 生成) │
│ 传输延迟 │ 0.1-0.3s │ 响应传输到客户端 │
├──────────────┼──────────┼────────────────────────────┤
│ 总延迟 │ 1-6s │ 无缓存时的典型延迟 │
│ 缓存命中延迟 │ 0.01-0.1s│ 几乎瞬时 │
└──────────────┴──────────┴────────────────────────────┘
优化: 缓存消除模型推理+Token 生成; 流式降低 TTFT; 批处理减少网络延迟5.6 与其他概念的关联
<- Token 与成本监控:缓存同时降低延迟和成本
-> Prompt 版本管理:缓存 key 包含 Prompt 版本
-> 完整生产系统:性能优化是生产系统的用户体验保障
6. Prompt 版本管理
6.1 定义
Prompt 版本管理(PromptRegistry)对 LLM 应用的 Prompt 模板进行版本化管理,支持版本注册、A/B 测试对比不同 Prompt 版本效果、以及版本间平滑迁移。
对应 Demo: demos/08_LLMOps生产化/06_Prompt版本管理.py
6.2 版本注册表
class PromptVersion:
"""Prompt 版本"""
prompt_id: str
version: str # 版本号 "1.2.0"
template: str # Prompt 模板
variables: list[str] # 模板变量
is_active: bool # 是否当前激活版本
class PromptRegistry:
"""Prompt 版本注册表"""
def register(self, prompt_id, template, variables, description=""):
"""注册新版本"""
version = f"{len(self.versions.get(prompt_id, [])) + 1}.0.0"
pv = PromptVersion(prompt_id, version, template, variables, False)
self.versions.setdefault(prompt_id, []).append(pv)
return pv
def get_active(self, prompt_id: str) -> PromptVersion:
"""获取当前激活版本"""
for v in reversed(self.versions.get(prompt_id, [])):
if v.is_active:
return v6.3 A/B 测试
class PromptABTest:
"""Prompt A/B 测试"""
def __init__(self, prompt_id, version_a, version_b):
self.version_a = version_a
self.version_b = version_b
self.traffic_split = 0.5 # 50%/50%
self.results = {"a": [], "b": []}
def assign(self, user_id: str) -> str:
"""分配用户到 A 或 B 组(同一用户始终同一组)"""
hash_val = hash(user_id) % 100 / 100
return self.version_a if hash_val < self.traffic_split else self.version_b
def analyze(self) -> ABTestResult:
a_avg = sum(self.results["a"]) / len(self.results["a"]) if self.results["a"] else 0
b_avg = sum(self.results["b"]) / len(self.results["b"]) if self.results["b"] else 0
return ABTestResult(
version_a_avg=a_avg, version_b_avg=b_avg, winner="a" if a_avg > b_avg else "b" if b_avg > a_avg else "tie"
)6.4 版本迁移
class PromptMigration:
"""版本迁移:渐进式灰度切换"""
STAGES = [0.01, 0.1, 0.5, 1.0] # 1% -> 10% -> 50% -> 100%
def advance(self):
"""推进迁移进度"""
idx = self.STAGES.index(self.rollout_percentage) if self.rollout_percentage in self.STAGES else 0
if idx < len(self.STAGES) - 1:
self.rollout_percentage = self.STAGES[idx + 1]
if self.rollout_percentage >= 1.0:
self.status = MigrationStatus.COMPLETED
def rollback(self):
"""回滚迁移"""
self.rollout_percentage = 0
self.status = MigrationStatus.ROLLED_BACK6.5 与其他概念的关联
<- 延迟与性能优化:缓存 key 包含 Prompt 版本
-> CI/CD 流水线:Prompt 变更通过 CI/CD 发布
-> 灰度发布与回滚:Prompt 迁移使用灰度策略
7. CI/CD 流水线
7.1 定义
CI/CD 流水线将 LLM 应用的构建、测试、评测、部署自动化,通过六阶段流水线和 EvalGate 门禁机制,确保每次变更都经过质量验证才能上线。
对应 Demo: demos/08_LLMOps生产化/07_CI-CD流水线.py
7.2 六阶段流水线
代码提交
│
▼
┌──────────┐
│ 1. BUILD │ 构建镜像/依赖安装/代码编译
└────┬─────┘
▼
┌──────────┐
│ 2. TEST │ 单元测试/集成测试/代码检查
└────┬─────┘
▼
┌──────────┐
│ 3. EVAL │ 评测执行(正确性/忠实度/轨迹/Agent)
└────┬─────┘
▼
┌──────────┐
│ 4. GATE │ 门禁检查(通过率>=95%? 安全=100%? 回归=0?)
└────┬─────┘
│ 通过 │ 不通过
▼ ▼
┌──────────┐ 阻止合并 + 发送报告
│ 5.DEPLOY │ 部署到预发/灰度环境
└────┬─────┘
▼
┌──────────┐
│6.VERIFY │ 线上验证(冒烟测试/健康检查)
└──────────┘7.3 评测门禁
class EvalGate:
"""评测门禁:决定是否允许部署"""
def __init__(self):
self.rules = [
("pass_rate", ">=", 0.95), # 通过率 >= 95%
("avg_score", ">=", 0.85), # 平均分 >= 0.85
("safety", "==", 1.0), # 安全性 = 100%
("regression", "==", 0), # 回归数 = 0
("p95_latency", "<=", 5.0), # P95 延迟 <= 5s
]
def check(self, eval_results: dict) -> GateResult:
failures = []
for metric, op, threshold in self.rules:
value = eval_results.get(metric, 0)
if not self._compare(value, op, threshold):
failures.append(f"{metric}: {value} {op} {threshold} 未满足")
return GateResult(passed=len(failures) == 0, failures=failures)7.4 流水线编排
class PipelineOrchestrator:
"""流水线编排器"""
async def run(self, commit_sha: str) -> PipelineReport:
stages = [
PipelineStage.BUILD,
PipelineStage.TEST,
PipelineStage.EVAL,
PipelineStage.GATE,
PipelineStage.DEPLOY,
PipelineStage.VERIFY,
]
results = {}
for stage in stages:
try:
result = await self._create_step(stage, commit_sha).run()
results[stage.value] = result
if stage == PipelineStage.GATE and not result.passed:
return PipelineReport(status="blocked", stages=results, message="评测门禁未通过")
except Exception as e:
return PipelineReport(status="failed", stages=results, error=str(e))
return PipelineReport(status="success", stages=results)7.5 与其他概念的关联
<- Prompt 版本管理:Prompt 变更通过 CI/CD 发布
-> 灰度发布与回滚:CI/CD 部署阶段使用灰度策略
-> 质量监控告警:CI/CD 门禁是质量监控的一部分
8. 灰度发布与回滚
8.1 定义
灰度发布与回滚通过金丝雀发布(四阶段渐进灰度)、蓝绿部署(双环境切换)和回滚管理器(快速回退),确保新版本上线风险可控。
对应 Demo: demos/08_LLMOps生产化/08_灰度发布与回滚.py
8.2 金丝雀发布(四阶段)
class CanaryRelease:
"""金丝雀发布:四阶段渐进灰度"""
STAGES = [
("stage_1", 0.01, 900), # 1% 流量,观察 15 分钟
("stage_2", 0.05, 1800), # 5% 流量,观察 30 分钟
("stage_3", 0.25, 3600), # 25% 流量,观察 1 小时
("stage_4", 1.00, 0), # 100% 流量,灰度完成
]
def route(self) -> str:
"""路由请求到新版本或旧版本"""
ratio = self.STAGES[self.current_stage][1]
return "canary" if random.random() < ratio else "stable"
def check_health(self) -> bool:
"""检查金丝雀版本健康度"""
recent = self.health_metrics[-100:]
if len(recent) < 100:
return True
success_rate = sum(1 for m in recent if m.success) / len(recent)
return success_rate >= 0.95 # 成功率 >= 95%
def promote(self):
"""进入下一阶段"""
if self.current_stage < len(self.STAGES) - 1:
self.current_stage += 1
def rollback(self):
"""回滚到旧版本"""
self.current_stage = 0
self.status = CanaryStatus.ROLLED_BACK金丝雀四阶段发布:
阶段1: 1% ──> 观察15min ──> 健康? ──是──> 阶段2
阶段2: 5% ──> 观察30min ──> 健康? ──是──> 阶段3
阶段3: 25% ──> 观察1h ──> 健康? ──是──> 阶段4
阶段4: 100% ──> 灰度完成
(任何阶段不健康 -> 回滚到 0%)8.3 蓝绿部署
class BlueGreenDeployment:
"""蓝绿部署:双环境切换"""
def __init__(self):
self.blue = {"version": "v1", "status": BlueGreenStatus.ACTIVE}
self.green = {"version": "v2", "status": BlueGreenStatus.IDLE}
def switch(self):
"""切换蓝绿环境"""
self.green["status"] = BlueGreenStatus.ACTIVE
self.blue["status"] = BlueGreenStatus.IDLE # 保留用于回滚
def rollback(self):
"""回滚到 blue 环境"""
self.blue["status"] = BlueGreenStatus.ACTIVE
self.green["status"] = BlueGreenStatus.IDLE8.4 回滚管理器
class RollbackManager:
"""回滚管理器"""
def create_checkpoint(self, version: str, config: dict):
"""创建回滚检查点"""
self.checkpoints.append(Checkpoint(version, config, time.time()))
def rollback(self, target_version: str = None) -> Checkpoint:
"""回滚到指定版本(默认上一个稳定版本)"""
if target_version:
cp = next((c for c in self.checkpoints if c.version == target_version), None)
else:
cp = self.checkpoints[-2] if len(self.checkpoints) >= 2 else None
if not cp:
raise NoCheckpointError("无可用回滚检查点")
self._deploy(cp)
return cp
def auto_rollback_on_failure(self, health_check, threshold: int = 3):
"""自动回滚:连续失败达到阈值时触发"""
failures = 0
while True:
if not health_check():
failures += 1
if failures >= threshold:
self.rollback()
return True
else:
failures = 0
time.sleep(10)8.5 与其他概念的关联
<- CI/CD 流水线:灰度是 CI/CD 部署阶段的策略
-> 质量监控告警:灰度健康度依赖质量监控
-> 完整生产系统:灰度回滚是生产系统的安全网
9. 质量监控告警
9.1 定义
质量监控告警(QualityMonitor)持续监控线上 LLM 应用的质量指标,基于 SLO 设定告警规则,检测质量漂移,通过 AlertManager 发送告警,在 Dashboard 可视化展示。
对应 Demo: demos/08_LLMOps生产化/09_质量监控告警.py
9.2 SLO(服务等级目标)
class SLO:
"""服务等级目标"""
name: str
target: float # 目标值
window: float # 时间窗口(秒)
metric: str # 指标名
# 典型 SLO 配置
SLOS = [
SLO("availability", 0.999, 86400, "success_rate"), # 可用性 99.9%
SLO("latency_p95", 5.0, 3600, "p95_latency"), # P95 延迟 < 5s
SLO("quality_score", 0.85, 3600, "eval_score"), # 质量分 >= 0.85
SLO("safety", 1.0, 3600, "safety_pass_rate"), # 安全性 100%
SLO("cost_per_request", 0.05, 3600, "avg_cost"), # 单次成本 < $0.05
]9.3 告警管理
class AlertRule:
"""告警规则"""
metric: str
condition: str # ">=", "<=", "=="
threshold: float
severity: AlertSeverity # INFO / WARNING / CRITICAL
cooldown: float # 冷却时间(避免重复告警)
class AlertManager:
"""告警管理器"""
def check(self, metrics: dict) -> list[Alert]:
new_alerts = []
for rule in self.rules:
value = metrics.get(rule.metric, 0)
if self._violated(value, rule):
if time.time() - self.last_triggered.get(rule.name, 0) < rule.cooldown:
continue # 冷却中
alert = Alert(
rule=rule.name,
severity=rule.severity,
message=f"{rule.metric}={value} 违反 SLO {rule.condition}{rule.threshold}",
)
self.active_alerts.append(alert)
self.last_triggered[rule.name] = time.time()
new_alerts.append(alert)
self._notify(alert) # 发送通知(邮件/Slack/PagerDuty)
return new_alerts9.4 质量监控器
class QualityMonitor:
"""质量监控器:持续监控线上质量"""
def record(self, request: RequestRecord):
self.records.append(request)
if len(self.records) % 100 == 0:
self._check_slos()
def _compute_metrics(self) -> dict:
recent = self.records[-1000:]
return {
"success_rate": sum(1 for r in recent if r.success) / len(recent),
"p95_latency": self._percentile([r.latency for r in recent], 95),
"eval_score": sum(r.eval_score for r in recent) / len(recent),
"safety_pass_rate": sum(1 for r in recent if r.safe) / len(recent),
"avg_cost": sum(r.cost for r in recent) / len(recent),
}9.5 监控仪表盘
┌──────────────────────────────────────────────────────────────┐
│ MonitoringDashboard(监控仪表盘) │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ 可用性 │ │ P95延迟 │ │ 质量分 │ │ 单次成本 │ │
│ │ 99.95% │ │ 3.2s │ │ 0.88 │ │ $0.032 │ │
│ │ SLO:99.9%│ │ SLO:<5s │ │ SLO:>.85 │ │ SLO:<.05 │ │
│ │ 达标 │ │ 达标 │ │ 达标 │ │ 达标 │ │
│ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │
│ 活跃告警: 0 │
│ 趋势图: 质量 0.90 ┤ ████████████████████████████ │
│ 0.85 ┤ ████████████████████████████ <- SLO 线 │
│ └──────────────────────────── │
│ 00:00 06:00 12:00 18:00 │
└──────────────────────────────────────────────────────────────┘9.6 与其他概念的关联
<- 灰度发布与回滚:灰度健康度依赖质量监控
<- CI/CD 流水线:SLO 是 CI/CD 门禁的依据
-> 完整生产系统:质量监控是生产系统的观测能力
10. 完整生产系统
10.1 定义
完整生产系统(ProductionLLMSystem)是将模型服务、流量路由、链路追踪、成本监控、性能优化、Prompt 管理、CI/CD、灰度发布、质量监控等所有组件有机集成的端到端生产级 LLM 系统。
对应 Demo: demos/08_LLMOps生产化/10_完整生产系统.py
10.2 系统架构
┌──────────────────────────────────────────────────────────────┐
│ ProductionLLMSystem(完整生产系统) │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │ APIGateway │ │ ModelRouter│ │StreamingHnd│ 服务层 │
│ └────────────┘ └────────────┘ └────────────┘ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │LoadBalancer│ │FallbackChn │ │CircuitBrkr │ 可靠性层 │
│ └────────────┘ └────────────┘ └────────────┘ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │AgentTracer │ │UsageTracker│ │BudgetMgr │ 可观测性层 │
│ └────────────┘ └────────────┘ └────────────┘ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │CacheManager│ │BatchProc │ │PromptReg │ 性能层 │
│ └────────────┘ └────────────┘ └────────────┘ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │CIPipeline │ │CanaryRel │ │QualityMon │ 运维层 │
│ └────────────┘ └────────────┘ └────────────┘ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ SystemStatus(统一状态面板) │ │
│ └──────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────────┘10.3 集成实现
class ProductionLLMSystem:
"""完整生产 LLM 系统:端到端集成"""
def __init__(self, config: SystemConfig):
# 服务层
self.gateway = APIGateway()
self.router = ModelRouter()
self.streaming = StreamingHandler()
# 可靠性层
self.load_balancer = LoadBalancer()
self.fallback = FallbackChain()
self.circuit_breaker = CircuitBreaker()
# 可观测性层
self.tracer = AgentTracer()
self.usage_tracker = UsageTracker()
self.budget_manager = BudgetManager(config.budget)
# 性能层
self.cache = CacheManager()
self.batch_processor = BatchProcessor()
self.prompt_registry = PromptRegistry()
# 运维层
self.ci_pipeline = CIPipeline()
self.canary = CanaryRelease()
self.quality_monitor = QualityMonitor()
async def handle_request(self, request) -> str:
"""处理用户请求(集成所有组件)"""
with self.tracer.trace_request(request) as trace:
self.gateway.authenticate(request) # 1. 网关认证
cached = self.cache.get_or_compute( # 2. 缓存检查
request.query, lambda q: None
)
if cached:
return cached # 缓存命中!
endpoint = self.router.route(request) # 3. 路由
try:
with self.tracer.trace_llm_call( # 4. 追踪+调用
endpoint.name, request.query, trace
):
response = await self.circuit_breaker.call(endpoint.call, request)
except Exception:
response = await self.fallback.call(request) # 5. 降级
self.usage_tracker.record( # 6. 记录用量
endpoint.name, response.input_tokens, response.output_tokens
)
self.budget_manager.check(response.cost) # 7. 预算检查
self.quality_monitor.record(
RequestRecord( # 8. 质量监控
success=True, latency=response.latency, eval_score=response.score, cost=response.cost
)
)
self.cache.set(request.query, response.content) # 9. 写缓存
return response.content
def get_status(self) -> SystemStatus:
"""获取系统状态"""
metrics = self.quality_monitor._compute_metrics()
return SystemStatus(
availability=metrics["success_rate"],
p95_latency=metrics["p95_latency"],
quality_score=metrics["eval_score"],
daily_cost=self.usage_tracker.summary(86400)["total_cost"],
active_alerts=len(self.quality_monitor.alert_manager.active_alerts),
active_version=self.prompt_registry.get_active("main").version,
canary_stage=self.canary.current_stage,
)10.4 请求全链路
用户请求完整链路:
-> 1. API Gateway: 认证 + 限流
-> 2. Tracer: 开始追踪
-> 3. Cache: 检查缓存(命中则直接返回)
-> 4. Model Router: 选择最优模型
-> 5. Load Balancer: 负载均衡
-> 6. Circuit Breaker: 熔断检查
-> 7. LLM Call: 调用模型(流式输出)
-> 8. [失败] Fallback Chain: 降级处理
-> 9. Usage Tracker: 记录 Token/成本
-> 10. Budget Manager: 预算检查
-> 11. Quality Monitor: 质量监控
-> 12. Cache: 写入缓存
-> 13. Tracer: 结束追踪
-> 14. 返回响应给用户10.5 与其他概念的关联
<- 模型服务部署:生产系统包含服务部署
<- 流量路由与降级:生产系统包含路由降级
<- 链路追踪:生产系统包含可观测性
<- Token 与成本监控:生产系统包含成本管理
<- 延迟与性能优化:生产系统包含性能优化
<- Prompt 版本管理:生产系统包含配置管理
<- CI/CD 流水线:生产系统包含交付管道
<- 灰度发布与回滚:生产系统包含安全发布
<- 质量监控告警:生产系统包含质量保障
概念关系总览
┌──────────────────────────────────────────┐
│ LLMOps 生产化技术体系 │
└──────────────────────────────────────────┘
服务层 可靠性层 可观测性层
│ │ │
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│模型服务 │──┐ │流量路由 │ ┌───│链路追踪 │
│部署 │ │ │与降级 │ │ └──────────┘
└──────────┘ │ └──────────┘ │ ┌──────────┐
│ │ └───│Token与 │
│ ▼ │成本监控 │
│ ┌──────────┐ └──────────┘
│ │熔断器 │ │
│ └──────────┘ │
│ │
┌─────────┴───────────────────────────────┐
│ 性能层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐│
│ │缓存 │ │批处理 │ │流式优化 ││
│ └──────────┘ └──────────┘ └──────────┘│
└──────────────────────┬──────────────────┘
│
┌──────────────────────┴──────────────────┐
│ 运维层 │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐│
│ │Prompt │ │CI/CD │ │灰度发布 ││
│ │版本管理 │ │流水线 │ │与回滚 ││
│ └──────────┘ └──────────┘ └──────────┘│
│ ┌──────────┐ │
│ │质量监控 │ │
│ │告警 │ │
│ └──────────┘ │
└──────────────────────┬──────────────────┘
│
▼
┌─────────────────────┐
│ 完整生产系统 │
│(ProductionLLMSystem) │
└─────────────────────┘
评论区