目 录CONTENT

文章目录

LLM - LLMOps Deploy

PySuper
2025-05-31 / 0 评论 / 0 点赞 / 4 阅读 / 0 字
温馨提示:
所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

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 endpoint

2.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
            raise

2.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_time

3.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 span

3.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_price

4.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 result

5.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 future

5.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 v

6.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_BACK

6.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.IDLE

8.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_alerts

9.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) │
              └─────────────────────┘

概念间的依赖关系

概念

核心职责

依赖概念

模型服务部署

提供模型 API 服务

流量路由与降级

高可用流量管理

模型服务部署

链路追踪

请求全链路可观测

模型服务部署

Token 与成本监控

成本追踪与预算

链路追踪

延迟与性能优化

降低延迟提升性能

Token 与成本监控

Prompt 版本管理

Prompt 配置管理

延迟与性能优化

CI/CD 流水线

自动化交付管道

Prompt 版本管理

灰度发布与回滚

安全上线与回退

CI/CD 流水线

质量监控告警

线上质量保障

灰度发布与回滚

完整生产系统

端到端集成

所有上述概念

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区