分布式论文阅读 I

GFS、Bigtable、Chubby 论文笔记

前言 #

二十年前,Google 面临的数据规模不断增长,以当时的单机性能,是无法应对巨大的存储需求的,要提升性能,我们只有两种方式:要么纵向扩展,将单机性能做到极致,要么横向扩展,将廉价机器组织成集群

前者的代价是昂贵的,并且边际效应十分明显,这迫使我们选择后者

好在分布式系统的理论已经足够完善,Google 的难点在于将理论变为工程现实,他们成功了,并发表了相关论文

我阅读的三篇论文中,GFS 实现了存储层的分片,Bigtable 实现了逻辑层的分散管理,Chubby 实现了分布式锁(帮助整个集群对运行环境达成共识)

这三个组件位于不同的抽象层,共同组成分布式存储的技术栈

我的写作目的是,记录这些组件的协作关系,并关注核心 idea 与工程优化,分辨出这些 problem 的本质,提取出其中的启发式知识

我不会深入讨论实现细节,下面我们先分别讨论三个组件的内部实现

审视一个系统时,必须先认清 “它提供了怎样的服务”,“如何与外界交互”,带着目的去读设计,这是我的阅读经验

前三部分是我阅读过程中随手记的笔记,语义可能不连贯

一 | GFS #

设计目的 #

GFS 本质是 “用多台机器存储巨大文件”

集群中一定存在机器故障,大文件不利于随机 IO

另外,GFS 的定位是,为上层应用提供可靠的存储抽象

结合这几点,我们先确定了 GFS 的设计特点:

  • 关注容错
  • 倾向于 append write
  • 设计简单的 File API

GFS 其实并没有什么理论上的创新,它只是给出了一个真实的工程实现,证明分布式设计是可用的

底层的分布式抽象是很好理解的

我们将一个巨大的 GFS File,按偏移量映射到不同的 chunk,分散到不同的 chunk-server 中存储,查询时按照索引,读取对应的 chunk,这就实现文件的分片

一旦单个 server 损坏,对应的 chunk 就变得不可用,因此我们在此基础上,对同一个逻辑 chunk 进行多个 server 备份,这就实现了容错

gfs

注意这里的 chunk 是一个逻辑对象,我们并不关心具体会存储到哪台机器,也不关心 server 是哪台机器上的线程

明确一点,GFS 的服务对象是 客户端 Lib,开发者调用客户端库接口 append 进行开发,而对 Lib 而言,GFS 的底层实现是透明的

因此 GFS 可以不保证强一致性,将其转移到 Lib 层进行实现,再进一步由开发者使用

一致性保证 #

GFS 的本质是维护 file offset(filename,chunk) 的映射

Master 就负责管理这些运行环境相关的 metadata,并将 client 的读写指令转发到对应的 chunk

然而网络是不稳定的,并发的写入操作可能会错序到达,在不同的 server 上呈现不同的顺序,造成不一致

为了解决这一点,我们在 chunk 集群中,指定一个 Primary server,由它来决定全局的写入顺序,这样控制流会先到达 primary server,进行排序后,再转发到各个 chunk server


但这么做仍有潜在的问题:不同的 append 写入会混合在一起,这样客户端视角下,offset 就不是连续的了,没有问题吗?

没有问题,前面提到过,客户端 Lib 是了解 GFS 细节的,因此在上层会做 offset 的进一步处理,保证正常读写

另外,如果 chunk 集群在同步过程中失败了,primary server 会通知上层 client 进行重试 append

数据与控制流分离 #

这样 Master 成为了单点瓶颈,如果所有数据都通过 Master 转发,显然会限制性能

因此 GFS 采用数据流解耦的设计

从 Master 角度,处理 Read 指令时,只需返回 chunk server 的地址,由客户端自行与 chunk 集群进行数据通信

从 Chunk 集群角度,Write 同步时,可以按最优顺序在 server 间传输数据,而不需要等待 primary 统一调度

这一点是如何实现的呢?

Write 时,我们实际上是先尽快将 Data 扩散到整个 chunk 集群,等待 Primary 安排好 offset 顺序后,再给对应的 Data 分配位置,准备回复客户端

Lease 设计 #

前面提到的 Master、Primary server 其实都是单主集群,一旦宕机,就无法恢复

因此我们需要设置租约机制(Lease),定期向 leader 发起询问,更新 lease 期限,如果迟迟收不到回复,认为租约超时,指定新的 leader


性能上,客户端可以随时缓存 GFS 的上下文信息,如 Chunk 地址、Primary server 地址


快照 #

GFS 采用 COW 实现快照

