目 录CONTENT

文章目录

LLM Ops 到底在做什么?从概念到落地全景图

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

一、为什么用 LLM Ops

传统 DevOps 解决的是"代码从开发到上线"的问题,核心链路:

但在大模型时代,这条链路变了。你交付的不是代码,而是一个概率性输出的模型服务。它有自己的问题:

  • 模型不是确定性的,同样的输入可能给出不同的输出

  • 模型有版本,但版本的"好坏"不像代码那样用单测就能判定

  • 模型调用的成本远高于普通 API,Token 是真金白银

  • 模型需要上下文管理、缓存、降级,传统网关搞不定

  • 模型的评测不是 pass/fail,而是"够不够好"

所以我们需要一套新的工程化体系——LLM Ops

二、LLM Ops 全景图

先上一张全景架构图,后面每个模块都会展开:

这 8 个模块构成了 LLM Ops 的完整闭环。下面逐一拆解。

三、数据管理

大模型的落地,数据是起点也是终点。训练需要数据,微调需要数据,评测需要数据,RAG 也需要数据。

3.1 数据集版本化

和代码一样,数据也需要版本管理。我们用 DVC(Data Version Control)来管理数据集:

# 初始化 DVC
dvc init

# 添加数据集
dvc add data/training_dataset_v1.jsonl

# 这会生成 .gitignore 和 .dvc 文件
git add data/training_dataset_v1.jsonl.dvc data/.gitignore
git commit -m "feat: add training dataset v1"

# 推送数据到远程存储(S3/OSS)
dvc remote add -d myremote s3://my-bucket/dvc-storage
dvc push

.dvc 文件长这样:

outs:
- md5: a3f2b8c9d1e4f5a6b7c8d9e0f1a2b3c4
  size: 52428800
  path: data/training_dataset_v1.jsonl

3.2 数据清洗 Pipeline

一个典型的数据清洗流水线:

import json
import re
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor


def clean_text(text: str) -> str:
    """清洗文本:去除特殊字符、多余空格、HTML标签"""
    text = re.sub(r'<[^>]+>', '', text)  # 去HTML标签
    text = re.sub(r'\s+', ' ', text)     # 合并空格
    text = text.strip()
    return text


def validate_sample(sample: dict) -> bool:
    """校验单条数据质量"""
    if not sample.get("prompt") or not sample.get("response"):
        return False
    if len(sample["prompt"]) < 10 or len(sample["response"]) < 20:
        return False
    if sample["prompt"].count("?") > 10:
        return False  # 疑似垃圾数据
    return True


def process_file(input_path: Path, output_path: Path):
    """处理单个数据文件"""
    valid_count = 0
    invalid_count = 0

    with open(input_path, 'r', encoding='utf-8') as fin, \
         open(output_path, 'w', encoding='utf-8') as fout:

        for line in fin:
            try:
                sample = json.loads(line.strip())
                sample["prompt"] = clean_text(sample["prompt"])
                sample["response"] = clean_text(sample["response"])

                if validate_sample(sample):
                    fout.write(json.dumps(sample, ensure_ascii=False) + '\n')
                    valid_count += 1
                else:
                    invalid_count += 1
            except (json.JSONDecodeError, KeyError):
                invalid_count += 1

    print(f"[{input_path.name}] valid: {valid_count}, invalid: {invalid_count}")


def run_pipeline(input_dir: str, output_dir: str):
    """并行执行数据清洗"""
    input_path = Path(input_dir)
    output_path = Path(output_dir)
    output_path.mkdir(parents=True, exist_ok=True)

    files = list(input_path.glob("*.jsonl"))
    with ThreadPoolExecutor(max_workers=4) as executor:
        for f in files:
            executor.submit(
                process_file,
                f,
                output_path / f.name
            )


if __name__ == "__main__":
    run_pipeline("data/raw", "data/cleaned")

3.3 数据血缘追踪

数据从哪来,经过了什么处理,用在了哪里——这些都需要追踪:

raw/web_crawl_20250501.jsonl
  → cleaned/cleaned_20250501.jsonl        (清洗)
  → filtered/filtered_20250501.jsonl      (质量过滤)
  → train_dataset_v3.jsonl                (合并多源)
  → fine-tune job #42                     (微调)
  → model my-llm:v3.2                     (产出模型)

我们用一个小工具来记录血缘关系:

from dataclasses import dataclass, field
from datetime import datetime
from typing import Optional
import hashlib
import json


@dataclass
class DataLineage:
    """数据血缘记录"""

    data_id: str
    source: str
    transformations: list = field(default_factory=list)
    parent_ids: list = field(default_factory=list)
    created_at: str = field(default_factory=lambda: datetime.now().isoformat())
    metadata: dict = field(default_factory=dict)

    def add_transform(self, name: str, params: dict = None):
        self.transformations.append(
            {
                "name": name,
                "params": params or {},
                "timestamp": datetime.now().isoformat(),
            }
        )

    def compute_hash(self, filepath: str) -> str:
        """计算文件哈希作为唯一ID"""
        h = hashlib.md5()
        with open(filepath, "rb") as f:
            for chunk in iter(lambda: f.read(8192), b""):
                h.update(chunk)
        return h.hexdigest()


# 使用示例
lineage = DataLineage(
    data_id="dataset_v3",
    source="web_crawl + manual_annotation",
    parent_ids=["raw_20250501", "raw_20250515"],
    metadata={"records": 50000, "language": "zh"},
)
lineage.add_transform("clean_text", {"remove_html": True})
lineage.add_transform("quality_filter", {"min_length": 20})
lineage.add_transform("deduplicate", {"method": "minhash"})

print(json.dumps(lineage.__dict__, indent=2, ensure_ascii=False))

四、模型管理

4.1 模型注册中心

类似 Docker Registry,模型也需要一个注册中心。我们用 MLflow:

