大规模基础设施


文档摘要

大规模基础设施 构建服务数百万用户的系统,远不止一台服务器那么简单。本文件覆盖可扩展性模式、分布式系统基础、微服务、数据流水线、数据库扩展、搜索与向量系统、可观测性、可靠性工程以及 CI/CD 一个每秒处理 1 个请求的模型,笔记本就能跑。要每秒处理 10 万个请求、保持 99.9% 可用性,就需要分布式系统、自动故障转移和精心设计的数据流水线。本文件讲的就是横亘两者之间的那些模式。 可扩展性 垂直扩展(vertical scaling,scale up):换台更强的机器。更多 CPU、更多内存、更大的 GPU。简单,但有硬上限(最大能买到的机器)和单一故障点。 水平扩展(horizontal scaling,scale out):加更多机器。每台分担一部分流量。

大规模基础设施

构建服务数百万用户的系统,远不止一台服务器那么简单。本文件覆盖可扩展性模式、分布式系统基础、微服务、数据流水线、数据库扩展、搜索与向量系统、可观测性、可靠性工程以及 CI/CD

  • 一个每秒处理 1 个请求的模型,笔记本就能跑。要每秒处理 10 万个请求、保持 99.9% 可用性,就需要分布式系统、自动故障转移和精心设计的数据流水线。本文件讲的就是横亘两者之间的那些模式。

可扩展性

  • 垂直扩展(vertical scaling,scale up):换台更强的机器。更多 CPU、更多内存、更大的 GPU。简单,但有硬上限(最大能买到的机器)和单一故障点。

  • 水平扩展(horizontal scaling,scale out):加更多机器。每台分担一部分流量。没有单机上限,但需要:负载均衡(第 01 节)、数据分区,以及处理分布式状态。

  • 无状态服务(stateless services) 天生就能水平扩展。在负载均衡器后面加更多实例即可。一个启动时加载权重、各请求互相独立的模型推理服务器就是无状态的——任何实例都能处理任何请求。

  • 有状态服务(stateful services)(数据库、KV-cache、特征存储)更难扩展。状态必须跨机器分区(分片,第 01 节),并复制以容错。

  • 可扩展性方程:对一个有 n 台服务器的系统:

    • 理想情况:吞吐量线性扩展(n 台服务器 → n\times 吞吐量)。
    • 实际情况:协调、负载均衡和数据传输带来的开销,让吞吐量呈次线性扩展。Amdahl 定律(第 13 章)在这里同样适用:串行部分(共享状态、协调)限制了加速比。

分布式系统

  • 分布式系统(distributed system) 是一组协同提供服务的机器。它面临几个根本性挑战:

  • 网络分区(network partitions):机器之间并不总能通信。网线被剪断、交换机故障、数据中心断电。系统必须能处理部分失败。

  • 时钟漂移(clock skew):不同机器的时钟不一致。"事件 A 在机器 1 上发生在 10:00:01"和"事件 B 在机器 2 上发生在 10:00:01",并不意味着它们同时发生。逻辑时钟(logical clocks)(Lamport 时间戳、向量时钟)在不依赖物理时钟的情况下建立顺序。

  • 共识(consensus):多台机器如何对一个值达成一致(比如谁是 leader)?Raft 是标准的共识算法。一组节点选出 leader,由 leader 处理所有写。leader 故障时,剩余节点选出新 leader。需要多数派(5 节点中的 3 个)才能运行,所以能容忍 \lfloor(n-1)/2\rfloor 个故障。

  • 分布式锁(distributed locks):确保只有一台机器执行某个关键操作。Redlock(基于 Redis)在多个 Redis 实例上获取锁。只要多数实例授予,锁就算拿到。用途:防止重复部署模型、确保只有一个训练任务在写 checkpoint。

微服务

微服务 ML 架构:API 网关把请求路由到特征、模型、日志服务,每个服务都有自己的数据库,之间通过消息队列连接

  • 微服务(microservices) 把一个系统拆分成若干小而独立可部署的服务。每个服务负责一个领域:
