GFS 是 Google 2003 年发表的论文,是分布式存储系统里较早的一批大规模工业化实践之一。当时的背景是 Google 需要一个适用于各种 application 的 Global 存储系统。

GFS 的目标

  1. 部署在廉价机器,需要自动监控、容错和恢复。
  2. 能够处理海量大文件,量级是数百万个文件、每个通常不小于 100MB,总共数百 TB 且快速增长。
  3. 总是批量读取数据,侧重于顺序读取和 Append 数据,因此没有太大的内存加速需求。
  4. 能够自动让数据在机器之间同步。

Design Overview

Interface

create, delete, open, close, read, write, snapshot, record append.

其中 record append 是重点:多个 client 可以并发向同一个文件追加,不需要额外的锁或者外部协调,GFS 自己保证每条记录的原子写入。多路归并的结果文件和生产者-消费者队列就是靠这个接口实现的。

Architecture

一个 master 和多个 chunkserver,master 包含 file namespace,file 进来的时候会被切分成若干个 chunks,每个 chunk 被分配一个不可变的 64bit 的 chunk handler。

client 可以和 master 交互找到对应的 chunk,然后再与 chunkserver 交互,获取保存在 linux file system 的 chunk data。

chunk size

64MB,比较大,因为面向的是大文件存储,对于热点文件不太友好,因为热点 chunk 需要处理很多节点的请求。

metadata

Master 存储的元数据主要为文件和 chunk namespace,文件到 chunk 的映射,每个 chunk 的 Replica 的位置。前两者可以通过日志备份进行持久化,日志过大后,可以压缩成类似于 B-Tree 的结构,并建立 checkpoint,恢复的时候先把最新的 checkpoint 从硬盘直接映射到内存,再重放这个 checkpoint 之后的那部分日志;而 chunk 的 Replica 的位置信息是不确定的,需要每个 chunkserver 告诉 Master。

Consistency Model

GFS 的一致性是比较宽松的,不保证所有副本在字节上完全一致。论文表 1 把结果分成三种情况:

  • 串行且成功的 write:区域是 defined,读到的就是这次写的完整内容。
  • 并发且成功的 write:区域是 consistent but undefined,各副本内容一致,但那是多个并发写交错拼出来的结果,读不出任何单次写的完整内容。
  • record append(不管串行还是并发):区域是 defined interspersed with inconsistent——每条成功追加的记录本身是 defined,但记录之间夹着 GFS 为了对齐偏移插入的 padding、以及失败重试留下的重复记录,这些夹缝区域是 inconsistent。

失败的变更同样让区域进入 inconsistent:clients 在访问同一个内容的时候可能会看到不同的数据。

容易混的一点:consistent but undefined 只用来说并发 write,不能拿去描述并发 record append。

System Interactions

Leases and Mutation Order

日志的顺序由 Master 指定的 Primary chunkserver 确定,一个写操作的流程如下,主要特点是 数据流和控制流分离

GFS 写入流程中 Client 先向 Master 获取副本位置,再向 Primary 和两个 Secondary Replica 传输数据并发送控制请求
GFS 写入的控制流与数据流。

GFS 的 client 先与 Master 交互,获得 chunkserver 的位置,client 再将数据推送到某个 chunkserver,chunkserver 会继续将数据串行地推送到其他 chunkserver,所有副本接收了数据后,client 会发送写请求给 primary chunkserver,标识之前传输的数据, Primary 会根据这些数据的所有修改根据请求的序列进行串行化处理,按顺序应用到本地,再将这些请求发送到其他 chunkserver。

若过程出现错误,直接返回给 client 让其重试。

Atomic Record Appends

原子化的 append,写入的偏移量由 GFS 自己选定并返回给 client,client 不能指定写到哪里,这一点和普通 write 完全不同。

如果这次追加会超出当前 chunk 的剩余空间,primary 会先把当前 chunk 用 padding 补满,通知所有 secondary 做同样的事,然后让 client 换到下一个 chunk 重试。为了控制 padding 造成的浪费,单次 append 的数据量上限是 chunk 的 1/4,即 16MB。

任一副本写失败时,GFS 返回错误让 client 重试,而重试是在新的偏移上重新追加一次,不是回去覆盖上次的残留。所以失败路径上会留下重复记录和 padding,各副本之间也不是字节一致的——GFS 只保证 at-least-once 的原子追加:成功返回的那个偏移上,每个副本都有这条完整记录。要去重就得靠 client 自己在记录里带上校验和或者唯一 id。

Snapshot

快照能快速地保存某个文件或者目录,实现原理与 COW 类似。Master 收到快照请求的时候,先撤销快照范围内这些 chunk 上已经发出去的租约——目的不是触发数据复制,而是让后续任何对这些 chunk 的写都必须重新回来向 master 申请租约,master 因此拿到了在放行写之前先做 COW 的机会。

撤租约之后,master 把快照操作记进日志,并把这些 chunk 的引用计数加一。等有 client 要写引用计数大于一的 chunk C 时,master 不会直接发租约,而是要求每一个持有 C 副本的 chunkserver 在本地创建一个新 chunk C’——本地复制不走网络,比跨机拷贝快得多。然后 master 把租约发给某个 C’ 副本,client 的写落在 C’ 上,原来的 C 留给快照。

Master Operation

Namespace management and locking

gfs 的锁是针对于命名空间的,并且 master 在操作某些数据前,需要得到对应命名空间的锁的集合。

首先 gfs 在逻辑上会通过前缀压缩路径查找表,master 在获取/a/b/c/d 的写锁时候,需要先获得/a 的读锁,再获得/a/b 的读锁…最后得到/a/b/c/d 的写锁。

这样的一个好处是可以执行某个目录下的并发操作。

gc

文件删除是懒回收的:master 收到删除请求后只是把文件重命名成一个带删除时间戳的隐藏名字,并保留 3 天。这 3 天里文件还能按新名字读到,也可以直接改回原名 undelete 恢复。

超过 3 天后,master 在定期扫描命名空间时才真正抹掉这条元数据,连带断开文件到 chunk 的映射。chunkserver 在心跳里汇报自己持有的 chunk,master 发现某些 chunk 已经不在元数据里,就告诉 chunkserver 这些是孤儿,可以删掉,空间到这一步才释放。