import mlflow
from mlflow.models import infer_signature


def register_model(
    model_name: str,
    model_path: str,
    metrics: dict,
    tags: dict = None,
    description: str = "",
):
    """注册模型到 MLflow Model Registry"""

    mlflow.set_tracking_uri("http://mlflow-server:5000")
    mlflow.set_experiment(model_name)

    with mlflow.start_run() as run:
        # 记录指标
        mlflow.log_metrics(metrics)

        # 记录标签
        if tags:
            mlflow.set_tags(tags)

        # 记录模型
        mlflow.log_artifacts(model_path, "model")

        # 注册到 Model Registry
        result = mlflow.register_model(f"runs:/{run.info.run_id}/model", model_name)

        # 添加描述
        client = mlflow.tracking.MlflowClient()
        client.update_model_version(
            name=model_name, version=result.version, description=description
        )

        print(f"Model {model_name} v{result.version} registered")
        print(f"Run ID: {run.info.run_id}")

    return result


# 使用示例
register_model(
    model_name="customer-service-llm",
    model_path="./model_output",
    metrics={"accuracy": 0.92, "bleu_score": 0.85, "rouge_l": 0.88, "latency_p95": 1.2},
    tags={
        "base_model": "Qwen2.5-7B",
        "fine_tune_method": "LoRA",
        "task": "customer_service",
    },
    description="客户服务微调模型v3,基于Qwen2.5-7B,LoRA微调,对话质量提升12%",
)

4.2 模型版本策略

模型版本不像代码,不能简单用语义化版本。我们定义一套版本规则:

版本号格式: {base_model}_{major}.{minor}.{patch}_{lora_version}

示例:
  qwen2.5-7b_3.2.1_lora-v5
  ├── qwen2.5-7b  : 基座模型
  ├── 3            : 第3次完整微调(数据集变更)
  ├── 2            : 第2次增量微调(少量数据补充)
  ├── 1            : bugfix(修复特定case)
  └── lora-v5      : LoRA适配器版本

版本流转规则:

# model_version_policy.yaml
staging:
  - 自动评测通过(准确率 > 阈值)
  - A/B测试流量 5%
  - 人工抽检 100条

production:
  - staging 阶段运行 ≥ 3天
  - P95 延迟 < 2s
  - 用户负反馈率 < 2%
  - 成本在预算内

rollback:
  - 准确率下降 > 5%
  - P99 延迟 > 5s
  - 用户负反馈率 > 5%
  - 自动回滚到上一版本

4.3 模型评测

模型评测不是跑个 BLEU 就完了,需要多维度评估:

from dataclasses import dataclass
from typing import List
import json
import time


@dataclass
class EvalResult:
    """评测结果"""

    model_version: str
    dataset: str
    metrics: dict
    latency_stats: dict
    sample_results: list
    timestamp: str


