100TB分布式存储系统设计
前言
系统设计通常从以下顺序考虑
- clarify requirement:清楚约束边界
- high-level architecture:高层设计
- Metadata design:元数据设计
- Sharding:分片
- Replication/EC:副本可用性/数据冗余
- Scaling/rebalancing:重平衡
- HotSpot:热点消除
- Failure scenairo
1. 确认需求
首先清楚约束边界:
- 100TB是当前容量还是未来峰值。
- 对象的平均大小
- 写多还是读多
- 需要跨机房容灾吗
- SLA是99.9还是99.99
假设:
容量:100TB usable
object avg size:16MB
workload:80% read / 20% write
availability:99.99%
durability:11 nines
2. 高层架构
我们设计四层架构
- Client 客户端可以执行命令。
- API Gateway层:负责协议对接,鉴权,请求路由,大文件切片。
1 | |
- Metadata Service:存储对象名称,存储桶,权限,对象到实际物理数据的映射关系,这是控制面。
- Data Node:实际存储的节点,以Chunk或者Slice形式。
3,元数据设计
假设每个分片为16MB,那么100TB/16MB = 6.5M Objects
每个metadata设计为
1 | |
最终元数据大小为
$$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 | |
读取时:
当用户下载这个 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. 热点处理
读热点
- 增加多级缓存:在Gateway层增加Redis缓存
- 动态副本:检测到某个 PG 访问频率极高时,自动在闲置节点上增加临时副本,分担压力。
- 使用CDN:边缘分发网络,不要使用从存储系统里面读。
- Key salting:对于热点数据,加盐打散,防止数据倾斜,对某些节点压力过大。
写热点
- 分片拆分 (Split): 如果某个元数据分片压力过大,触发分片自动拆分。
- 异步刷盘与 Buffer: 写入先落到高性能的 WAL(预写日志)或 SSD 缓存层,再异步同步到 HDD 存储层。
8. 失效情景
- Node Crash:通过心跳或者故障检测标记dead并触发repair。
- Rack Failure:实现rack-aware placement 保证不丢。
- 用 3 / 5 node consensus cluster,比如etcd/Raft保证线性一致。
9. 优化点
- 如果是小文件较多的系统,多个小文件合并为一个大文件,采用SSTable存储,减少磁盘 IOPS 压力和元数据数量。
- 强一致性保证:写入时采用 Quorum 机制(W+R > N),或者依赖 Raft/Paxos 协议同步元数据状态。
- 垃圾回收:标记删除。对象删除后并不立即物理删除,而是后台异步扫描删除过期的 Chunk,防止频繁删除导致性能下降。