创建快照时,我们强制撤销 chunk 现有的 lease,保证该快照对后续所有操作可见,并在 master 标记 快照对应的 chunk 地址

修改快照时,我们创建一个新 chunk 作为副本

二 | Bigtable #

GFS 仅仅是提供了一个抽象文件系统,为了扩展性,Google 以 GFS 为基础,又实现了 Bigtable

Bigtable 采用 K-V 映射作为存储对象,这样就可以根据 key Range 划分到不同的 server 上,实现横向扩展


Bigtable 将数据存储在不同的 Tablet Server 上

与 GFS 类似,Bigtable 也采用单 Master 架构,由 master 管理整个系统的控制流

当 client 发起操作时,master 会告知其对应的 tablet 地址,随后客户端就可以与 tablet server 直接交换数据流,而不需要经过 master


而一个 tablet 地址的定位,则采用三层索引架构,类似于 DNS 服务,查询结果会存入每一层缓存,以加速后续查询


分布式系统中,故障是不可忽略的

我们采用 Chubby 来管理 master 与 server 的存活状态

这些状态十分重要,所以我们才会采用共识算法来实现强一致性


而 Tablet server 则是采用 LSM-Tree 作为存储架构,调用 GFS 提供的接口,将数据持久化到 SSTable,其实现可以参考 Leveldb


GFS 提供物理存储服务,Bigtable 提供逻辑 KV 服务,Chubby 提供 metadata 共识服务

这样我们就将不同的层次解耦了

另外,这些系统运行的基本单位都是“进程”,意味着 GFS 和 Bigtable 进程可能会运行在同一台机器,甚至共用同一个硬盘

但我们不需要关心这些细节,分层抽象的目的,就是实现这种“位置透明”

Tablet 只是进程 #

Tablet server 只是逻辑上管理 KV 的进程而已,并不实际拥有数据

所有的 File 都存在 GFS 里,因此即使 Tablet 所在的机器 crash,我们也可以随时将该 Tablet 管理的 range 分配给其他 server,不需要复制任何东西


所以说,Bigtable 只提供了逻辑上的管理服务,要实现 tablet 分裂,我们只需要在 GFS File 上建立新的索引即可,而不需要修改具体数据

三 | Chubby #

Chubby 作为一个 分布式锁 服务,为 Google 内部的 GFS、Bigtable 系统,提供单 Master 的选举逻辑

以选举 master 为例,要保证 master 唯一,我们只需要在 Chubby 中获取唯一的 lock 即可,这个 lock 必须在整个集群中达成共识

接下来要讲两点,如何管理 lock,如何共识

对于 lock 管理,Chubby 提供 类 Unix 文件系统 接口,也就是 “路径 + 文件”的形式,锁服务就是对 “文件” 进行操作,同时 lock 小文件中还可以存储一些 metadata

对于 共识算法,Chubby 内部的 lock 采用 5 server 的 Paxos,其中 有一台 server 会成为 leader(注:与上面的 master 区分开), client 就是与这台 leader 机器做交互

但是这与外部集群有何关系呢?

Paxos 仅仅保证了 Chubby 内部的共识,与 client 方的 env 显然没有关系,为什么说外部集群也能对这个 lock 产生唯一的共识呢?

这是一个更简单的问题。这里 Chubby 实际担任了第三方服务的角色,也就是说整个集群,都要向 Chubby 这唯一的服务方做查询,这是一个单点服务的场景,既然所有查询都经过同一个节点,那么结果自然是被大家共识的

再讲清楚一点,Paxos 这里只是保证 Chubby 内部数据的”可靠性“,共识只发生在 Chubby 内部,外部的共识靠“单点查询”的系统设计来保证

这就指明了 Chubby 的使用场景:高可靠性

我们只会用 Chubby 来维护一些像是 master 身份的,需要整个集群达成共识的,运行环境,而不能用于存储大量用户数据

Session 维持 #

Chubby 与 client 通过 session 通信,如果有一方突然宕机了,无法发出 close 信号,另一方如何得知呢?

Chubby 选择 lease 机制,双方约定一个固定的时间段,并定期交换 heartbeat RPC,更新 lease 期限,如果 lease 超时,就自动释放 session 资源

这里我们采用 “双向计时”,也就是在 Chubby 与 client 双方,同时维护两个不同的计时器,考虑到 RPC 延迟,client 方的 lease 会更短一些

注意这里的两个计时器,它们不一定是同步的,我们仅在接收 heartbeat 时 对齐时钟


Chubby 方 lease 的作用是,client 失联时,能及时释放 session 资源

client 方 lease 的作用是,判断 Chubby 已经断开 session,不会再发出非法指令