class ModelEvaluator:
    """模型评测器"""

    def __init__(self, model_client, eval_dataset_path: str):
        self.model_client = model_client
        self.eval_dataset_path = eval_dataset_path
        self.eval_dataset = self._load_dataset(eval_dataset_path)

    def _load_dataset(self, path: str) -> list:
        dataset = []
        with open(path, "r", encoding="utf-8") as f:
            for line in f:
                dataset.append(json.loads(line.strip()))
        return dataset

    def evaluate(self, model_version: str) -> EvalResult:
        """执行评测"""
        latencies = []
        results = []
        correct = 0

        for sample in self.eval_dataset:
            start = time.time()
            response = self.model_client.chat(
                model=model_version,
                messages=[{"role": "user", "content": sample["prompt"]}],
                temperature=0.0,  # 评测时固定温度
            )
            latency = time.time() - start
            latencies.append(latency)

            # 自动评测
            is_correct = self._check_response(response, sample["expected"])
            if is_correct:
                correct += 1

            results.append(
                {
                    "prompt": sample["prompt"][:50],
                    "expected": sample["expected"][:50],
                    "actual": response[:50],
                    "correct": is_correct,
                    "latency": latency,
                }
            )

        # 统计指标
        latencies.sort()
        eval_result = EvalResult(
            model_version=model_version,
            dataset=self.eval_dataset_path,
            metrics={
                "accuracy": correct / len(self.eval_dataset),
                "total": len(self.eval_dataset),
                "correct": correct,
            },
            latency_stats={
                "p50": latencies[len(latencies) // 2],
                "p95": latencies[int(len(latencies) * 0.95)],
                "p99": latencies[int(len(latencies) * 0.99)],
                "avg": sum(latencies) / len(latencies),
            },
            sample_results=results[:20],  # 保留前20条详细结果
            timestamp=time.strftime("%Y-%m-%d %H:%M:%S"),
        )

        return eval_result

    def _check_response(self, response: str, expected: str) -> bool:
        """简单的响应匹配检查"""
        # 实际项目中会用 LLM-as-Judge 或更复杂的评测逻辑
        return expected.lower().strip() in response.lower().strip()

五、应用编排

模型训练好了,但怎么用?这就是应用编排要解决的问题。

从简单的 Prompt 到复杂的 Agent,从单轮对话到多步推理,应用编排是 LLM 落地的关键一环。

5.1 Prompt 工程与版本管理

Prompt 是 LLM 的"代码",也需要版本管理和 A/B 测试。我们用 LangSmith 来管理:

from langsmith import Client
from datetime import datetime
import json


class PromptManager:
    """Prompt 版本管理器"""

    def __init__(self, project_name: str):
        self.client = Client()
        self.project_name = project_name

    def create_prompt(
        self,
        name: str,
        template: str,
        variables: list[str],
        tags: list[str] = None,
        description: str = "",
    ) -> str:
        """创建 Prompt 模板"""

        prompt_data = {
            "name": name,
            "template": template,
            "variables": variables,
            "tags": tags or [],
            "description": description,
            "created_at": datetime.now().isoformat(),
            "project": self.project_name,
        }

        # 保存到 LangSmith
        prompt_id = self.client.create_prompt(
            prompt_name=name, prompt_template=template, tags=tags
        )

        print(f"✅ Prompt '{name}' created with ID: {prompt_id}")
        return prompt_id

    def get_prompt(self, name: str, version: str = "latest") -> dict:
        """获取指定版本的 Prompt"""
        return self.client.get_prompt(name, version=version)

    def compare_versions(self, name: str, v1: str, v2: str):
        """对比两个版本的效果"""
        prompt_v1 = self.get_prompt(name, v1)
        prompt_v2 = self.get_prompt(name, v2)

        # 获取评测数据
        metrics_v1 = self.client.get_prompt_metrics(name, v1)
        metrics_v2 = self.client.get_prompt_metrics(name, v2)

        return {
            "v1": {"prompt": prompt_v1, "metrics": metrics_v1},
            "v2": {"prompt": prompt_v2, "metrics": metrics_v2},
            "winner": v1 if metrics_v1["score"] > metrics_v2["score"] else v2,
        }


# 使用示例
pm = PromptManager(project_name="customer-service")

# 创建客服 Prompt
pm.create_prompt(
    name="customer-greeting",
    template="""
        你是一个专业的客服助手,名字叫小智。

        用户信息:
        - 姓名:{customer_name}
        - VIP等级:{vip_level}
        - 历史订单数:{order_count}
        - 最近一次咨询:{last_inquiry}

        请用友好、专业的语气回答用户问题。如果问题超出你的能力范围,请引导用户联系人工客服。

        用户问题:{question}

        你的回答:
    """,
    variables=["customer_name", "vip_level", "order_count", "last_inquiry", "question"],
    tags=["customer-service", "greeting", "v2.0"],
    description="客服问候模板 v2.0,增加了历史订单数和最近咨询信息",
)

5.2 Chain 编排:多步骤推理

使用 LangChain 构建多步骤推理链:

from langchain.chains import LLMChain, SequentialChain
from langchain.prompts import PromptTemplate
from langchain_openai import ChatOpenAI


class CustomerServiceChain:
    """客服多步骤处理链"""

    def __init__(self, llm):
        self.llm = llm
        self.chain = self._build_chain()

    def _build_chain(self) -> SequentialChain:
        """构建处理链"""

        # Step 1: 意图识别
        intent_prompt = PromptTemplate(
            input_variables=["question"],
            template="""
            分析用户问题的意图,从以下类别中选择一个:
- 订单查询
- 退换货
- 产品咨询
- 投诉建议
- 其他

用户问题:{question}

意图分类:""",
        )
        intent_chain = LLMChain(llm=self.llm, prompt=intent_prompt, output_key="intent")

        # Step 2: 信息提取
        extract_prompt = PromptTemplate(
            input_variables=["question", "intent"],
            template="""从用户问题中提取关键信息。

                意图:{intent}
                问题:{question}

                请提取以下信息(JSON格式):
                - 订单号(如果有)
                - 产品名称(如果有)
                - 时间信息(如果有)
                - 其他关键信息

                提取结果:
            """,
        )
        extract_chain = LLMChain(
            llm=self.llm, prompt=extract_prompt, output_key="extracted_info"
        )

        # Step 3: 生成回复
        response_prompt = PromptTemplate(
            input_variables=["question", "intent", "extracted_info"],
            template="""
                基于以下信息生成专业的客服回复:

                    用户问题:{question}
                    意图分类:{intent}
                    提取信息:{extracted_info}

                    要求:
                    1. 语气友好、专业
                    2. 如果信息不足,礼貌地询问补充信息
                    3. 提供明确的解决方案或后续步骤

                客服回复:
            """,
        )
        response_chain = LLMChain(
            llm=self.llm, prompt=response_prompt, output_key="response"
        )

        # 组合成顺序链
        overall_chain = SequentialChain(
            chains=[intent_chain, extract_chain, response_chain],
            input_variables=["question"],
            output_variables=["intent", "extracted_info", "response"],
            verbose=True,
        )

        return overall_chain

    def process(self, question: str) -> dict:
        """处理用户问题"""
        return self.chain({"question": question})


# 使用示例
llm = ChatOpenAI(model="gpt-4", temperature=0.7)
service_chain = CustomerServiceChain(llm)

result = service_chain.process("我的订单 #12345 已经3天了还没发货,怎么回事?")
print(f"意图:{result['intent']}")
print(f"提取信息:{result['extracted_info']}")
print(f"回复:{result['response']}")

5.3 Agent 编排:自主决策与工具调用

使用 LangGraph 构建具有自主决策能力的 Agent:

from typing import TypedDict, Annotated, Sequence
from langchain_core.messages import BaseMessage, HumanMessage, AIMessage
from langgraph.graph import StateGraph, END
from langgraph.prebuilt import ToolExecutor, ToolInvocation
import operator


# 定义状态
class AgentState(TypedDict):
    messages: Annotated[Sequence[BaseMessage], operator.add]
    next_action: str


# 定义工具
def search_order(order_id: str) -> dict:
    """查询订单信息"""
    # 实际项目中调用订单系统 API
    return {
        "order_id": order_id,
        "status": "已发货",
        "tracking_number": "SF1234567890",
        "estimated_delivery": "2025-06-03",
    }


def search_product(product_name: str) -> dict:
    """查询产品信息"""
    return {"name": product_name, "price": 299.00, "stock": 156, "rating": 4.8}


def create_ticket(issue: str, priority: str = "normal") -> dict:
    """创建工单"""
    return {
        "ticket_id": "TK20250601001",
        "status": "已创建",
        "priority": priority,
        "estimated_response": "2小时内",
    }


# 工具列表
tools = [search_order, search_product, create_ticket]
tool_executor = ToolExecutor(tools)


# 定义节点
def call_model(state: AgentState):
    """调用模型决策下一步"""
    messages = state["messages"]

    # 这里简化处理,实际项目中会调用 LLM 判断
    last_message = messages[-1].content

    if "订单" in last_message and "#" in last_message:
        return {
            "next_action": "search_order",
            "messages": [AIMessage(content="正在查询订单信息...")],
        }
    elif "产品" in last_message or "价格" in last_message:
        return {
            "next_action": "search_product",
            "messages": [AIMessage(content="正在查询产品信息...")],
        }
    elif "投诉" in last_message or "问题" in last_message:
        return {
            "next_action": "create_ticket",
            "messages": [AIMessage(content="正在为您创建工单...")],
        }
    else:
        return {
            "next_action": "respond",
            "messages": [AIMessage(content="我理解您的问题,让我为您处理...")],
        }


def execute_tool(state: AgentState):
    """执行工具调用"""
    action = state["next_action"]
    last_message = state["messages"][-2].content  # 用户消息

    if action == "search_order":
        # 提取订单号
        import re

        order_id = re.search(r"#(\w+)", last_message).group(1)
        result = search_order(order_id)
        message = f"订单 #{result['order_id']} 状态:{result['status']},快递单号:{result['tracking_number']},预计送达:{result['estimated_delivery']}"

    elif action == "search_product":
        result = search_product("示例产品")
        message = f"产品信息:{result['name']},价格:¥{result['price']},库存:{result['stock']}件,评分:{result['rating']}⭐"

    elif action == "create_ticket":
        result = create_ticket(last_message, priority="high")
        message = f"已为您创建工单 {result['ticket_id']},优先级:{result['priority']},预计响应时间:{result['estimated_response']}"

    else:
        message = "我会尽快为您处理这个问题。"

    return {"messages": [AIMessage(content=message)]}


def should_continue(state: AgentState) -> str:
    """判断是否继续"""
    action = state.get("next_action", "")
    if action in ["search_order", "search_product", "create_ticket"]:
        return "execute"
    else:
        return "end"


# 构建图
workflow = StateGraph(AgentState)

# 添加节点
workflow.add_node("agent", call_model)
workflow.add_node("execute", execute_tool)

# 设置入口
workflow.set_entry_point("agent")

# 添加边
workflow.add_conditional_edges(
    "agent", should_continue, {"execute": "execute", "end": END}
)
workflow.add_edge("execute", END)

# 编译
app = workflow.compile()


# 使用示例
def run_agent(user_input: str):
    """运行 Agent"""
    result = app.invoke(
        {"messages": [HumanMessage(content=user_input)], "next_action": ""}
    )

    for msg in result["messages"]:
        print(f"{msg.__class__.__name__}: {msg.content}")


# 测试
run_agent("我的订单 #12345 什么时候能到?")

5.4 RAG 架构:知识增强生成

构建企业级 RAG 系统:

from langchain.vectorstores import Chroma
from langchain.embeddings import OpenAIEmbeddings
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.document_loaders import DirectoryLoader
from langchain.chains import RetrievalQA
from langchain_openai import ChatOpenAI
import chromadb


class EnterpriseRAG:
    """企业级 RAG 系统"""

    def __init__(
        self,
        collection_name: str,
        embedding_model: str = "text-embedding-3-small",
        persist_directory: str = "./chroma_db",
    ):
        self.collection_name = collection_name
        self.embeddings = OpenAIEmbeddings(model=embedding_model)
        self.persist_directory = persist_directory

        # 初始化向量数据库
        self.vectorstore = Chroma(
            collection_name=collection_name,
            embedding_function=self.embeddings,
            persist_directory=persist_directory,
        )

    def ingest_documents(self, docs_path: str, chunk_size: int = 1000):
        """导入文档到向量数据库"""

        # 加载文档
        loader = DirectoryLoader(docs_path, glob="**/*.md")
        documents = loader.load()

        # 分块
        text_splitter = RecursiveCharacterTextSplitter(
            chunk_size=chunk_size,
            chunk_overlap=200,
            separators=["\n\n", "\n", "。", "!", "?", ".", "!", "?", " ", ""],
        )
        chunks = text_splitter.split_documents(documents)

        # 添加元数据
        for i, chunk in enumerate(chunks):
            chunk.metadata["chunk_id"] = i
            chunk.metadata["source_file"] = chunk.metadata.get("source", "unknown")

        # 存入向量数据库
        self.vectorstore.add_documents(chunks)
        self.vectorstore.persist()

        print(f"✅ 已导入 {len(documents)} 个文档,分割为 {len(chunks)} 个块")

    def query(
        self, question: str, top_k: int = 3, score_threshold: float = 0.7
    ) -> dict:
        """查询知识库"""

        # 1. 检索相关文档
        docs_with_scores = self.vectorstore.similarity_search_with_score(
            question, k=top_k
        )

        # 2. 过滤低分文档
        filtered_docs = [
            (doc, score) for doc, score in docs_with_scores if score >= score_threshold
        ]

        if not filtered_docs:
            return {
                "answer": "抱歉,我在知识库中没有找到相关信息。",
                "sources": [],
                "confidence": 0.0,
            }

        # 3. 构建上下文
        context_parts = []
        sources = []
        for doc, score in filtered_docs:
            context_parts.append(doc.page_content)
            sources.append(
                {
                    "file": doc.metadata.get("source_file", "unknown"),
                    "chunk_id": doc.metadata.get("chunk_id", -1),
                    "score": float(score),
                }
            )

        context = "\n\n---\n\n".join(context_parts)

        # 4. 生成回答
        llm = ChatOpenAI(model="gpt-4", temperature=0.3)

        prompt = f"""基于以下知识库内容回答问题。如果知识库中没有相关信息,请明确说明。

知识库内容:
{context}

用户问题:{question}

要求:
1. 回答要准确、专业
2. 如果有多个相关信息,请综合回答
3. 引用具体的知识库内容
4. 如果不确定,请说明

你的回答:"""

        answer = llm.invoke(prompt).content

        return {
            "answer": answer,
            "sources": sources,
            "confidence": sum(s["score"] for s in sources) / len(sources),
        }

    def hybrid_search(self, question: str, top_k: int = 5) -> list:
        """混合检索:向量检索 + 关键词检索"""

        # 向量检索
        vector_results = self.vectorstore.similarity_search(question, k=top_k)

        # 关键词检索(简化版,实际项目中用 BM25)
        import jieba

        keywords = jieba.lcut(question)
        keyword_results = []

        for keyword in keywords:
            results = self.vectorstore.similarity_search(keyword, k=2)
            keyword_results.extend(results)

        # 合并去重
        all_results = vector_results + keyword_results
        unique_results = list({doc.page_content: doc for doc in all_results}.values())

        return unique_results[:top_k]


# 使用示例
rag = EnterpriseRAG(collection_name="company_knowledge")

# 导入文档
rag.ingest_documents("./docs/knowledge_base")

# 查询
result = rag.query("公司的退货政策是什么?")
print(f"回答:{result['answer']}")
print(f"置信度:{result['confidence']:.2f}")
print(f"来源:{result['sources']}")

5.5 应用编排最佳实践

from dataclasses import dataclass
from enum import Enum


class OrchestrationPattern(Enum):
    """编排模式"""

    SIMPLE_PROMPT = "simple_prompt"  # 单次 Prompt
    CHAIN = "chain"  # 顺序链
    AGENT = "agent"  # 自主 Agent
    RAG = "rag"  # 知识增强
    HYBRID = "hybrid"  # 混合模式


@dataclass
class OrchestrationDecision:
    """编排决策"""

    pattern: OrchestrationPattern
    reason: str
    estimated_cost: float
    estimated_latency: float


def choose_orchestration(
    task_complexity: str,  # simple / medium / complex
    need_external_knowledge: bool,
    need_tools: bool,
    latency_requirement: float,  # seconds
) -> OrchestrationDecision:
    """选择合适的编排模式"""

    # 决策树
    if task_complexity == "simple" and not need_external_knowledge and not need_tools:
        return OrchestrationDecision(
            pattern=OrchestrationPattern.SIMPLE_PROMPT,
            reason="任务简单,单次调用即可",
            estimated_cost=0.002,
            estimated_latency=0.5,
        )

    if need_external_knowledge and not need_tools:
        return OrchestrationDecision(
            pattern=OrchestrationPattern.RAG,
            reason="需要外部知识,使用 RAG",
            estimated_cost=0.005,
            estimated_latency=1.2,
        )

    if task_complexity == "medium" and not need_tools:
        return OrchestrationDecision(
            pattern=OrchestrationPattern.CHAIN,
            reason="中等复杂度,使用多步骤链",
            estimated_cost=0.008,
            estimated_latency=2.0,
        )

    if need_tools or task_complexity == "complex":
        return OrchestrationDecision(
            pattern=OrchestrationPattern.AGENT,
            reason="需要工具调用或任务复杂,使用 Agent",
            estimated_cost=0.015,
            estimated_latency=3.5,
        )

    # 默认混合模式
    return OrchestrationDecision(
        pattern=OrchestrationPattern.HYBRID,
        reason="综合场景,使用混合模式",
        estimated_cost=0.012,
        estimated_latency=2.5,
    )


# 使用示例
decision = choose_orchestration(
    task_complexity="complex",
    need_external_knowledge=True,
    need_tools=True,
    latency_requirement=5.0,
)

print(f"推荐模式:{decision.pattern.value}")
print(f"原因:{decision.reason}")
print(f"预估成本:¥{decision.estimated_cost}")
print(f"预估延迟:{decision.estimated_latency}s")

六、推理服务

6.1 vLLM 部署

vLLM 是目前最主流的推理引擎,PagedAttention 让吞吐量翻倍:

# docker-compose-vllm.yaml
version: '3.8'

services:
  vllm-server:
    image: vllm/vllm-openai:latest
    container_name: vllm-server
    ports:
      - "8000:8000"
    command: >
      --model /models/qwen2.5-7b
      --served-model-name customer-service-llm
      --host 0.0.0.0
      --port 8000
      --tensor-parallel-size 2
      --gpu-memory-utilization 0.9
      --max-model-len 8192
      --dtype auto
      --enable-prefix-caching
      --trust-remote-code
    volumes:
      - ./models:/models
    deploy:
      resources:
        reservations:
          devices:
            - driver: nvidia
              count: 2
              capabilities: [gpu]
    environment:
      - CUDA_VISIBLE_DEVICES=0,1
    restart: unless-stopped
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
      interval: 30s
      timeout: 10s
      retries: 3

6.2 推理服务架构

一个完整的推理服务不只是 vLLM,还需要网关、缓存、降级:

6.3 推理路由与降级

from enum import Enum
from dataclasses import dataclass
from typing import Optional
import time
import random


class ModelTier(Enum):
    PRIMARY = "primary"  # 主力模型:Qwen2.5-72B
    SECONDARY = "secondary"  # 备选模型:Qwen2.5-7B
    FALLBACK = "fallback"  # 降级模型:Qwen2.5-1.5B


@dataclass
class ModelEndpoint:
    name: str
    tier: ModelTier
    url: str
    max_tokens: int
    cost_per_1k_tokens: float
    is_healthy: bool = True
    current_load: float = 0.0


class InferenceRouter:
    """推理路由器:根据场景选择合适的模型"""

    def __init__(self):
        self.endpoints = [
            ModelEndpoint(
                name="qwen-72b-primary",
                tier=ModelTier.PRIMARY,
                url="http://vllm-01:8000",
                max_tokens=8192,
                cost_per_1k_tokens=0.012,
            ),
            ModelEndpoint(
                name="qwen-7b-secondary",
                tier=ModelTier.SECONDARY,
                url="http://vllm-02:8000",
                max_tokens=4096,
                cost_per_1k_tokens=0.002,
            ),
            ModelEndpoint(
                name="qwen-1.5b-fallback",
                tier=ModelTier.FALLBACK,
                url="http://vllm-fallback:8000",
                max_tokens=2048,
                cost_per_1k_tokens=0.0005,
            ),
        ]

    def route(
        self,
        prompt: str,
        max_tokens: int = 2048,
        priority: str = "balanced",  # quality / cost / balanced
        enable_fallback: bool = True,
    ) -> ModelEndpoint:
        """路由决策"""

        # 过滤不可用的端点
        available = [ep for ep in self.endpoints if ep.is_healthy]

        if not available:
            raise RuntimeError("No healthy model endpoints available")

        # 根据优先级排序
        if priority == "quality":
            # 质量优先:从大到小
            candidates = sorted(
                available, key=lambda x: x.cost_per_1k_tokens, reverse=True
            )
        elif priority == "cost":
            # 成本优先:从小到大
            candidates = sorted(available, key=lambda x: x.cost_per_1k_tokens)
        else:
            # 均衡:考虑负载
            candidates = sorted(available, key=lambda x: x.current_load)

        # 选择第一个满足 max_tokens 要求的
        for ep in candidates:
            if ep.max_tokens >= max_tokens:
                return ep

        # 降级:如果都满足不了,选最大的
        if enable_fallback:
            return max(available, key=lambda x: x.max_tokens)

        raise RuntimeError(f"No endpoint supports max_tokens={max_tokens}")

七、可观测性

7.1 三大支柱

┌─────────────────────────────────────────────────────────┐
│                  LLM 可观测性三大支柱                     │
├─────────────────────────────────────────────────────────┤
│                                                         │
│  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐    │
│  │    指标      │  │    日志      │  │    追踪      │    │
│  │  Metrics    │  │    Logs     │  │   Traces    │    │
│  │             │  │             │  │             │    │
│  │ ·QPS        │  │ ·请求日志   │  │ ·调用链路   │    │
│  │ ·延迟P50/95 │  │ ·响应内容   │  │ ·Token分布  │    │
│  │ ·Token用量  │  │ ·错误详情   │  │ ·模型版本   │    │
│  │ ·成本       │  │ ·审核日志   │  │ ·缓存命中   │    │
│  │ ·缓存命中率 │  │             │  │             │    │
│  └─────────────┘  └─────────────┘  └─────────────┘    │
│                                                         │
│  ┌─────────────────────────────────────────────────┐   │
│  │              告警 & 看板                          │   │
│  │  ·Grafana Dashboard                              │   │
│  │  ·企业微信/飞书告警                               │   │
│  │  ·成本日报/周报                                   │   │
│  └─────────────────────────────────────────────────┘   │
└─────────────────────────────────────────────────────────┘

7.2 核心指标采集

from prometheus_client import Counter, Histogram, Gauge, start_http_server
import time


# 定义指标
LLM_REQUEST_TOTAL = Counter(
    "llm_request_total", "Total LLM API requests", ["model", "status", "endpoint"]
)

LLM_REQUEST_DURATION = Histogram(
    "llm_request_duration_seconds",
    "LLM request duration in seconds",
    ["model"],
    buckets=[0.5, 1.0, 2.0, 3.0, 5.0, 10.0, 30.0],
)

LLM_TOKEN_USAGE = Counter(
    "llm_token_usage_total",
    "Total tokens used",
    ["model", "token_type"],  # token_type: input / output
)

LLM_COST_TOTAL = Counter(
    "llm_cost_total_yuan", "Total cost in yuan", ["model", "project"]
)

LLM_CACHE_HIT = Counter(
    "llm_cache_hit_total", "Cache hit count", ["cache_type"]  # exact / semantic
)

LLM_ACTIVE_REQUESTS = Gauge(
    "llm_active_requests", "Currently active requests", ["model"]
)


class LLMMetricsMiddleware:
    """LLM 指标中间件"""

    def __init__(self, model_name: str, project: str = "default"):
        self.model_name = model_name
        self.project = project

    def __call__(self, func):
        def wrapper(*args, **kwargs):
            LLM_ACTIVE_REQUESTS.labels(model=self.model_name).inc()
            start_time = time.time()

            try:
                result = func(*args, **kwargs)

                # 记录成功请求
                LLM_REQUEST_TOTAL.labels(
                    model=self.model_name, status="success", endpoint="chat"
                ).inc()

                # 记录 Token 用量
                if hasattr(result, "usage"):
                    LLM_TOKEN_USAGE.labels(
                        model=self.model_name, token_type="input"
                    ).inc(result.usage.prompt_tokens)

                    LLM_TOKEN_USAGE.labels(
                        model=self.model_name, token_type="output"
                    ).inc(result.usage.completion_tokens)

                    # 记录成本
                    cost = self._calculate_cost(result.usage)
                    LLM_COST_TOTAL.labels(
                        model=self.model_name, project=self.project
                    ).inc(cost)

                return result

            except Exception as e:
                LLM_REQUEST_TOTAL.labels(
                    model=self.model_name, status="error", endpoint="chat"
                ).inc()
                raise

            finally:
                duration = time.time() - start_time
                LLM_REQUEST_DURATION.labels(model=self.model_name).observe(duration)
                LLM_ACTIVE_REQUESTS.labels(model=self.model_name).dec()

        return wrapper

    def _calculate_cost(self, usage) -> float:
        """计算单次调用成本(元)"""
        # 不同模型不同价格,这里示例
        pricing = {
            "qwen2.5-72b": {"input": 0.004, "output": 0.012},
            "qwen2.5-7b": {"input": 0.001, "output": 0.002},
            "deepseek-v3": {"input": 0.002, "output": 0.006},
        }
        p = pricing.get(self.model_name, {"input": 0.002, "output": 0.006})
        input_cost = (usage.prompt_tokens / 1000) * p["input"]
        output_cost = (usage.completion_tokens / 1000) * p["output"]
        return input_cost + output_cost


# 启动 metrics server
start_http_server(9090)

7.3 Grafana Dashboard 配置

核心看板指标:

{
  "dashboard": {
    "title": "LLM Ops Monitor",
    "panels": [
      {
        "title": "QPS by Model",
        "type": "timeseries",
        "targets": [
          {
            "expr": "rate(llm_request_total[5m])"
          }
        ]
      },
      {
        "title": "P95 Latency",
        "type": "stat",
        "targets": [
          {
            "expr": "histogram_quantile(0.95, rate(llm_request_duration_seconds_bucket[5m]))"
          }
        ]
      },
      {
        "title": "Token Usage (24h)",
        "type": "piechart",
        "targets": [
          {
            "expr": "sum(llm_token_usage_total) by (token_type)"
          }
        ]
      },
      {
        "title": "Daily Cost Trend",
        "type": "timeseries",
        "targets": [
          {
            "expr": "sum(increase(llm_cost_total_yuan[1d])) by (model)"
          }
        ]
      },
      {
        "title": "Error Rate",
        "type": "gauge",
        "targets": [
          {
            "expr": "rate(llm_request_total{status=\"error\"}[5m]) / rate(llm_request_total[5m])"
          }
        ],
        "thresholds": {
          "steps": [
            { "value": 0, "color": "green" },
            { "value": 0.05, "color": "yellow" },
            { "value": 0.1, "color": "red" }
          ]
        }
      }
    ]
  }
}

八、成本管理

8.1 成本监控告警

# cost_alert_rules.yaml
groups:
  - name: llm_cost_alerts
    rules:
      # 单日成本超限
      - alert: DailyCostExceeded
        expr: sum(increase(llm_cost_total_yuan[24h])) > 5000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "日成本超限"
          description: "过去24小时LLM调用成本超过5000元"

      # 单个模型成本异常
      - alert: ModelCostAnomaly
        expr: |
          (sum(increase(llm_cost_total_yuan[1h])) by (model)
           / sum(increase(llm_cost_total_yuan[1h] offset 24h)) by (model))
           > 2.0
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: "模型成本同比异常"
          description: "{{ $labels.model }} 过去1小时成本是昨日同期的2倍以上"

      # 单用户配额超限
      - alert: UserQuotaExceeded
        expr: sum(increase(llm_cost_total_yuan[24h])) by (project) > 1000
        for: 5m
        labels:
          severity: info
        annotations:
          summary: "项目配额超限"
          description: "项目 {{ $labels.project }} 24小时内成本超过1000元"

8.2 成本优化策略

策略

原理

预期节省

语义缓存

相似问题直接返回缓存结果

30-50%

模型路由

简单问题用小模型,复杂问题用大模型

40-60%

Prompt 压缩

压缩 system prompt 减少 input token

15-25%

批量推理

合并请求批量处理

20-30%

流式截断

检测到废话提前截断

10-15%

九、CI/CD

9.1 LLM 流水线架构

9.2 GitLab CI 配置

# .gitlab-ci.yml - LLM 模型发布流水线

stages:
  - data_validation
  - fine_tune
  - evaluation
  - staging_deploy
  - production_deploy

variables:
  MODEL_NAME: "customer-service-llm"
  BASE_MODEL: "Qwen/Qwen2.5-7B-Instruct"
  MLFLOW_TRACKING_URI: "http://mlflow-server:5000"

# 阶段1:数据校验
data_validation:
  stage: data_validation
  image: python:3.11
  script:
    - pip install pandas pydantic
    - python scripts/validate_dataset.py --input data/train.jsonl
    - python scripts/data_quality_score.py --threshold 0.85
  rules:
    - changes:
        - data/**

# 阶段2:微调训练
fine_tune:
  stage: fine_tune
  image: nvidia/cuda:12.1.0-runtime-ubuntu22.04
  script:
    - pip install torch transformers peft trl datasets
    - python scripts/fine_tune.py
      --base_model $BASE_MODEL
      --dataset data/train.jsonl
      --method lora
      --output_dir ./model_output
      --epochs 3
      --learning_rate 2e-4
    - python scripts/log_to_mlflow.py --model_path ./model_output
  artifacts:
    paths:
      - ./model_output/
    expire_in: 7 days
  rules:
    - when: manual  # 训练需要手动触发
  tags:
    - gpu-a100

# 阶段3:自动评测
evaluation:
  stage: evaluation
  image: python:3.11
  script:
    - pip install mlflow openai
    - python scripts/evaluate_model.py
      --model_path ./model_output
      --eval_dataset data/eval.jsonl
      --metrics accuracy,bleu,rouge_l
      --threshold accuracy:0.90
    - python scripts/latency_benchmark.py
      --model_path ./model_output
      --max_p95_latency 2.0
  needs:
    - fine_tune
  rules:
    - when: on_success

# 阶段4:Staging 部署
staging_deploy:
  stage: staging_deploy
  image: bitnami/kubectl
  script:
    - kubectl config use-context staging
    - kubectl set image deployment/vllm-server
        vllm=$MODEL_NAME:$CI_COMMIT_SHORT_SHA
        -n llm-staging
    - kubectl rollout status deployment/vllm-server -n llm-staging
    - python scripts/smoke_test.py --env staging
  needs:
    - evaluation
  environment:
    name: staging
    url: https://llm-staging.internal.company.com
  rules:
    - when: on_success

# 阶段5:Production 部署
production_deploy:
  stage: production_deploy
  image: bitnami/kubectl
  script:
    # 灰度发布:先切5%流量
    - kubectl config use-context production
    - kubectl apply -f k8s/canary-deployment.yaml
    - sleep 300  # 观察5分钟
    - python scripts/canary_metrics_check.py --threshold 0.02
    # 全量发布
    - kubectl apply -f k8s/production-deployment.yaml
    - kubectl rollout status deployment/vllm-server -n llm-production
  needs:
    - staging_deploy
  environment:
    name: production
    url: https://llm.internal.company.com
  rules:
    - when: manual  # 生产发布必须手动确认
  only:
    - main

十、安全合规

10.1 内容审核

大模型的输出不可控,必须加审核层:

from enum import Enum
from dataclasses import dataclass
from typing import Optional


class RiskLevel(Enum):
    SAFE = "safe"
    LOW = "low"
    MEDIUM = "medium"
    HIGH = "high"
    BLOCKED = "blocked"


@dataclass
class ModerationResult:
    risk_level: RiskLevel
    categories: list
    confidence: float
    reason: str
    filtered_content: Optional[str] = None


class ContentModerator:
    """内容审核器"""

    # 敏感词列表(实际项目中从数据库/配置中心加载)
    SENSITIVE_PATTERNS = {
        "politics": ["政治敏感词1", "政治敏感词2"],
        "violence": ["暴力关键词1", "暴力关键词2"],
        "privacy": ["身份证", "银行卡号", "手机号"],
        "competitor": ["竞品名称1", "竞品名称2"],
    }

    def moderate_input(self, text: str) -> ModerationResult:
        """审核用户输入"""
        # 1. 正则规则引擎
        for category, patterns in self.SENSITIVE_PATTERNS.items():
            for pattern in patterns:
                if pattern in text:
                    return ModerationResult(
                        risk_level=RiskLevel.HIGH,
                        categories=[category],
                        confidence=0.95,
                        reason=f"匹配敏感规则: {category}",
                    )

        # 2. PII 检测
        import re

        pii_patterns = {
            "phone": r"1[3-9]\d{9}",
            "id_card": r"\d{17}[\dXx]",
            "bank_card": r"\d{16,19}",
        }
        for pii_type, pattern in pii_patterns.items():
            if re.search(pattern, text):
                return ModerationResult(
                    risk_level=RiskLevel.MEDIUM,
                    categories=["pii"],
                    confidence=0.90,
                    reason=f"检测到PII信息: {pii_type}",
                )

        # 3. 调用审核模型(如通义内容安全API)
        # result = call_moderation_api(text)

        return ModerationResult(
            risk_level=RiskLevel.SAFE,
            categories=[],
            confidence=0.99,
            reason="通过所有审核规则",
        )

    def moderate_output(self, text: str) -> ModerationResult:
        """审核模型输出"""
        # 输出审核逻辑类似,但关注点不同
        # 主要检查:幻觉、偏见、不当内容
        return self.moderate_input(text)

10.2 数据脱敏

import re
from typing import Callable


class DataMasker:
    """数据脱敏器"""

    MASKERS: dict[str, Callable] = {}

    @classmethod
    def register(cls, name: str):
        """注册脱敏规则"""

        def decorator(func):
            cls.MASKERS[name] = func
            return func

        return decorator

    @classmethod
    def mask(cls, text: str, rules: list[str] = None) -> str:
        """执行脱敏"""
        rules = rules or list(cls.MASKERS.keys())
        for rule in rules:
            if rule in cls.MASKERS:
                text = cls.MASKERS[rule](text)
        return text


@DataMasker.register("phone")
def mask_phone(text: str) -> str:
    """手机号脱敏: 13812345678 → 138****5678"""
    return re.sub(r"(1[3-9]\d)\d{4}(\d{4})", r"\1****\2", text)


@DataMasker.register("email")
def mask_email(text: str) -> str:
    """邮箱脱敏: test@example.com → t***@example.com"""
    return re.sub(r"(\w)\w+(@\S+)", r"\1***\2", text)


@DataMasker.register("id_card")
def mask_id_card(text: str) -> str:
    """身份证脱敏: 310101199001011234 → 3101**********1234"""
    return re.sub(r"(\d{4})\d{10}(\d{4})", r"\1**********\2", text)


# 使用示例
raw_text = "用户手机号13812345678,邮箱test@example.com,身份证310101199001011234"
masked = DataMasker.mask(raw_text)
print(masked)
# 输出: 用户手机号138****5678,邮箱t***@example.com,身份证3101**********1234

十一、总结

最后,给一个 LLM Ops 的成熟度模型,方便你评估自己团队在哪个阶段:

Level 0 ── 无意识
  │  · 直接调用 API,没有任何工程化
  │  · 没有版本管理,没有监控
  │
Level 1 ── 规范化
  │  · 有基本的 Prompt 管理
  │  · 有简单的监控和告警
  │  · 有成本追踪
  │
Level 2 ── 自动化
  │  · 自动评测流水线
  │  · CI/CD for 模型
  │  · 语义缓存、模型路由
  │
Level 3 ── 平台化
  │  · 统一 AI 平台
  │  · 多租户、资源隔离
  │  · 自助式模型服务
  │
Level 4 ── 智能化
  │  · 自动模型选择与调优
  │  · 自适应缓存策略
  │  · AIOps for LLM
  │
Level 5 ── 生态化
     · 模型市场
     · 跨团队协作
     · 产业级 AI 供应链

大部分团队在 Level 0-1,头部团队在 Level 2-3,Level 4 以上目前几乎没有。

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区