┌─────────────┐ ┌──────────────┐ ┌─────────────┐ │ API Gateway │→ │ Feature Svc │→ │ Feature DB │ └─────────────┘ └──────────────┘ └─────────────┘ │ ├────────→ ┌──────────────┐ ┌─────────────┐ │ │ Model Svc │→ │ Model Store │ │ └──────────────┘ └─────────────┘ │ └────────→ ┌──────────────┐ ┌─────────────┐ │ Logging Svc │→ │ Log Store │ └──────────────┘ └─────────────┘
  • 优点:独立部署(更新模型服务时不用碰特征服务)、独立扩缩容(按请求负载扩模型服务器,按特征读取速率扩特征服务器)、技术自由(模型服务用 Python,特征服务用 Go)。

  • 缺点:网络开销(每次服务调用都是一次网络往返)、复杂度(调试要跨越多个服务)、数据一致性(没有跨服务事务)。

  • 服务发现(service discovery):API 网关怎么找到模型服务?方案有:基于 DNS(每个服务注册一个 DNS 名)、K8s service(内置)、或服务注册中心(Consul、Eureka)。

  • Saga 模式:对于跨多个服务的操作(创建用户 + 分配资源 + 发欢迎邮件),用一个 saga:一连串本地事务,任何一步失败时执行补偿动作。

数据流水线

  • ML 系统消耗海量数据。数据流水线(data pipelines) 负责搬运、转换和服务这些数据:

批处理

  • 按固定间隔(每小时、每天)处理大批量数据。

  • MapReduce:最早的批处理范式。Map(对每条记录独立转换)→ Shuffle(按键分组)→ Reduce(每组聚合)。概念上简单,但写起来啰嗦。

  • Apache Spark:现代批处理引擎。内存计算(对迭代算法比 MapReduce 快 100 倍)。支持 SQL、DataFrame 和 ML 流水线。大规模特征工程的事实标准。

  • 示例:为推荐系统计算用户特征。输入:过去 30 天 10 亿条用户活动事件。输出:1 亿条用户特征向量。每天作为一个 Spark 作业运行,输出到特征存储。

流处理

  • 数据一到就实时处理(亚秒级延迟)。

  • Apache Flink:领先的流处理引擎。精确一次(exactly-once)处理、事件时间处理(按事件发生时间而非到达时间处理)、窗口(滚动、滑动、会话窗口)。

  • Kafka Streams:Kafka 内置的轻量级流处理。适合无需另起集群的简单转换(过滤、聚合)。

  • 示例:实时欺诈检测。每笔信用卡交易是一个 Kafka 事件。一个 Flink 作业计算滚动统计量(交易频率、地点变化),并在 100ms 内标记异常。

Lambda 架构

  • 把批处理和流处理结合起来。批层(batch layer) 提供精确、全面的结果(但有延迟)。速度层(speed layer) 提供近似的实时结果。一个服务层(serving layer) 把两者合并。

  • 在实践中,很多团队现在用 Kappa 架构:只用流处理,把流当成唯一事实来源。流是可以重放的(Kafka 保留事件),所以通过重放流就能模拟批处理。

ML 训练基础设施

  • 训练一个前沿模型(100B+ 参数)是一个大规模基础设施问题:上千张 GPU 跑上几个月,消耗数兆瓦电力,产生 PB 级数据,花费数千万美元。基础设施决定了训练是成是败。

GPU 集群

  • 一个训练集群是一组通过高速网络互联的 GPU 服务器。关键组件:

GPU 集群:每个节点有 8 张 GPU 通过 NVLink 互联,节点之间用 InfiniBand 按胖树拓扑连接,规模从 64 张扩展到 16000+ 张

  • GPU 服务器(节点):每台服务器有 4-8 张 GPU。典型配置:8 × H100 GPU、2 × AMD EPYC CPU、2 TB 内存、30 TB NVMe SSD。节点内的 GPU 用 NVLink 互联(H100 上每张 900 GB/s),比 PCIe 快 30 倍。

  • 集群规模:小训练集群有 64-256 张 GPU(8-32 个节点)。前沿模型训练集群有 4000-32000 张 GPU(500-4000 个节点)。Meta 的 Llama 3 用了 16384 张 H100。Google 在 8000+ 芯片的 TPU pod 上训练。

  • 粗算一下:训练一个 70B 模型约需 $2M 算力。训练一个 400B+ 前沿模型约需 $50-100M。集群硬件本身按 H100 价格(每张 $30K × 16000 张 = $480M)约要 $500M-$1B。

