100TB分布式存储系统设计

前言

系统设计通常从以下顺序考虑

  1. clarify requirement:清楚约束边界
  2. high-level architecture:高层设计
  3. Metadata design:元数据设计
  4. Sharding:分片
  5. Replication/EC:副本可用性/数据冗余
  6. Scaling/rebalancing:重平衡
  7. HotSpot:热点消除
  8. Failure scenairo

1. 确认需求

首先清楚约束边界:

  1. 100TB是当前容量还是未来峰值。
  2. 对象的平均大小
  3. 写多还是读多
  4. 需要跨机房容灾吗
  5. SLA是99.9还是99.99

假设:
容量:100TB usable
object avg size:16MB
workload:80% read / 20% write
availability:99.99%
durability:11 nines

2. 高层架构

我们设计四层架构

  1. Client 客户端可以执行命令。
  2. API Gateway层:负责协议对接,鉴权,请求路由,大文件切片。
1
2
3
4
PUT object
GET object
DELETE object
LIST bucket
  1. Metadata Service:存储对象名称,存储桶,权限,对象到实际物理数据的映射关系,这是控制面。
  2. Data Node:实际存储的节点,以Chunk或者Slice形式。

3,元数据设计

假设每个分片为16MB,那么100TB/16MB = 6.5M Objects
每个metadata设计为

1
2
3
4
5
6
object key         64B 64 字节可以覆盖 80% 以上的标准文件名路径
size 8B 如果只使用4B,2^32=4GB,对于大文件来说不够用
version 8B 乐观锁(Optimistic Locking)和多版本控制(MVCC)
replica pointers 24B 3副本的直接寻址,每个副本占用8B,
checksum 32B 32B正好是SHA-256算法的输出长度,用于完整性校验
timestamp 8B 64bit的时间戳是Unix最精确的纳秒级别长度

最终元数据大小为

$$6.5M \times 144B = 1 GB$$

元数据可以轻松加载进Master的内存


元数据可以存储在TiKV或者RockDB中,与GFS不同,元数据不仅仅存储于Master中,为了避免:

  • 内存瓶颈:海量小文件会导致每个文件需要一个Object Key。
  • 单点故障:虽然有辅助的Shadow Master,但是切换不是瞬间,存在数据不一致风险。
  • 写入瓶颈:在高并发写入时,Master 为了保证一致性,会对目录树加锁。大量并发请求会阻塞在 Master 的单机锁上。同时控制流会经过Master节点,Master的网卡和CPU会成为瓶颈

虽然 GFS 的设计很经典,但在 100TB 且存在海量小文件的场景下,单机 Master 会面临 内存撑爆 (Memory Bottleneck) 和 写入并发受限 (Write Bottleneck) 的问题。同时,单机 Master 是 SPOF (单点故障),不符合金融级高可用的要求。因此,我倾向于使用类似 TiKV 的分布式 KV 来存储元数据,利用其 Raft 协议保证强一致性和自动容错,并通过分片机制解决元数据的扩展性问题。

4. 数据切片

对于小对象(比如 < 100MB),客户端发起一个简单的 HTTP PUT 请求,将整个对象作为 Payload 发送。

当对象很大(比如 10GB 或 1TB)时,一次性 PUT 会面临极高的风险:因此我们需要引入切片。

Multipart Upload 的具体工作流程(S3 标准)
这个过程通常分为三个阶段,按此顺序进行:

第一阶段:初始化 (Initiate)

客户端告诉服务器:“我要上传对象 movie.mp4,请准备好。”
服务器生成一个全局唯一的 Upload ID 并返回给客户端。这个 ID 用于跟踪后续所有分片。

第二阶段:上传分片 (Upload Parts)

客户端将 10GB 文件切成(例如)1000 个 10MB 的分片。
并行上传: 客户端可以启动 10 个线程,并发地上传不同的分片。
分片标识: 每个分片请求都带上 Upload ID 和一个 Part Number(序号,用于排序)。
校验: 每一个分片上传成功,服务器会返回一个 ETag(通常是该分片的 MD5)。

