云计算百科
云计算领域专业知识百科平台

第40章:【高级篇综合实战】从源码到生产打造可观测高可用 RAGFlow 平台

1 项目背景

业务场景

「云帆科技」的 RAGFlow 平台经过一年的运营,已从单机房 Docker Compose 演进到多机房 K8s 集群,服务 20 个业务部门、2000 名员工、20 万份文档、1000 万 Chunk。随着规模增长和架构复杂度提升,团队面临一系列"平台级"挑战:

  • 百万级文档的检索延迟从 500ms 涨到 3 秒——需要对 Dealer 检索链路做源码级优化。
  • 自定义扩展散落各处:5 个部门的 WPS Parser、Jira Tool、审批组件缺乏统一管理——升级 RAGFlow 时兼容性未知。
  • 跨机房部署的数据同步:北京和上海两个机房,用户就近访问但数据需要保持一致。
  • SRE 落地:需要 SLO 定义、容量规划、灰度发布、备份恢复、成本监控、事故演练——全套运维体系。
  • CTO 要求:将过去 39 章的知识融会贯通,交付一个"从源码到生产"的完整平台方案,包含架构图、源码扩展、部署清单、监控大盘、压测报告和事故预案。

    痛点

    平台级交付 vs 单点优化的核心差异:

  • 单点优化 vs 全链路最优:优化了检索但忽略了 LLM 调用,体验改善有限。
  • 一次性部署 vs 持续演进:平台不是"搭好就完了"——需要灰度发布、回滚、A/B 测试能力。
  • 个人英雄 vs 团队 SOP:出问题不能再靠"小李看一眼"——需要监控、告警、Runbook。
  • 功能可用 vs 品质可量化:需要 SLO(P95<8s, 可用率>99.9%),而不只是"感觉还行"。
  • 2 项目设计

    小胖:(把 39 章的笔记铺了一桌子)“大师!39 章学完了,每一章的知识点我都理解,但怎么把这些东西串成一个生产级平台?我感觉自己像攒了一堆乐高积木,但不知道怎么拼成一座城堡。”

    大师:“这就是高级篇综合实战的意义——把前面所有的知识点编织成一张完整的网。我们按’平台交付六大领域’来组织:”

    平台交付六大领域:

    领域1:架构设计
    – 多机房部署拓扑图
    – 流量路由与就近接入
    – 数据同步策略(MySQL 主从、MinIO 跨区域复制)

    领域2:源码扩展
    – 自定义 Parser(WPS 格式)
    – 自定义 Tool(Jira 工单集成)
    – 检索耗时埋点与分段监控

    领域3:可观测性
    – Prometheus + Grafana 租户级仪表盘
    – Jaeger 全链路追踪
    – 5 条核心告警 + 分级通知

    领域4:SRE 落地
    – SLO 定义与监控(P95<8s / 可用率 99.9%)
    – 容量规划模型
    – 灰度发布 + 金丝雀流量

    领域5:运维治理
    – 备份恢复(MySQL binlog + MinIO mirror)
    – 事故预案与演练
    – 成本监控(LLM 费用 / 基础设施)

    领域6:验收标准
    – 百万级文档 P95 检索 < 500ms
    – 千万级 Chunk 索引不退化
    – 核心链路 100% 可追踪
    – 解析失败率 < 1%

    小胖:“多机房部署——北京和上海两个机房——文档存哪?用户怎么就近访问?”

    大师:“多机房的关键策略是’读写分离 + 就近接入’:”

    多机房架构:

    用户请求
    │
    ┌──────┴──────┐
    │ GeoDNS │ (根据来源 IP 路由到最近机房)
    └──────┬──────┘
    ┌───────────┴───────────┐
    ▼ ▼
    ┌──────────────────┐ ┌──────────────────┐
    │ 北京机房 (主) │ │ 上海机房 (从) │
    │ │ │ │
    │ MySQL (主, 读写) │◀──▶│ MySQL (从, 只读) │ binlog 异步复制
    │ MinIO (主) │◀──▶│ MinIO (从) │ mirror 同步
    │ Infinity (主) │ │ Infinity (只读) │ 索引独立
    │ API Server ×3 │ │ API Server ×3 │
    │ Task Executor ×2 │ │ Task Executor ×1 │ (解析仅在北京执行)
    └──────────────────┘ └──────────────────┘

    策略:
    – 写操作(上传文档、创建数据集)→ 全部路由到北京
    – 读操作(检索、问答)→ 就近机房处理
    – 文档解析 → 仅在北京(Task Executor 在上海部署但不消费队列)

    小白:“那如何保证核心链路 100% 可追踪?”

    大师:“TraceID 贯穿整个请求链路——从 API 入口到 LLM 返回:”

    # 平台级 TraceID 注入——覆盖所有核心链路
    import uuid
    from contextvars import ContextVar

    trace_id_var = ContextVar("trace_id", default="no-trace")

    class PlatformTracer:
    """平台级全链路追踪"""

    @classmethod
    def inject_trace_id(cls, request):
    """在 API 入口注入 TraceID"""
    trace_id = request.headers.get("X-Trace-ID") or str(uuid.uuid4())[:12]
    trace_id_var.set(trace_id)

    # 注入到日志
    logger.info(f"[trace={trace_id}] {request.method} {request.path}")

    return trace_id

    @classmethod
    def trace_span(cls, span_name, **attrs):
    """创建一个追踪 Span,自动携带 trace_id"""
    trace_id = trace_id_var.get()
    span = tracer.start_span(span_name)
    span.set_attribute("trace_id", trace_id)
    for k, v in attrs.items():
    span.set_attribute(k, str(v)[:200])
    return span

    # 在关键节点埋点
    # 1. API 入口
    # 2. Dealer 检索(向量搜索、关键词搜索、RRF 融合)
    # 3. Rerank 调用
    # 4. LLM 调用
    # 5. 文档解析(仅 Task Executor)

    3 项目实战

    环境准备

    目标:基于前 39 章知识,在 K8s 集群上交付完整的"可观测高可用 RAGFlow 平台"。

    前提:K8s 集群(至少 6 节点)、Prometheus + Grafana + Jaeger + Loki 可观测性套件已部署。

    分步实现

    步骤1:多机房 K8s 部署架构

    # platform-production.yaml – 生产级部署清单

    # 北京机房:主集群
    —
    apiVersion: apps/v1
    kind: Deployment
    metadata:
    name: ragflow–api–beijing
    namespace: ragflow–prod
    spec:
    replicas: 5
    template:
    spec:
    affinity:
    podAntiAffinity:
    requiredDuringSchedulingIgnoredDuringExecution:
    – labelSelector:
    matchLabels:
    app: ragflow–api
    topologyKey: kubernetes.io/hostname # 分散到不同物理节点
    containers:
    – name: api
    image: infiniflow/ragflow:v0.27.0
    env:
    – name: DEPLOY_REGION
    value: "beijing"
    – name: MYSQL_HOST
    value: mysql–primary.beijing.svc.cluster.local
    – name: TRACE_ENABLED
    value: "true"
    resources:
    requests: {memory: "2Gi", cpu: "1"}
    limits: {memory: "4Gi", cpu: "2"}

    —
    # 上海机房:从集群(只读)
    apiVersion: apps/v1
    kind: Deployment
    metadata:
    name: ragflow–api–shanghai
    namespace: ragflow–prod
    spec:
    replicas: 3
    template:
    spec:
    containers:
    – name: api
    image: infiniflow/ragflow:v0.27.0
    env:
    – name: DEPLOY_REGION
    value: "shanghai"
    – name: MYSQL_HOST
    value: mysql–replica.shanghai.svc.cluster.local
    – name: READ_ONLY_MODE
    value: "true" # 只读模式(禁止写操作)
    – name: TRACE_ENABLED
    value: "true"

    # GeoDNS 配置(Cloudflare / Route53 示例)
    # 北京用户 → 解析到 beijing.ragflow.yunfan.com
    # 上海用户 → 解析到 shanghai.ragflow.yunfan.com
    # 其他用户 → 就近或轮询

    步骤2:源码扩展集成——检索耗时埋点

    # rag/nlp/search.py 增加检索耗时埋点(整合第36章知识)

    class Dealer:
    async def search_with_trace(self, question, **kwargs):
    """带全链路追踪的检索"""
    with PlatformTracer.trace_span("dealer.search",
    question=question[:100],
    dataset_count=len(self.dataset_ids)) as span:

    # 子 Span: 向量检索
    with PlatformTracer.trace_span("dealer.vector_search") as vs_span:
    t0 = time.time()
    vector_results = await self.vector_search(query_vector)
    vs_span.set_attribute("results_count", len(vector_results))
    vs_span.set_attribute("latency_ms", (time.time() – t0) * 1000)

    # 子 Span: 关键词检索
    with PlatformTracer.trace_span("dealer.keyword_search") as ks_span:
    t0 = time.time()
    keyword_results = await self.keyword_search(question)
    ks_span.set_attribute("results_count", len(keyword_results))
    ks_span.set_attribute("latency_ms", (time.time() – t0) * 1000)

    # 子 Span: RRF 融合
    with PlatformTracer.trace_span("dealer.fusion"):
    fused = self.fusion(vector_results, keyword_results)

    # 子 Span: Rerank
    if self.config.get("rerank_model"):
    with PlatformTracer.trace_span("dealer.rerank") as r_span:
    t0 = time.time()
    fused = await self.rerank(question, fused)
    r_span.set_attribute("latency_ms", (time.time() – t0) * 1000)

    span.set_attribute("total_results", len(fused))
    return fused[:self.config.get("top_k", 10)]

    步骤3:Grafana 平台级监控大屏

    {
    "dashboard": {
    "title": "RAGFlow 平台总览 – 生产",
    "rows": [
    {
    "title": "全平台核心指标",
    "panels": [
    {
    "title": "总 QPS (按机房)",
    "targets": [
    {"expr": "sum(rate(chat_requests_total{region='beijing'}[1m]))", "legendFormat": "北京"},
    {"expr": "sum(rate(chat_requests_total{region='shanghai'}[1m]))", "legendFormat": "上海"}
    ]
    },
    {
    "title": "P95 延迟 (按机房)",
    "targets": [
    {"expr": "histogram_quantile(0.95, chat_latency_seconds_bucket{region='beijing'})", "legendFormat": "北京"},
    {"expr": "histogram_quantile(0.95, chat_latency_seconds_bucket{region='shanghai'})", "legendFormat": "上海"}
    ]
    },
    {
    "title": "LLM 错误率 (按供应商)",
    "targets": [
    {"expr": "rate(llm_errors_total{factory='YunFan'}[5m]) / rate(llm_calls_total{factory='YunFan'}[5m])"},
    {"expr": "rate(llm_errors_total{factory='OpenAI'}[5m]) / rate(llm_calls_total{factory='OpenAI'}[5m])"}
    ]
    },
    {
    "title": "租户 SLO 达标率",
    "targets": [
    {"expr": "sum(tenant_slo_compliant) / count(tenant_slo_compliant)"}
    ],
    "thresholds": [{"value": 0.95, "color": "green"}, {"value": 0.9, "color": "yellow"}, {"value": 0, "color": "red"}]
    }
    ]
    },
    {
    "title": "成本中心",
    "panels": [
    {
    "title": "LLM 费用趋势 (日)",
    "targets": [{"expr": "llm_cost_daily_total"}]
    },
    {
    "title": "基础设施成本 (月)",
    "targets": [{"expr": "kube_cost_total"}]
    }
    ]
    }
    ]
    }
    }

    步骤4:灰度发布与金丝雀部署

    # canary-deployment.yaml
    # 金丝雀发布:5% 流量走新版本,95% 走旧版本

    apiVersion: v1
    kind: Service
    metadata:
    name: ragflow–api–canary
    spec:
    selector:
    app: ragflow–api
    version: canary

    —
    apiVersion: networking.k8s.io/v1
    kind: Ingress
    metadata:
    name: ragflow–canary
    annotations:
    nginx.ingress.kubernetes.io/canary: "true"
    nginx.ingress.kubernetes.io/canary-weight: "5" # 5% 流量
    spec:
    rules:
    – host: ragflow.yunfan.com
    http:
    paths:
    – backend:
    service:
    name: ragflow–api–canary
    port:
    number: 9380

    # 灰度发布流程
    # Step 1: 部署金丝雀版本
    kubectl apply -f canary-deployment.yaml

    # Step 2: 观察金丝雀指标(错误率、延迟)
    # 若金丝雀版本错误率 > 2%,立即回滚

    # Step 3: 逐步提高金丝雀权重 5% → 20% → 50% → 100%
    kubectl annotate ingress ragflow-canary \\
    nginx.ingress.kubernetes.io/canary-weight="20"

    # Step 4: 全量切换后将金丝雀版本提升为稳定版
    kubectl set image deployment/ragflow-api-stable \\
    api=infiniflow/ragflow:v0.27.0

    # Step 5: 删除金丝雀部署
    kubectl delete -f canary-deployment.yaml

    步骤5:备份恢复与事故预案

    # disaster_recovery.sh – 灾难恢复脚本
    #!/bin/bash
    echo "=== RAGFlow 灾难恢复预案 ==="

    # === 备份策略 ===
    # MySQL: 每日全量 + binlog 实时备份
    mysqldump -h mysql-primary –all-databases –single-transaction \\
    | gzip > /backup/mysql/ragflow_$(date +%Y%m%d).sql.gz

    # MinIO: 实时 mirror 到备用节点
    mc mirror –watch ragflow-minio/ragflow backup-minio/ragflow-backup &

    # Infinity/ES: 每日快照
    curl -X PUT "http://infinity:23820/db/ragflow/_snapshot/daily_$(date +%Y%m%d)"

    # === 恢复流程 ===
    # 场景1: MySQL 数据损坏
    # 1. 停止 API Server (防止写入脏数据)
    # 2. 从最近备份恢复
    # zcat /backup/mysql/ragflow_20250615.sql.gz | mysql -h mysql-primary
    # 3. 应用 binlog 到故障前一刻
    # 4. 重启 API Server

    # 场景2: 整个北京机房宕机
    # 1. GeoDNS 将所有流量切换到上海机房
    # 2. 上海 MySQL 从库提升为主库
    # mysql -h mysql-replica.shanghai -e "STOP SLAVE; RESET SLAVE ALL;"
    # 3. 上海 Task Executor 开始消费 Redis 队列
    # 4. 通知用户服务降级(仅读可用,写操作暂不可用)

    # 场景3: LLM API 全部不可用
    # 1. 自动切换到本地 Ollama 模型(已在 LLMBundle 降级链中配置)
    # 2. 前端展示"当前使用备用模型,回答质量可能下降"
    # 3. 主模型恢复后自动切回

    echo "预案检查完成"

    步骤6:事故演练——混沌工程

    # chaos_test.py – 混沌工程测试
    import random
    import subprocess
    import time

    class ChaosTest:
    """RAGFlow 平台混沌工程测试"""

    @classmethod
    def kill_random_api_pod(cls):
    """随机杀死一个 API Server Pod——验证自动恢复"""
    pods = subprocess.run(
    ["kubectl", "get", "pods", "-n", "ragflow-prod", "-l", "app=ragflow-api",
    "-o", "jsonpath={.items[*].metadata.name}"],
    capture_output=True, text=True
    ).stdout.split()
    victim = random.choice(pods)
    print(f"💣 杀死 Pod: {victim}")
    subprocess.run(["kubectl", "delete", "pod", victim, "-n", "ragflow-prod"])

    @classmethod
    def inject_network_delay(cls, target, delay_ms=2000):
    """注入网络延迟——测试超时处理"""
    print(f"🐌 注入 {delay_ms}ms 延迟到 {target}")
    subprocess.run([
    "kubectl", "exec", target, "-n", "ragflow-prod", "–",
    "tc", "qdisc", "add", "dev", "eth0", "root", "netem",
    f"delay", f"{delay_ms}ms"
    ])

    @classmethod
    def drain_redis_memory(cls):
    """填满 Redis 内存——测试 OOM 恢复"""
    print("💧 填满 Redis 内存…")
    subprocess.run([
    "kubectl", "exec", "redis-0", "-n", "ragflow-prod", "–",
    "redis-cli", "DEBUG", "POPULATE", "1000000", "testkey", "1000"
    ])

    @classmethod
    async def run_chaos_experiment(cls):
    """运行混沌实验——同时监控业务指标"""
    # 1. 记录实验前的基准指标
    baseline = cls._get_current_metrics()

    # 2. 注入故障
    cls.kill_random_api_pod()
    await asyncio.sleep(10)

    # 3. 检查是否触发告警
    alerts = cls._get_firing_alerts()

    # 4. 等待系统自动恢复
    await asyncio.sleep(30)

    # 5. 检查恢复后的指标
    recovered = cls._get_current_metrics()

    # 6. 输出实验报告
    print("\\n=== 混沌实验报告 ===")
    print(f"故障注入: 杀死随机 API Pod")
    print(f"触发告警: {len(alerts)} 条 ({[a['name'] for a in alerts]})")
    print(f"恢复时间: {recovered['timestamp'] – baseline['timestamp']:.0f}s")
    print(f"可用率影响: {recovered['availability'] – baseline['availability']:.2%}")
    print(f"结论: {'✅ 通过' if recovered['availability'] > 0.995 else '❌ 未通过'}")

    步骤7:验收自检清单

    # acceptance_checklist.py
    def run_acceptance_checks():
    """平台交付验收自检"""
    checks = []

    # 1. 核心链路可追踪
    trace_id = str(uuid.uuid4())
    resp = requests.get(f"{BASE}/api/v1/chats/xxx", headers={"X-Trace-ID": trace_id})
    jaeger_trace = get_jaeger_trace(trace_id)
    checks.append(("核心链路 TraceID 追踪", jaeger_trace is not None))

    # 2. 百万级文档 P95 检索 < 500ms
    p95 = benchmark_search_p95(doc_count=1000000)
    checks.append(("百万文档 P95 检索", p95 < 500))

    # 3. 自动故障恢复
    kill_api_pod()
    time.sleep(30)
    pod_recovered = count_running_api_pods() >= 3
    checks.append(("API Pod 自动恢复", pod_recovered))

    # 4. 跨机房就近接入
    beijing_latency = measure_latency_from("beijing")
    shanghai_latency = measure_latency_from("shanghai")
    checks.append(("就近接入延迟", beijing_latency < 100 and shanghai_latency < 100))

    # 5. 解析失败率 < 1%
    failure_rate = get_parse_failure_rate_last_24h()
    checks.append(("解析失败率", failure_rate < 0.01))

    # 输出验收报告
    print("\\n=== 平台验收报告 ===")
    all_pass = True
    for name, result in checks:
    status = "✅" if result else "❌"
    print(f" {status} {name}: {'通过' if result else '未通过'}")
    if not result:
    all_pass = False
    print(f"\\n总体: {'✅ 验收通过' if all_pass else '❌ 验收未通过'}")

    完整代码清单

    路径说明
    column/chapter40/ 本章部署清单、脚本、监控配置
    api/ragflow_server.py 平台入口(TraceID 注入)
    rag/nlp/search.py Dealer 检索埋点版
    agent/canvas.py Canvas 追踪版

    4 项目总结

    优点 & 缺点

    维度RAGFlow 自建平台云厂商 RAG 服务开源 + 自运维商业 RAG 平台
    定制化 ★★★ 完全可控 ★★☆ 有限 ★★★ 完全可控 ★☆☆ 黑盒
    运维成本 ★★☆ 需要SRE团队 ★★★ 平台承担 ★★☆ 需要团队 ★★★ 平台承担
    扩展性 ★★★ 源码级扩展 ★★☆ 插件有限 ★★★ 源码级 ★☆☆ 受限
    可观测性 ★★★ 完全定制 ★★☆ 平台提供 ★★★ 完全定制 ★★★ 内置
    跨机房 ★★☆ 需自建 ★★★ 平台提供 ★★☆ 需自建 ★★★ 平台提供
    成本 ★★★ 可控 ★★☆ 按量 ★★★ 可控 ★☆☆ 昂贵

    适用场景

  • 企业级知识管理平台:20+ 部门共享,百万级文档,千级并发用户。
  • 多机房部署:跨地域团队,就近访问+数据同步。
  • 需要源码级定制:WPS Parser、Jira Tool、审批组件等已开发的自定义扩展。
  • SRE 体系完善的组织:有监控、告警、事故预案、混沌工程的团队。
  • 成本敏感且对数据安全要求高:全部私有化部署,无数据出境。
  • 不适用场景:

  • 3 人小团队:全套平台运维(K8s+Prometheus+Jaeger+多机房)需要至少 2 个专职 SRE。
  • 不需要定制化:如果标准 RAGFlow 功能完全满足需求,不需要高级篇的源码扩展。
  • 注意事项

  • 灰度发布必须有回滚预案:金丝雀版本的新功能可能引入"静默错误"(不报错但结果错误)——需要业务指标监控而不仅是技术指标。
  • 跨机房的数据同步延迟:MySQL binlog 主从复制通常有 50-200ms 延迟。如果用户在 A 机房上传文档后立即在 B 机房检索——可能查不到。
  • LLM 费用失控:千万级 Chunk 的 Embedding 费用可能高达数千美元——需要费用预算告警。
  • 混沌实验必须在非业务高峰期执行:实验本身会导致短暂的服务劣化——提前通知用户。
  • 交付物清单(本章产出):

    交付物格式用途
    多机房架构图 draw.io/Mermaid 技术评审、向领导汇报
    K8s 部署清单 YAML 一键部署到新集群
    源码扩展补丁 Python diff 检索埋点、TraceID 注入
    监控大屏 JSON Grafana JSON 导入即用的平台仪表盘
    SLO 定义文档 Markdown 对齐业务与技术的承诺
    容量规划报告 CSV + 图表 3/6/12 个月扩容计划
    备份恢复 Runbook Shell + 文档 运维值班的标准操作流程
    混沌实验报告 Markdown 证明平台韧性
    验收自检脚本 Python 每次上线前的质量门
    成本监控日报 Python + Cron 按部门/供应商的每日费用

    与基础篇综合实战(第16章)的对比:

    维度第16章(基础篇交付)第40章(高级篇交付)
    部署方式 Docker Compose 单机 K8s 多机房集群
    用户规模 500 人 / 1 个部门 2000 人 / 20 个部门
    文档规模 50 份 / 5000 Chunk 20 万份 / 1000 万 Chunk
    可用性 单点故障 99.9% SLO
    可观测性 docker logs Prometheus + Grafana + Jaeger + Loki
    扩展管理 无 企业扩展注册中心
    安全 基础权限 五层隔离 + 沙箱 + 审计
    运维 人工重启 自动恢复 + 灰度发布 + 混沌工程
    交付周期 2 周(学习+搭建) 3 个月(逐步演进)
    5. 备份恢复必须定期演练:不能等到真出事了才第一次执行恢复流程——建议每季度做一次恢复演练。

    常见踩坑经验

    故障现象根因解决方法
    金丝雀版本错误率低但业务指标下降 新版 Retriever 返回了不同排序,导致答案质量变化 金丝雀监控同时看技术指标和业务指标
    跨机房查询延迟高达 500ms 北京的请求错误路由到了上海(绕路了) 检查 GeoDNS 配置和 CDN 节点
    灾难恢复演练后数据不一致 MinIO mirror 不是实时同步,有分钟级延迟 演练后执行数据一致性校验脚本
    混沌实验触发真实 P0 告警 告警没有区分"混沌实验"和"真实故障" 混沌实验开始时在告警系统中标记维护窗口
    LLM 费用日预算频繁触发 某部门的 Prompt 中包含了"请详细回答,不少于 2000 字" 审计 Prompt 长度,设置 max_tokens 硬限制
    灰度发布后旧 Pod 未完全终止 HPA 在灰度期间扩容了旧版本 Pod 灰度期间暂停 HPA,或对旧版本 Deployment 设置 replicas=0
    SLO 达标但用户投诉"太慢" P95 < 8s 但部分长尾用户(P99+)可能等到 20s+ 增加 P99 监控指标,对长尾请求做超时降级
    扩展注册中心的版本检查误报 版本号比较逻辑过于严格(如 “>=0.26” 被解释为 “仅 0.26”) 使用 semver 库做语义化版本比较

    平台持续演进的年度路线图

    本章是全专栏的收官,但平台建设永远在路上。建议的年度演进路线:

    Q1(当前季度):稳定
    – 完成灰度发布流水线
    – 建立 SLO 监控与月度报告
    – 首次混沌实验 + 事故演练
    – 制定 LLM 费用预算和告警

    Q2:扩展
    – 新接 5 个业务部门
    – 上线企业扩展市场
    – 多机房写入能力试点(双主方案)
    – 引入 LLM 缓存层(相同问题直接返回缓存结果,降低成本 30%+)

    Q3:优化
    – 千万级 Chunk 检索性能优化(冷热分离 / 向量量化)
    – 自研 embedding 模型微调(基于公司文档领域数据)
    – 全链路 P99 < 5s
    – 日成本降低 20%(通过缓存 + Prompt 优化 + 模型降级)

    Q4:智能
    – 用户行为分析驱动的自动优化(高频问题缓存、自动发现低质量切片)
    – 多模态 RAG(图片+表格+视频字幕的统一检索)
    – Agent 自主决策的复杂工作流(零人工干预的报销审批)
    – 平台 AIOps(自动根因分析、自动扩容、自动降级)

    步骤8:SLO 定义与监控——平台可用性承诺

    目标:定义平台级的 SLO,并建立实时监控与月度报告。

    # platform_slo.py – 平台 SLO 定义与监控
    from datetime import datetime, timedelta

    class PlatformSLO:
    """平台级服务水平目标"""

    # SLO 定义
    SLOS = {
    "chat_availability": {
    "target": 0.999, # 99.9% 可用性
    "window": "30d", # 滚动 30 天窗口
    "description": "问答 API 成功率",
    },
    "chat_p95_latency": {
    "target": 8000, # P95 < 8 秒
    "window": "7d",
    "description": "问答 P95 延迟 (ms)",
    },
    "parse_success_rate": {
    "target": 0.99, # 99% 解析成功率
    "window": "7d",
    "description": "文档解析成功率",
    },
    "search_p95_latency": {
    "target": 500, # P95 < 500ms(含 Rerank)
    "window": "7d",
    "description": "检索 P95 延迟(含 Rerank)(ms)",
    },
    }

    @classmethod
    def calculate_compliance(cls, metrics):
    """计算 SLO 达标情况"""
    report = {}
    for slo_name, slo_def in cls.SLOS.items():
    actual = metrics.get(slo_name, 0)
    target = slo_def["target"]
    compliant = actual <= target if "latency" in slo_name else actual >= target
    burn_rate = cls._calculate_burn_rate(actual, target)

    report[slo_name] = {
    "target": target,
    "actual": actual,
    "compliant": compliant,
    "burn_rate": burn_rate, # 预算消耗速率
    "remaining_budget": 1 – burn_rate,
    "status": "🟢" if compliant else ("🟡" if burn_rate < 0.8 else "🔴"),
    }
    return report

    @classmethod
    def _calculate_burn_rate(cls, actual, target):
    """计算错误预算消耗速率"""
    # 如果 SLO 是 99.9%,允许 0.1% 的错误
    # 如果当前错误率是 0.15%,burn_rate = 0.15/0.1 = 1.5(超预算 50%)
    if isinstance(target, float) and target < 1:
    error_budget = 1 – target
    current_error = 1 – actual
    return current_error / max(error_budget, 0.0001)
    return 0

    @classmethod
    def generate_monthly_slo_report(cls):
    """生成月度 SLO 报告"""
    metrics = cls._fetch_30d_metrics()
    report = cls.calculate_compliance(metrics)

    print("=" * 70)
    print(f" RAGFlow 平台 SLO 月度报告 – {datetime.now():%Y年%m月}")
    print("=" * 70)
    print(f"{'SLO指标':<24} {'目标':<12} {'实际':<12} {'状态':<6} {'预算消耗':<10}")
    print("-" * 70)

    all_compliant = True
    for name, data in report.items():
    print(f"{data['description']:<24} {data['target']:<12} {data['actual']:<12.4f} "
    f"{data['status']:<6} {data['burn_rate']:.0%}")
    if not data['compliant']:
    all_compliant = False

    print("-" * 70)
    print(f"总体达标: {'✅ 全部达标' if all_compliant else '❌ 存在未达标项'}")
    return report

    # 集成到 Grafana 监控
    # 在 Grafana 中创建 SLO 面板,显示每个 SLO 的实时达标状态
    # 设置告警:当 burn_rate > 1.0 且持续 1 小时 → 触发 P1 告警

    步骤9:成本监控与优化——LLM 费用日预算告警

    目标:建立 LLM 费用的日预算监控,超支自动告警。

    # cost_monitor.py – LLM 费用监控
    from datetime import datetime, timedelta

    class CostMonitor:
    """LLM 费用日预算监控"""

    DAILY_BUDGETS = {
    "OpenAI": 50, # $50/天
    "DeepSeek": 20, # $20/天
    "YunFan": 0, # 自研模型免费
    "Ollama": 0, # 本地模型免费
    }

    ALERT_THRESHOLDS = {
    "warning": 0.7, # 70% 预算 → 提醒
    "critical": 0.9, # 90% 预算 → 告警
    "block": 1.0, # 100% 预算 → 阻止(切换到备用模型)
    }

    def __init__(self):
    self.today_costs = defaultdict(float)
    self.today = datetime.now().date()

    def record_call(self, factory, tokens_in, tokens_out):
    """记录一次 LLM 调用及其费用"""
    # 简化计费模型
    cost_per_1k_input = {"OpenAI": 0.01, "DeepSeek": 0.001, "YunFan": 0, "Ollama": 0}
    cost_per_1k_output = {"OpenAI": 0.03, "DeepSeek": 0.002, "YunFan": 0, "Ollama": 0}

    cost = (tokens_in / 1000 * cost_per_1k_input.get(factory, 0) +
    tokens_out / 1000 * cost_per_1k_output.get(factory, 0))

    self.today_costs[factory] += cost

    # 检查预算
    budget = self.DAILY_BUDGETS.get(factory, 10)
    spent_ratio = self.today_costs[factory] / max(budget, 1)

    if spent_ratio > self.ALERT_THRESHOLDS["critical"]:
    self._send_alert(factory, spent_ratio, self.today_costs[factory], budget)
    if spent_ratio > self.ALERT_THRESHOLDS["block"]:
    return False # 阻止调用
    elif spent_ratio > self.ALERT_THRESHOLDS["warning"]:
    self._send_alert(factory, spent_ratio, self.today_costs[factory], budget,
    severity="warning")

    return True

    def get_daily_report(self):
    """生成每日费用报告"""
    print(f"\\n=== LLM 费用日报 ({self.today}) ===")
    total = sum(self.today_costs.values())
    for factory, budget in self.DAILY_BUDGETS.items():
    cost = self.today_costs.get(factory, 0)
    pct = cost / max(budget, 1) * 100 if budget > 0 else 0
    bar = "█" * int(pct / 10) + "░" * (10 – int(pct / 10))
    print(f" {factory:<12} ${cost:>8.2f} / ${budget:>6} [{bar}] {pct:.0f}%")
    print(f" {'总计':<12} ${total:>8.2f}")

    步骤10:扩展统一管理——企业扩展市场

    目标:解决 5 个部门自定义扩展分散管理的问题。

    # extension_registry.py – 企业扩展注册中心
    import json
    import hashlib
    from pathlib import Path

    class ExtensionRegistry:
    """企业扩展统一注册中心"""

    EXTENSIONS_DIR = Path("/opt/ragflow/extensions")
    REGISTRY_FILE = EXTENSIONS_DIR / "registry.json"

    def __init__(self):
    self.registry = self._load_registry()

    def register_extension(self, name, ext_type, module_path, class_name,
    version, author, description):
    """注册一个扩展"""
    ext_id = hashlib.md5(f"{name}:{ext_type}:{module_path}".encode()).hexdigest()[:8]

    self.registry["extensions"][ext_id] = {
    "id": ext_id,
    "name": name,
    "type": ext_type, # parser / tool / component
    "module_path": module_path,
    "class_name": class_name,
    "version": version,
    "author": author,
    "description": description,
    "registered_at": datetime.now().isoformat(),
    "status": "active",
    "compatible_ragflow_versions": [">=0.26.0, <0.28.0"],
    }
    self._save_registry()
    return ext_id

    def check_compatibility(self, ragflow_version):
    """检查所有扩展与当前 RAGFlow 版本的兼容性"""
    issues = []
    for ext_id, ext in self.registry["extensions"].items():
    if ext["status"] != "active":
    continue
    # 简化版本检查(实际应用 semver 库)
    compatible = any(
    check_version(ragflow_version, vc)
    for vc in ext["compatible_ragflow_versions"]
    )
    if not compatible:
    issues.append({
    "extension": ext["name"],
    "type": ext["type"],
    "compatible_versions": ext["compatible_ragflow_versions"],
    "current_version": ragflow_version,
    })

    if issues:
    print(f"⚠ 发现 {len(issues)} 个兼容性问题:")
    for i in issues:
    print(f" – {i['extension']} ({i['type']}): "
    f"兼容 {i['compatible_versions']}, 当前 {i['current_version']}")

    return issues

    def list_extensions(self, ext_type=None):
    """列出所有已注册扩展"""
    exts = self.registry["extensions"].values()
    if ext_type:
    exts = [e for e in exts if e["type"] == ext_type]

    print(f"\\n=== 企业扩展清单 ({len(exts)} 个) ===")
    print(f"{'名称':<20} {'类型':<10} {'版本':<8} {'状态':<8} {'作者':<12}")
    print("-" * 60)
    for ext in exts:
    print(f"{ext['name']:<20} {ext['type']:<10} {ext['version']:<8} "
    f"{ext['status']:<8} {ext['author']:<12}")

    步骤11:全链路压测报告——验证平台容量

    目标:基于第28章知识,对生产级平台做全链路压测。

    # full_chain_benchmark.sh – 全链路压测脚本
    #!/bin/bash
    echo "=== RAGFlow 全链路压测 ==="
    echo "场景: 2000 并发用户, 北京+上海双机房"
    echo ""

    # 测试矩阵
    USERS=(100 500 1000 2000)
    DURATION=300 # 每轮 5 分钟

    for users in "${USERS[@]}"; do
    echo "— 并发: $users —"

    # 启动 locust (分布式模式, 北京+上海各一个 worker)
    locust -f locustfile.py \\
    –headless \\
    –users $users \\
    –spawn-rate 50 \\
    –run-time ${DURATION}s \\
    –master-host locust-master \\
    –host http://ragflow-api:9380 \\
    –csv "benchmark_${users}users"

    # 等待稳定
    sleep 30

    # 收集指标
    python3 collect_metrics.py –users $users –output "metrics_${users}users.json"

    echo " P50: $(jq '.chat_latency.p50' metrics_${users}users.json)ms"
    echo " P95: $(jq '.chat_latency.p95' metrics_${users}users.json)ms"
    echo " P99: $(jq '.chat_latency.p99' metrics_${users}users.json)ms"
    echo " 成功率: $(jq '.success_rate' metrics_${users}users.json)"
    echo " HPA 副本数: $(jq '.hpa_replicas' metrics_${users}users.json)"
    echo ""
    done

    echo "=== 压测完成, 报告已保存 ==="
    # 生成压测报告
    python3 generate_benchmark_report.py

    思考题

  • 当前平台的多机房方案是"北京写、两地读"。如果业务要求两地都能写(如文档可以在上海上传),该如何设计"多主写入 + 冲突解决"方案?MySQL 的双主复制有什么陷阱?

  • 千万级 Chunk 下,Infinity 的检索 P95 延迟可能从 50ms 增长到 500ms。你已做了向量量化、索引分片、查询缓存等优化——但仍不满足 SLO(P95<200ms)。是否应该考虑"冷热分离"——6 个月前的文档自动降级到慢存储(如本地 HD),但检索时可能查不到——如何在 Prompt 中告知 LLM 检索结果不完整?

  • (答案提示见附录 D。)

    延伸阅读与资源

    10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用 后端工程师转型AI第一课-Ollama 与私有化大模型实战 大型语言模型(LLM) vLLM 高性能推理落地实战 Agent开发之LlamaIndex 实战修炼与源码进阶 大语言模型Transformers 实战修炼与源码剖析

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 第40章:【高级篇综合实战】从源码到生产打造可观测高可用 RAGFlow 平台
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!