所以 client 的计时会更短,更保守,保证安全性

不过即便 client lease 已经到期,chubby lease 也需要再过一段时间才会释放资源,如果在这段时间内,client 又接收到 chubby 的 heartbeat RPC,说明租约还有效,此时 client 会恢复已经释放的缓存,续约后正常运行

这段真空期,我们称作 jeopardy delay

我们这里用 KeepAlive RPC 替代 heartbeat,前者是阻塞式的,Master 收到 RPC 后不会立刻回复,而是等到租约快到期时才回复,这样就减少了 RPC 发送频率


而在 Master 方,我们也有一段 grace period,租约过期后,会延迟一段时间再释放资源

既然没有收到 RPC,说明已经断开连接了,保留资源又有什么作用呢?

这里的 Master,其实指的是新选举出的 master

新 master 显然不会收到旧的 RPC,当旧 Master 失效后,client 会尝试寻找新的 master,一旦 client 与新 master 建立 session,旧的资源就可以重用了


显然 client 方是知道存在 grace period 的,所以租约过期后,它至少也会等待这段时间结束后,再释放资源

lock 设计 #

chubby 的 lock 是建议性锁,client 可以绕过锁强行访问文件


这么做,放松了对上层读写的限制,需要上层自行 acquire lock 进行管理


另外,我们对每个 lock 维护一个 seq 序号,这样能够阻断乱序到达的指令,明确当前锁的持有者


Chubby 还引入了 lock delay

当 client A 宕机时,它的旧指令可能还在网络中,而 Chubby 会释放 A 对应的锁

旧指令 可能在 没有锁 的情况下执行,导致语义错误

为了处理这种问题,chubby 选择对异常释放的 lock,保留一段“安全期”,限制任何人都不能获取该锁,直到旧指令在网络中消亡(TTL)

这样 其他 client 不会获取到该锁,也不会发出新指令,就不会产生意外冲突


event 机制 #

我们将 KeepAlive RPC 作为定期通信的载体,它会携带 event 信息,完成通知机制

这是一种 Push 机制,不需要 client 显式搭载请求


这么做虽然会有延迟,但 Chubby 的关注点是“可靠性”,而不是性能,因此延迟是可以接受的


缓存设计 #

Chubby 维护的是强一致性缓存,读请求只从 client 缓存中取数据

写操作发生时,Chubby 会用 event 机制,通知所有的 client 清理相关缓存

收到所有 client 的回复后,chubby 才会执行写操作

此后,一旦发生读操作, client 会从 chubby 加载新数据,到本地缓存,然后再返回读请求的结果


这么做,降低了读操作的延迟,将成本分摊到写操作的延迟上

同时,各个 client 间的缓存是 强一致 的,不必担心可靠性


如何扩展规模 #

如果 client 的规模过大,将会产生大量的 KeepAlive RPC,我们可以延长 Master 方的阻塞时长,减少 RPC 频率


为减少 RPC,我们还可以设置一层 proxy,这样 client 与 session 就是多对一的关系,master 只需要维护 proxy 的 session

由 proxy 代理,转发 master 回复的 event 信息即可


另外还可以对 chubby 的文件做 分片,由不同的 master 集群维护

四 | 直觉 #

这里我想整理下前面的各种优化机制,将工程问题 抽象回 分布式系统的固有特征

机制 作用
lease 维持 session、 primary chunkserver
delay 为网络延迟提供安全期
heartbeat 更新 lease,确认存活
缓存 加速 read
seq 为 lock 提供逻辑时钟
event 通知事件,保证缓存一致性
数算分离 分离数据流和控制流

时钟问题 #

机器间的网络延迟是固定存在的,我们必须为延迟提供容错

seq 是由第三方集中分配的时钟,这样各个节点可以对此达成共识

lease 则是以物理时钟为参考,双方在约定的时间段内,承认彼此的存在

delay 牺牲了可用性,资源失效后,仍会再等待一段时间,直到旧指令在网络中全部消亡,确保不会产生冲突

可用性 #

缓存使得读操作不需要再通过网络,自然提升了性能,但也带来了一致性问题

分离数据流,本质是利用了机器间的大量闲置带宽,让 client 与 server 直接通信,而不需要经过 master

组件协作 #

组件 一致性模型 权衡 (Trade-off)
GFS 弱一致性 (Relaxed Consistency) 追求极致的追加写(Append)吞吐量
Bigtable 单行强一致性 方便上层业务逻辑开发,简化事务
Chubby 强一致性 (Paxos) 牺牲性能,换取绝对可靠的“元数据共识”