第三阶段:完成/合并 (Complete)

当所有分片上传完毕,客户端发送一个“完成”请求,并提交一份清单(包含所有 Part Number 和对应的 ETag)。
逻辑合并: 服务器核对清单无误后,宣布对象已创建。


当你完成 Multipart Upload 后,元数据服务器(如 TiKV)里记录的不再是一个简单的物理地址,而是一个 有序列表:
对象:movie.mp4 的元数据

1
2
3
4
Part 1 -> Node_A, Disk_1, Offset_100 (Replica 1,2,3)
Part 2 -> Node_B, Disk_5, Offset_500 (Replica 1,2,3)
...
Part N -> Node_Z, Disk_2, Offset_900 (Replica 1,2,3)

读取时:
当用户下载这个 10GB 对象时,网关从元数据中拿到这个分片列表,按顺序从不同的数据节点读取内容,并以 Stream(流)的形式返回给用户。用户感觉是在读一个连续的文件,但实际数据分布在全集群的不同机器上。

这样做:

  • 支持断点续传:比如任意Chunk下载失败,只需要重新传这个分片,无需全部上传。
  • 吞吐量大:客户端可以利用 100Gbps 的内网,同时向 50 台不同的数据存储节点(OSD)发送数据,写入性能随集群规模线性增长。
  • 支持未知文件大小:比如上传数据流Stream,可以一边生成分片一边上传。
  • 实现Append:虽然对象存储通常不支持原地修改,但通过分片,可以在不改变原有分片的情况下,上传新的分片来重组对象。

对于每个Chunk的副本分片到不同机架上:
采用机架感知(Rack-aware)的放置策略。通过将三个副本分散到三个不同的机架(Rack A, B, C),我们可以有效规避 ToR 交换机或 PDU 掉电引起的单点故障。在实现上,利用类似 CRUSH 的算法进行拓扑映射,并优先保证副本在物理空间上的最大隔离。同时,在读取数据时,利用拓扑亲和性(Topology Affinity)优先读取距离客户端最近的副本,以降低跨交换机的网络流量。

5. Replication vs EC

小规模使用3-Replication,特点是

  • 快速读取
  • 写简单
  • 回复简单
  • 100%可用需要300%物理存储

大规模使用EC纠删码,比如RS(10,4),但是10+4方案为了保证机架感知,需要14台机架,确认物理机架可用数量至少14台。10+4意味着10个数据块+4个校验块,空间利用率$\frac{10}{10+4}$,允许最多挂掉4个节点。


通常来说

  • hot objects 采用 replication:如果存热点数据的某台机器挂了,系统直接从另一个副本读,用户完全感知不到延迟增加。
  • cold objects 采用 EC:冷数据几乎不被访问,所以 EC 带来的“读解码延迟”是可以接受的。

动态转储:

更深入设计,一个优秀的对象存储可以采用混合容错的方式,热点数据可以转储为冷数据。

  • 策略 1:基于时间 (TTL): 对象创建后的前 7 天(热期)用 3 副本;7 天后自动转为 EC。
  • 策略 2:基于访问频率 (Access Frequency): 使用LRU策略,系统监控访问次数,如果一个对象在 30 天内访问次数少于 5 次,后台线程自动触发“副本转 EC”的 Job。
  • 策略 3:写缓冲 (Write Buffer): 所有新对象先用副本方式快速写入(保证写入低延迟),等凌晨系统负载低时,后台批量异步转换为 EC 存储。

6. 扩容与重平衡

