1 项目背景
业务场景
「云帆科技」的 RAGFlow 平台经过一年的运营,已从单机房 Docker Compose 演进到多机房 K8s 集群,服务 20 个业务部门、2000 名员工、20 万份文档、1000 万 Chunk。随着规模增长和架构复杂度提升,团队面临一系列"平台级"挑战:
CTO 要求:将过去 39 章的知识融会贯通,交付一个"从源码到生产"的完整平台方案,包含架构图、源码扩展、部署清单、监控大盘、压测报告和事故预案。
痛点
平台级交付 vs 单点优化的核心差异:
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 项目总结
优点 & 缺点
| 定制化 | ★★★ 完全可控 | ★★☆ 有限 | ★★★ 完全可控 | ★☆☆ 黑盒 |
| 运维成本 | ★★☆ 需要SRE团队 | ★★★ 平台承担 | ★★☆ 需要团队 | ★★★ 平台承担 |
| 扩展性 | ★★★ 源码级扩展 | ★★☆ 插件有限 | ★★★ 源码级 | ★☆☆ 受限 |
| 可观测性 | ★★★ 完全定制 | ★★☆ 平台提供 | ★★★ 完全定制 | ★★★ 内置 |
| 跨机房 | ★★☆ 需自建 | ★★★ 平台提供 | ★★☆ 需自建 | ★★★ 平台提供 |
| 成本 | ★★★ 可控 | ★★☆ 按量 | ★★★ 可控 | ★☆☆ 昂贵 |
适用场景
不适用场景:
注意事项
交付物清单(本章产出):
| 多机房架构图 | draw.io/Mermaid | 技术评审、向领导汇报 |
| K8s 部署清单 | YAML | 一键部署到新集群 |
| 源码扩展补丁 | Python diff | 检索埋点、TraceID 注入 |
| 监控大屏 JSON | Grafana JSON | 导入即用的平台仪表盘 |
| SLO 定义文档 | Markdown | 对齐业务与技术的承诺 |
| 容量规划报告 | CSV + 图表 | 3/6/12 个月扩容计划 |
| 备份恢复 Runbook | Shell + 文档 | 运维值班的标准操作流程 |
| 混沌实验报告 | Markdown | 证明平台韧性 |
| 验收自检脚本 | Python | 每次上线前的质量门 |
| 成本监控日报 | Python + Cron | 按部门/供应商的每日费用 |
与基础篇综合实战(第16章)的对比:
| 部署方式 | 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 实战修炼与源码剖析
网硕互联帮助中心





评论前必须登录!
注册