网络拓扑

  • GPU 节点之间的网络是最关键的基础设施组件。如果 GPU 之间不能足够快地交换梯度,它们就会干等着通信完成。

  • InfiniBand 是 GPU 集群网络的标准。NVIDIA 的 Quantum-2 InfiniBand 每端口 400 Gb/s。每个节点通常有 8 个 InfiniBand 端口(每张 GPU 一个),整节点共 400 GB/s 的二分带宽。

  • RDMA(远程直接内存访问,Remote Direct Memory Access):InfiniBand 支持 RDMA,可以在不同节点的 GPU 显存之间直接传数据,绕过 CPU。这把延迟从约 100μs(TCP)降到约 1μs,对高效的梯度 all-reduce(第 6 章)至关重要。

  • 网络拓扑很重要胖树(fat tree,Clos 网络) 提供完整二分带宽(任意两张 GPU 之间都能以全速通信)。更便宜的拓扑(rail-optimized(轨道优化)3D 环面(3D torus))带宽较低但更省钱。拓扑必须和并行策略匹配:

    • 数据并行:所有 GPU 之间做 all-reduce → 需要高二分带宽(胖树)。
    • 张量并行:节点内通信 → NVLink 搞定(不需要走网络)。
    • 流水线并行:相邻流水线阶段之间通信 → 只在特定节点对之间需要带宽(rail-optimized 就够了)。
  • 以太网替代方案RoCE v2(融合以太网上的 RDMA)在标准以太网基础设施上提供 RDMA。比 InfiniBand 便宜,但延迟更高、拥塞更多。Google 在某些 TPU pod 网络里用 RoCE。Ultra Ethernet Consortium 正在为 AI 工作负载开发无损以太网。

训练用存储

  • 训练需要三层存储:

    • 数据集存储:训练语料(1-100 TB 文本,或 PB 级多模态数据)。存在分布式文件系统或对象存储里。必须支持高吞吐顺序读(数据加载器大批量读取数据)。LustreGPFS 是常见的 HPC 文件系统;云上的替代品有 FSx for Lustre(AWS)和 Filestore(GCP)。

    • checkpoint 存储:周期性保存的训练状态(模型权重 + 优化器状态 + 调度器状态)。对一个混合精度加 Adam 优化器的 70B 模型:每个 checkpoint 约 560 GB(70B × 4 字节 × 2,优化器部分)。一个 3 个月的训练每小时存一次 = 约 2000 个 checkpoint = 1.1 PB。实际中只保留最近 N 个,旧的删掉。存储必须足够快,以至于存 checkpoint 不会显著拖慢训练。

    • 日志和指标:实验跟踪数据(loss 曲线、学习率调度、梯度范数)。相对小,但必须实时写入。W&B、MLflow 或 TensorBoard 负责这部分。

  • 存储瓶颈:一个 16000 张 GPU 的集群加载一个训练 batch,需要持续以约 100 GB/s 读数据。如果文件系统撑不住这个吞吐,GPU 就会干等数据。数据流水线优化(预取、缓存、用 WebDataset 或 Mosaic Streaming 优化格式)至关重要。

作业调度

  • 一个 GPU 集群要服务多个团队和项目。作业调度器(job scheduler) 把 GPU 分配给训练作业:

  • SLURM:标准的 HPC 作业调度器。用户提交作业时指定 GPU 数、内存和时间上限。SLURM 分配资源并管理队列。支持基于优先级的调度、抢占,以及团队间的公平份额分配。

  • 带 GPU 调度的 Kubernetes(第 18 章第 02 节):云原生方案。K8s 的 GPU 设备插件把 GPU 暴露为可调度资源。VolcanoRun:ai 在此基础上加了 ML 专属调度功能:gang 调度(一次把一个作业需要的所有 GPU 全部分配下来,而不是一张一张地分)、优先级队列、GPU 时间共享。

  • 调度难题

    • 碎片化:一个 1000 张 GPU 的集群可能有 200 张空闲,但分散在 50 个节点上(每节点 4 张空闲)。一个需要 128 张连续 GPU 的作业就没法跑,尽管总数够。去碎片化(迁移作业以整合空闲 GPU)或拓扑感知调度(分配连接良好的 GPU)能缓解这个问题。
    • 优先级与抢占:紧急实验应该抢占低优先级作业。但抢占一个已经跑了 2 天的训练作业会浪费算力。调度器必须在优先级和效率之间权衡。
    • 公平份额:即使某团队提交的作业超过它的份额,各团队在长期来看也应当拿到它该得的那部分算力。