当系统加入新节点后,如何迁移数据:

  • 设置水位线:磁盘使用率达到80%后自动触发
  • 不能全量迁移:传统哈希中,由$Hash(key)% N$ 变为$Hash(key)% (N+1)$时,所有桶几乎都要重平衡,这样影响过大
  • PG 迁移: 以 Placement Group 为单位进行迁移,而不是单个对象。为了不迁移所有Chunk,我们引入了中间层,这个中间层通常被称为 Placement Group (PG)、Bucket 或 Shard (分片)。当新节点加入时,重平衡器(Rebalancer)只需要计算 PG 的归属变化,而不是每个对象的归属变化。
  • 流量控制 (Throttling): 重平衡任务必须是后台低优先级的。需通过令牌桶限制带宽,防止由于数据迁移导致正常 IO 请求(用户请求)的延迟抖动。
  • 一致性保持: 迁移过程中,利用元数据的版本号或 Proxy 模式,确保读取请求能重定向到新位置。
  • 虚拟节点:让 1 个物理节点在环上化身为 100 个甚至 1000 个虚拟节点。用于解决数据倾斜的问题,比如一个机架上占用90%,而另一个机架只有1%的数据
    负载均衡: 虚拟节点越多,环被切分得越碎,数据分布越均匀。
    异构支持: 性能强的机器(如 64TB 硬盘)可以分配 200 个虚拟节点,性能弱的(如 16TB)只分配 50 个。
  • 一致性哈希;使用ketama算法,设计一个一致性哈希环,
  • 令牌桶和调度算法:让迁移的优先级尽可能低,以确保正常访问的带宽足够。

一致性哈希算法

Ketama算法

  • 步骤:
    哈希空间环化: 将哈希输出范围想象成一个圆环。
  • 节点映射: 使用相同的哈希算法(如 MurmurHash3,性能高且分布均匀)将服务器节点的 ID 或 IP 映射到环上的多个点(虚拟节点)。
  • 数据映射: 计算对象的 Key 的哈希值,也映射到环上。
  • 查找规则: 从对象在环上的位置开始,顺时针行走,遇到的第一个节点就是该数据的归属。
  • 如何实现不变:
    当你新增节点 N_new 时,它只会插入到圆环的某个位置。
    只有在 N_new 逆时针方向到上一个节点之间的那一小段弧线上的对象,其顺时针遇到的第一个节点才变成了 N_new
    圆环上其他所有区域的对象,顺时针遇到的节点依然是原来的那个,因此位置完全不变。

JumpConsistentHash

原理: 它通过一个随机数生成器(由 Key 作为种子)来模拟“抛硬币”。
当节点数量由N变为N+1时,只有$\frac{1}{N+1}$的概率跳转到
优点: 内存占用近乎为零,速度极快。
缺点: 不支持手动指定某个节点的权重(除非通过虚拟节点变通),且只能按顺序增减节点。

Google Maglev

这是 Google 的负载均衡器使用的方案,适用于对分布均匀度要求极高的场景。

7. 热点处理

读热点

  1. 增加多级缓存:在Gateway层增加Redis缓存
  2. 动态副本:检测到某个 PG 访问频率极高时,自动在闲置节点上增加临时副本,分担压力。
  3. 使用CDN:边缘分发网络,不要使用从存储系统里面读。
  4. Key salting:对于热点数据,加盐打散,防止数据倾斜,对某些节点压力过大。

写热点

  1. 分片拆分 (Split): 如果某个元数据分片压力过大,触发分片自动拆分。
  2. 异步刷盘与 Buffer: 写入先落到高性能的 WAL(预写日志)或 SSD 缓存层,再异步同步到 HDD 存储层。

8. 失效情景

  1. Node Crash:通过心跳或者故障检测标记dead并触发repair。
  2. Rack Failure:实现rack-aware placement 保证不丢。
  3. 用 3 / 5 node consensus cluster,比如etcd/Raft保证线性一致。

9. 优化点

  1. 如果是小文件较多的系统,多个小文件合并为一个大文件,采用SSTable存储,减少磁盘 IOPS 压力和元数据数量。
  2. 强一致性保证:写入时采用 Quorum 机制(W+R > N),或者依赖 Raft/Paxos 协议同步元数据状态。
  3. 垃圾回收:标记删除。对象删除后并不立即物理删除,而是后台异步扫描删除过期的 Chunk,防止频繁删除导致性能下降。

100TB分布式存储系统设计
https://yicizhang00.github.io/posts/分布式/系统设计/100TB分布式存储设计/
作者
Yici Zhang
发布于
2026年6月23日
许可协议