容错

  • 在数千张 GPU 跑几个月的规模下,硬件故障不是例外,而是常态。一个 16000 张 GPU 集群的平均故障间隔(MTBF)是以小时计,而不是以月计。

  • 常见故障:GPU 显存错误(ECC 可纠正和不可纠正)、NVLink 故障(节点内 GPU 间通信)、InfiniBand 链路故障(节点间通信)、节点崩溃(内核 panic、电源故障)和存储故障(磁盘或控制器故障)。

  • checkpoint 是第一道防线。每 N 步保存一次完整的训练状态(模型、优化器、数据加载器位置)。一旦发生故障:定位故障节点,替换或剔除它,从最近的 checkpoint 重启训练。一次故障的代价,就是从上一个 checkpoint 到故障点之间的那段算力。

  • checkpoint 频率的权衡:频繁 checkpoint(每 10 分钟一次)故障时浪费的算力更少,但拖慢训练(存 560 GB 要花时间)。不频繁 checkpoint(每 2 小时一次)更快,但故障时最多浪费 2 小时算力。多数团队每 20-60 分钟 checkpoint 一次。

  • 弹性训练(elastic training):现代框架(PyTorch Elastic、DeepSpeed)支持在不重启的情况下改变训练规模。500 个节点里坏了 2 个,训练就以 498 个节点继续。坏节点被替换后,训练会自动把它们重新纳入。

  • 健康监控:持续监控所有 GPU(温度、显存错误、计算吞吐)、网络链路(丢包、延迟)和存储(吞吐、错误率)。异常时自动告警。有些集群会周期性跑 GPU 健康检查(一段短计算测试),以便在硬件真正失效前就提前发现劣化。

  • 大规模下:训练 Meta 的 Llama 3(16384 张 H100,54 天)期间经历了约 466 次作业中断。有效训练时间只占墙上时钟时间的约 90%——10% 损耗在故障和恢复上。能做到 90%(而不是 50% 或 70%)的基础设施,正是把"能训练前沿模型的组织"和"不能的"区分开来的东西。

成本与效率

  • 训练基础设施的成本主要由 GPU 小时数决定:
组件 占总成本百分比
GPU 算力 70-80%
网络(InfiniBand) 10-15%
存储 5-10%
散热和供电 5-10%
  • GPU 利用率(Model FLOPs Utilisation,MFU) 衡量 GPU 理论峰值性能中真正用于有效计算的比例。H100 峰值是 989 TFLOPS(FP8)。达到 40-50% MFU 算不错,50-60% 算优秀。差距来源于:通信开销(all-reduce、流水线气泡)、显存带宽限制,以及 checkpoint 和数据加载期间的空闲。

  • 提升 MFU:让计算和通信重叠(第 6 章)、用高效注意力(Flash Attention,第 16 章)、优化数据加载(防止 GPU 饿着)、降低 checkpoint 开销(异步 checkpoint,先写快速 NVMe,再后台拷到持久存储)。

  • 自建还是租用:小规模下(<256 张 GPU),云更便宜(无前期投入,按小时付费)。大规模下(>1000 张 GPU,持续用 6 个月以上),自建硬件更便宜(3 年总拥有成本 TCO 低约 2-3 倍)。多数 AI 公司两者混用:自建集群做长期训练,云做突发容量和实验。

数据库扩展

  • 读副本(read replicas):把读查询路由到主库的副本上。主库负责写,副本负责读。因为大多数工作负载都是读多写少(95%+ 是读),这让读吞吐随副本数线性扩展。

  • 分区(partitioning,即第 01 节的分片):把数据分散到多个数据库。每个分区独立,可以并行读写。难点是跨分区查询(连接不同分片的数据)。

  • 连接池(connection pooling):数据库的连接数有限。连接池(PostgreSQL 的 PgBouncer)在多个请求间复用连接,避免几百个服务实例同时抢着连导致的连接耗尽。

搜索与向量系统

文本搜索

  • 倒排索引(inverted index):文本搜索的基石。对每个词,存一个包含该词的文档列表。一次查询就是求各查询词列表的交集。Elasticsearch 是标准选择:分布式、实时,支持全文搜索、聚合和地理空间查询。

  • BM25:标准的文本检索打分函数。按词频、逆文档频率和文档长度归一化来给文档打分。简单但有效——对于关键词密集的查询,它今天仍然能与神经方法一较高下。

向量搜索

  • 向量数据库 存嵌入向量(高维向量),并支持快速的近似最近邻(approximate nearest neighbour,ANN)搜索。给定一个查询向量,找出 k 个最相似的已存向量。

  • FAISS(Facebook AI Similarity Search):一个用于 ANN 搜索的库(不是数据库)。支持多种索引类型:

    • Flat:精确搜索,O(n)。用于小数据集或作为真值。
    • IVF(倒排文件,Inverted File):把向量分到若干簇,只搜索最近的几个簇。每次查询 O(n/k)
    • HNSW(分层可导航小世界,Hierarchical Navigable Small World):基于图。构建一个分层图,从粗到细地导航。极快且准确,是大多数应用的默认选择。
    • 乘积量化(Product Quantisation,PQ):把向量压缩成紧凑编码,以节省内存搜索。用精度换内存。
  • 托管式向量数据库:Pinecone、Weaviate、Milvus、Qdrant。它们处理 FAISS 不擅长的扩缩容、复制和实时更新。

  • 对 RAG 而言(检索增强生成):用户查询 → 用文本编码器嵌入 → 在向量库里搜相关文档 → 把检索到的文档拼到 LLM 提示词前面。检索的质量直接决定了 LLM 回答的质量。

可观测性

  • 可观测性(observability) 指的是从系统的外部输出推断出系统内部正在发生什么的能力。三大支柱:

日志

  • 结构化日志(structured logs,JSON) 可搜索、可解析。非结构化日志("ERROR: something failed")做不到。永远要记录:时间戳、服务名、请求 ID(用于跨服务追踪)、严重级别,以及相关上下文。

  • ELK 栈(Elasticsearch、Logstash、Kibana):标准的日志流水线。Logstash 收集并转换日志,Elasticsearch 给它们建索引,Kibana 负责可视化和搜索。

指标

  • 指标(metrics) 是随时间变化的数值测量:请求速率、错误率、延迟分位数、GPU 利用率、队列深度。Prometheus 从服务抓取指标,Grafana 在带告警的仪表盘上可视化它们。

  • RED 方法(针对服务):Rate(每秒请求数)、Errors(错误率)、Duration(延迟)。每个服务都要监控这三项。

  • USE 方法(针对资源):Utilisation(利用率)、Saturation(饱和度,队列深度)、Errors(错误)。每种资源(CPU、GPU、内存、磁盘、网络)都要监控这三项。

追踪

  • 分布式追踪(distributed tracing) 跟踪一个请求在多个服务间的旅程。一个用户请求打到 API 网关 → 特征服务 → 模型服务 → 后处理。一条追踪(trace) 记录每一跳的耗时,显示延迟花在了哪里。

  • OpenTelemetry:追踪、指标和日志的开放标准。只插桩一次,导出到任何后端(Jaeger、Zipkin、Datadog)。

可靠性

  • SLO(服务等级目标,Service Level Objective):目标可靠性。"99.9% 的请求在 <200ms 内完成。"这给出了一个具体的错误预算:0.1% 的请求(约每月 43 分钟)可以慢或失败。

  • SLI(服务等级指标,Service Level Indicator):测量值。"过去 5 分钟的 99 分位延迟。"

  • SLA(服务等级协议,Service Level Agreement):带后果的合同承诺。"如果可用性跌破 99.95%,客户获得补偿。"

  • 错误预算(error budgets):如果你的 SLO 是 99.9% 而你实际做到了 99.99%,你就有预算去做有风险的变更(部署新模型、迁移数据库)。如果你只做到 99.85%,那就冻结一切变更,专注可靠性。错误预算把可靠性从一个抽象目标变成了一种可度量的资源。

  • 混沌工程(chaos engineering):故意注入故障(杀一台服务器、加网络延迟、弄坏数据),测试你的系统能不能正确应对。Netflix 的 Chaos Monkey 会随机终止生产实例。系统没倒下,说明它有韧性;倒下了,你就在用户之前发现了一个 bug。

CI/CD

  • 持续集成(Continuous Integration,CI):每次代码变更都自动构建和测试。每次 push 触发:lint、类型检查、单元测试、集成测试。任何一项失败,变更就被拒。这能在 bug 进生产之前抓住它。

  • 持续部署(Continuous Deployment,CD):把通过 CI 的变更自动部署上去。部署策略:

    • 蓝绿(blue-green):跑两套完全一样的环境(蓝 = 当前,绿 = 新)。瞬间把流量从蓝切到绿。如果绿出问题,切回蓝(瞬间回滚)。

    • 金丝雀(canary):把一小部分流量(1-5%)导到新版本。监控错误。指标好的话再逐步加流量。这能限制坏部署的影响范围。

    • 功能开关(feature flags):把新代码部署上去,但用开关藏起来。先对一部分用户启用开关(内部测试者,然后 beta 用户,然后所有人)。把"部署"(代码已经上线)和"发布"(用户能看到功能)解耦。

  • 对 ML 而言:CI/CD 还包括模型专属步骤。一次模型变更触发:单元测试(形状测试、梯度检查)、在留出集上评估(准确率不能退化)、影子部署(新模型和老模型并行跑,对比输出)、以及逐步放量(从 1% → 100% 的金丝雀)。


发布者: 作者: HenryNdubuaku 转发
评论区 (0)
U