Lecture 22: Remote and Distributed File Systems (NFS, AFS, GFS) — 远程与分布式文件系统

目录 · ← l21 · l23 →

Lecture 22: Remote and Distributed File Systems (NFS, AFS, GFS) — 远程与分布式文件系统

讲义对应:CS 425 FA2026 Lecture 22。本章对应课程 Lecture 24「Distributed File Systems」,主要素材为 L24.A.FA25.pdf(29 页:文件系统基础、Unix 文件 API、Vanilla DFS 的三服务架构、NFS、AFS);并整合 Lecture 24-B「Consistency Models」L24.B.FA25.pdf:一致性谱系、线性一致/顺序一致/因果一致/最终一致、会话级模型 MR/MW/RMW)作为 NFS/AFS/GFS 一致性口径的前置,Lecture 19-20「RPC」L19-20.FA25.pdf:Sun RPC、at-least-once 与 at-most-once、幂等操作)作为 NFS/AFS 每一次 read/write 都是 RPC 的前置,Lecture 9-11「Cassandra」L9-11.FA25.pdf:最终一致、quorum、read repair、CAP)作为一致性谱系另一端的参照,Lecture 4「MapReduce and Hadoop」L4.FA25.pdf:GFS/HDFS 作为 MapReduce 的存储底座、3 副本、chunk 大小、数据本地性、2-1 机架副本布局)作为 HDFS 生态角色的前置。GFS 部分(22.2.13-22.2.20、22.3.3-22.3.6、22.4.1)以 Ghemawat et al. 的 SOSP 2003 论文与 HDFS 论文为主干——课程讲义只把 GFS/HDFS 当作 MapReduce 的分布式文件系统一笔带过,本篇按论文与工程实践补齐。 教材对应:Coulouris 5th Ed. Ch. 12(Distributed File Systems):12.1 文件系统与分布式文件系统的动机、12.2 文件服务架构(flat file service / directory service / client service,即本讲的 “Vanilla DFS”)、12.3 Sun NFS、12.4 Andrew File System、12.5 近期进展(缓存粒度、回调与租约);前置:Ch. 4/Ch. 5 的 RPC 与远程调用语义、Ch. 6 的一致性模型。 阅读材料:Sanjay Ghemawat, Howard Gobioff, Shun-Tak Leung, The Google File System, SOSP 2003;John H. Howard et al., Scale and Performance in a Distributed File System, ACM TOCS 1988(AFS 论文);Russel Sandberg et al., Design and Implementation of the Sun Network Filesystem, USENIX 1985;Konstantin Shvachko et al., The Hadoop Distributed File System, MSST 2010;RFC 3530(NFS Version 4 Protocol);可选:Satyanarayanan, The Evolution of Coda, ACM TOCS 2002;Weil et al., Ceph: A Scalable, High-Performance Distributed File System, OSDI 2006。

22.1 概述

本章回答一个具体而根本的问题:当文件不在本地磁盘上,而在另一台(或很多台)机器上时,操作系统应该给它什么接口、什么缓存、什么一致性? 讲义先造了一个假想的 Vanilla DFS,把分布式文件系统拆成三个角色(Flat file service / Directory service / Client service),并在设计过程中当场做出两个影响此后四十年的决策:接口用”绝对偏移”而不是”文件描述符”(换取幂等性,从而可以安全地按 at-least-once 语义重试),以及服务器不保存打开文件表(换取无状态,从而崩溃后无需恢复)。这两个决策像两条分岔的路标,指向此后所有真实系统的取舍:

  • Sun NFS(1980s,Sun Microsystems,今天仍在用)选了”无状态 + 块级远程访问 + 打开时验证“:服务器极简、崩溃恢复极快,代价是一致性弱(只保证 close-to-open),且锁与缓存失效必须靠旁路协议(NLM/NSM)补救。
  • AFS(CMU,Andrew File System)选了”整文件缓存 + 服务器回调(callback promise)“:用服务器的状态换更强、更快的一致性,并把服务器负载从”块访问次数”降到”打开/关闭次数”,从而支撑数以万计的客户端。
  • GFS / HDFS(Google 2003 / Apache Hadoop)干脆换了一个接口:不做通用文件系统,只服务”少量超大文件、一次写多次读、追加(append)为主、批量吞吐优先“的工作负载,用 64 MB 大 chunk + 单 master 元数据 + 主副本租约(lease)+ 数据 pipeline + 松弛一致性(relaxed consistency) 换到数千节点的吞吐。

本章在课程中的位置很清楚:向下依赖 Lecture 19-20 的 RPC(NFS/AFS 的每一次 read/write 本质上都是一次 RPC,因此 at-least-once 与幂等性直接决定语义);横向与 Lecture 24-B 的一致性谱系对接(NFS 的 close-to-open 落在谱系的弱端,AFS 用回调换到稍强,GFS 主动退到”松弛一致”);向上是 Lecture 4 的 MapReduce、HDFS、以及今天所有大数据存储的地基。

一句话概括贯穿全章的黄金法则(22.6 再次点题):

分布式文件系统的设计由工作负载假设决定:NFS 假设通用工作负载 ⇒ 选择无状态 + 块级远程访问;AFS 假设读多写少的中小文件 ⇒ 选择整文件缓存 + 回调;GFS 假设少量巨型文件的追加写 ⇒ 选择大 chunk + 单 master 元数据 + 松弛一致性。三者没有优劣,只有假设是否匹配。

22.2 核心概念与分布式机制图解

22.2.1 文件系统到底提供什么(What a File System Provides)

  • 定义与目的文件系统(file system) 是比”磁盘块”更高一层的抽象:它把裸块设备包装成文件(file)目录(directory/folder),让用户和进程不必直接和磁盘块、内存块打交道。它一次性提供了四件事:命名(命名空间)访问控制存储分配(哪些块属于哪个文件)元数据管理(属性);在本章关心的分布式场景里还要加上第五件:缓存与一致性

  • 一个文件里有什么:文件 = 头部(header)+ 数据块(Block 0 … Block N-1)。头部(inode)保存属性:

    属性说明
    时间戳creation / read / write / header 修改时间(NFS 的 close-to-open 一致性完全建立在这几个时间戳上)
    文件类型例如 .c.java,决定由哪个程序解释内容
    属主(ownership)例如 edison
    访问控制表(ACL)谁能以什么模式(读/写/执行)访问这个文件
    引用计数(reference count)有多少个目录项指向这个文件;可以 > 1(硬链接),减到 0 时文件才可以被删除
  • 目录也是文件:目录”就是文件”,只不过它的”数据”是它所包含文件的元信息以及指向这些文件在磁盘上位置的指针。理解这一点很重要:目录操作(lookupunlink)本质上就是对一种特殊文件的读写,因此分布式文件系统可以把”目录服务”和”平坦文件服务”分开实现(22.2.4)。

  • Unix 文件 API(讲义原口径)——分布式文件系统的接口设计几乎全部由此派生:

    调用语义关键点
    filedes = open(name, mode)打开文件,返回文件描述符访问前必须打开;内核为 fd 建内部数据结构
    filedes = creat(name, mode)创建文件并返回 fdmode = r / w / x
    close(filedes)关闭释放 fd(在 NFS 中,”close” 还是一致性协议的一个关键时刻)
    error = read(filedes, buffer, num_bytes)读写指针当前位置读读完自动前进 num_bytes;返回字节数 / 0 (EOF) / −1 (错误,errno)
    error = write(filedes, buffer, num_bytes)写入指针位置同样自动前进
    pos = lseek(filedes, offset, whence)移动读写指针whence 决定绝对还是相对
    link(old, new) / unlink(old)增加/减少硬链接link 使引用计数 +1;unlink −1,为 0 才删除
    stat/fstat(name, buffer)取文件的头部(属性)NFS 中的 GETATTR 就是它的远程版本
  • 机制图解:文件、目录、inode 与读写指针的关系。

      目录也是文件:目录里的"数据"是 (名字 -> inode) 的映射
   +-------------------+          +--------------------------------------+
   |dir:/usr/edison    |          | inode #4711  (header)                |
   |---                |          | type=.c  owner=edison refcount=2     |
   |"a.c"   -> #4711   |--------->| atime / mtime / ctime                |
   |"my_inv"-> #4712   |          | ACL: edison:rw, bob:r                |
   +-------------------+          +--------------------------------------+
                                             |  指向
                                             v
   +---------------------------------------------------------------------+
   | 进程视角:fd = 3  -->  [ 文件表项: 读写指针 offset=512, mode=r ]    |
   | read(fd, buf, 128)  ==>  读 [512,640),然后 offset 自动变成 640     |
   | ("自动前进"使 Unix 的 read/write 不是幂等的 —— 见 22.3.1)         |
   +---------------------------------------------------------------------+
  • 讲义的三个思考题(先想再看答案)
    1. 调用 lseek 会引发一次磁盘寻道吗? 不会。lseek 只修改内核中该 fd 的读写指针(内存里的文件表项),既不读也不写数据,因此不产生磁盘 I/O;真正的寻道发生在随后的 read/write 访问到新位置时。这正是一个”元数据操作 vs 数据操作“的区分——分布式文件系统里同样的区分决定了请求走 master 还是走数据服务器(22.2.14)。
    2. 文件被删除后还能存在硬链接吗? 不能。硬链接是”目录项 → inode”的绑定,unlink 使引用计数递减,计数归零时 inode 被回收;此时任何指向它的名字都已不存在,也就不可能有硬链接存在。反过来说,只要还有硬链接(计数 > 0),文件就没被真正删除
    3. 符号链接(symbolic/soft link)呢? 可以。符号链接是”另一个文件,内容是路径名”,它不增加引用计数。目标被删除后,符号链接仍然存在,只是解析失败(悬空链接 dangling link)。
  • 关键假设与系统模型:以上都是单机、磁盘可靠、内核可信的世界。本章要做的事就是把这个模型搬到网络上,于是每一条假设都要重新谈判:网络会延迟和丢包(Lecture 19-20)、机器会崩溃(Lecture 6)、不同客户端会并发写同一份数据(Lecture 24-B)。

22.2.2 为什么需要分布式文件系统(Why a Distributed File System)

定义与目的分布式文件系统(DFS, Distributed File System)文件存放在服务器机器上,客户端机器通过对服务器做 RPC 来完成文件操作,同时让客户端”感觉不到”文件是远程的。讲义把它归结为三个”期望属性”:

  • 透明性(Transparency):客户端访问 DFS 文件就像访问本地(Unix)文件——同一套 API,客户端代码不修改;位置、复制等细节对客户端不可见。
  • 支持并发客户端(Concurrent clients):多个客户端进程可以并发读写同一个文件。
  • 复制(Replication):为了容错(也为了读扩展)。

而在工程实践中,引入 DFS 的动机可以拆成五条,每条对应一个不同的技术抓手:

动机想要的东西主要技术抓手代价
共享数据多用户/多机器协作同一份数据中央文件服务 + 挂载一致性(缓存失效)
容量扩展单机磁盘装不下 / 加盘要停机多服务器条带化、chunk 化元数据管理
可用性与容错机器坏了数据还在、服务不断副本、校验和、自动修复存储成本 ×3、恢复带宽
性能扩展并发访问的聚合吞吐随节点数增长数据与元数据路径分离、pipeline一致性模型被迫松弛
集中管理备份、配额、审计、访问控制一处生效单一命名空间 + 集中认证元数据服务器成为关键路径

讲义的三个核心概念(必须区分清楚)

  1. One-copy update semantics(单副本更新语义):当文件有多个副本时,从客户端角度看,其内容与”只有一个副本”时没有区别。这是所有复制文件系统追求的”理想一致性”,也是绝大多数系统只能近似实现的目标——NFS 放弃了它(没有副本,只有缓存),AFS 用回调近似它,GFS 明确放弃它(松弛一致性)。
  2. At-most-once vs At-least-once 操作语义(从 Lecture 19-20 直接继承):

    语义含义适用例子
    At-most-once操作最多执行一次不能重复的操作append(追加重复 ⇒ 数据错)、非幂等的 x = x + 1
    At-least-once操作至少执行一次(可能重复执行)幂等(idempotent)操作绝对位置读文件、把某位置写成定值

    讲义特别强调”Choose carefully“:NFS 之所以敢于在超时后无脑重试,正是因为它把接口设计成幂等的(绝对偏移 + 无读写指针)。而 AFS 与 GFS 的 append 天生不幂等,所以它们必须用别的手段(回调、序列号 + at-least-once 语义 + 应用层去重)来处理重试。

  3. 安全(Security in DFS)认证(authentication) = 验证”你确实是你说的人”;授权(authorization) = 验证”这个人能不能访问这个文件、以什么模式”。两种流行的授权表示:

    方式组织方式优点缺点
    访问控制表(ACL)每个文件一份”允许谁、以什么模式”的列表撤销容易(改文件);符合 Unix 习惯文件多时存储/管理开销大;检查时要找到文件
    能力列表(Capability list)每个用户一份”能访问哪些文件、什么模式”的列表检查快(查用户自己的列表)撤销困难(要追回能力票据);可被转发

    工程上的做法是两者结合:把权限切成一个个 capability,每个 (user, file) 一份;认证由 Kerberos/TLS 完成(NFSv4 集成了 Kerberos),授权由服务器逐请求检查(因为无状态服务器不能记住”你已经通过检查了”,这一点在 22.3.1 会再出现)。

  • 直观解释(”它是什么?”):本地文件系统像自己家的书架——拿书不用打招呼,但也只有你能用。分布式文件系统像图书馆 + 借书证:书在馆里(服务器),你在家看书(客户端缓存),图书馆负责登记谁借了什么(元数据)。三位主角的差别恰好是三种”借书规矩”(22.2.22 展开):NFS 是”每次去图书馆抄一页”(块级远程访问),AFS 是”把整本书借回家、有人要改书就打电话叫你送回来”(整文件缓存 + 回调),GFS 是”书切成三份存三个库房,只有一个库房负责登记改动”(大 chunk + 主副本租约)。

22.2.3 DFS 的七个核心设计问题(贯穿全章的主线)

任何分布式文件系统都必须回答下面七个问题;本章的三位主角给出的答案几乎处处相反,而分歧的根源都是对工作负载的不同假设。

#设计问题问题的两端NFSAFSGFS/HDFS
1接口模型上传/下载整个文件(upload/download) ↔ 细粒度远程访问(remote access)远程访问(read/write 字节区间)混合:接口是远程访问,实现是整文件上传/下载远程访问((handle, offset, len)),但语义按”追加”优化
2状态性有状态(服务器记录打开文件/锁/缓存) ↔ 无状态无状态(+ NLM/NSM 旁路)有状态(回调承诺表)master 有状态(元数据、租约),chunkserver 几乎无状态
3缓存在哪里缓存(客户端/服务器/两者);粒度(整文件/块);一致性机制(验证/回调/租约)客户端块级缓存 + 打开时验证;服务器也缓存客户端整文件缓存(持久化)+ 服务器回调失效客户端不缓存数据(只缓存元数据);服务器用 OS 的页缓存
4命名位置透明 ↔ 位置独立位置透明(挂载点后无服务器名),不独立(fh 绑死服务器)位置透明且位置独立/afs/cell/...,卷可迁移)位置透明(/user/log),复制位置对客户端不可见
5复制单副本(靠缓存提性能) ↔ 多副本(靠复制提可用性)早期基本无副本卷级只读副本每 chunk 3 副本,机架感知(2+1)
6并发控制无(应用自己解决) ↔ 服务器提供原子性与锁弱:无原子性,锁靠 NLM弱:无原子性,应用层锁Record Append 提供至少一次原子追加;其余靠单写者约定
7性能与可扩展性元数据瓶颈 ↔ 数据瓶颈服务器数据路径是瓶颈服务器请求数随打开/关闭次数增长master 只走元数据 ⇒ 数据路径可扩到数百 chunkserver

四种”缓存 × 接口模型”组合(这张表把设计空间闭合起来,本章后面三个系统分别落在其中三格):

 客户端缓存(靠近应用,命中 = 0 RTT)服务器缓存(靠近磁盘,用于摊薄磁盘 I/O)
下载/上传模型(整文件搬运)AFS:整文件 + 持久缓存 + 回调承诺;打开时取、关闭时写回AFS/Vice 在内存中缓存热文件与目录;GFS 的 chunkserver 依赖 Linux 页缓存
远程访问模型(细粒度读写)NFS:8 KB 块级缓存 + 属性缓存(Tc/Tm/t)+ 打开时验证;GFS/HDFS 客户端不缓存数据(只缓存元数据)NFS 服务器:delayed write(约 30 s 刷盘)write-through;GFS:chunkserver 的 Linux buffer cache
  • 关键假设与系统模型:本章默认 crash-stop / crash-recovery 故障模型(机器会崩溃并重启,但不会说谎),网络可能丢失、重复、乱序(因此需要至多/至少一次语义与幂等性),没有全局时钟(因此租约要用本地时钟近似,并假设时钟漂移率有界)。拜占庭故障不在本章范围内(详见 Lecture 26)。

22.2.4 Vanilla DFS:三个角色与 API(讲义的”最小可用 DFS”)

  • 定义与目的:讲义先不碰 NFS/AFS 的真实细节,而是设计一个”原味”的分布式文件系统(Vanilla DFS),用来把职责边界接口形态这两件事讲清楚。它由三类进程组成:
        +-------------------------------------------------------------+
        |                      SERVER 机器                             |
        |                                                             |
        |   +---------------------+        +----------------------+   |
        |   |  Flat file service  |<-------|  Directory service   |   |
        |   |  (文件内容 + 属性) |  是它的 |  (目录树、名字解析) |   |
        |   |  按 file_id 操作     |  客户  |  lookup/add/un_name  |   |
        |   +---------------------+        +----------+-----------+   |
        +----------------------------------------------|--------------+
                                                       |  RPC
        +----------------------------------------------|--------------+
        |                      CLIENT 机器              v              |
        |                        +--------------------------+         |
        |   Process  -----------> |      Client service      |         |
        |   (用 fd 访问文件)      | (把本地 fd 翻译成 RPC)  |         |
        |                        +--------------------------+         |
        +-------------------------------------------------------------+
  • Flat file service API(讲义原口径)

    调用语义为什么这样设计
    Read(file_id, buffer, position, num_bytes)绝对位置 positionnum_bytesbuffer没有自动前进的读写指针 ⇒ 同一请求重复执行结果相同 ⇒ 幂等 ⇒ 可以用 at-least-once
    Write(file_id, buffer, position, num_bytes)在绝对位置写同上,重复写同一段内容结果相同
    create/delete(file_id)创建/删除需要目录服务配合(引用计数)
    get_attributes/set_attributes(file_id, buffer)读/写头部这就是 NFS 的 GETATTR/SETATTR

    关键设计点file_id 不是文件描述符,而是”这个文件在那个文件系统里的唯一 ID“。为什么不用 fd? 讲义给出两条环环相扣的理由:

    1. 需要操作幂等(从而能用 at-least-once 语义安全重试);
    2. 需要服务器无状态(不保存打开文件表),因为崩溃后没有状态需要恢复——重启即恢复服务。 对比之下,Unix 的文件系统操作既不幂等(读写指针会前进)也不是无状态的(内核有打开文件表)
  • Directory service API

    调用语义
    file_id = lookup(dir, file_name)名字 → file_id(之后拿它去 flat file service 访问文件)
    add_name(dir, file_name, buffer)增加目录项,引用计数 +1
    un_name(dir, file_name)删除目录项,引用计数 −1,为 0 才可删文件
    get_names(dir, pattern)类似 ls -al 后面接 grep/find

    注意 Directory serviceFlat file service客户(它把自己的目录内容当作文件来读写),这是一种漂亮的递归分层:目录服务不需要自己实现持久化。

  • 直观解释(”它是什么?”):Vanilla DFS 就像一个只有柜台、没有记忆的档案馆:你把”哪一排哪个格子”(file_id + position)写在申请单上交给柜台,柜台照单取件,它不记得你昨天来过、也不记得你上次读到哪。好处是柜台工作人员随时可以被替换(崩溃重启无感),坏处是任何需要”记住上下文”的功能(锁、缓存一致性、顺序读优化)都得另设一个柜台——这正是 NFS 后来的处境(NLM 锁管理器、NSM 状态监视器)。

  • 关键假设与系统模型:可靠性交给 RPC 层(at-least-once + 幂等);服务器无状态意味着不支持服务器端的打开语义与锁;引用计数意味着删除是”最后一个 unlink”触发的,这与 Unix 一致。

22.2.5 NFS 架构:VFS + RPC + 挂载

  • 定义与目的NFS(Network File System) 由 Sun Microsystems 在 1980 年代提出,今天仍被广泛使用(Linux 的 nfs、企业 NAS、Kubernetes 的 nfs volume 都是它)。它的目标就是讲义列出的”透明性”:本地文件与远程文件在客户端看起来完全一样

  • 机制图解:NFS 的整体架构(对应讲义中的 NFS Architecture 图,并标注了缓存所在的位置)。

   CLIENT 系统                                                     SERVER 系统

+-------------------------------------------+                    +-------------------------------+
|Process               Process              |                    |(无状态:不保存打开文件表)     |
|   |                     |                 |                    |                               |
|   v                     v                 |  === Sun RPC ===>  |Virtual File System (VFS)      |
|Virtual File System (VFS)                  |      + XDR         |(服务器侧同样是 VFS)           |
|- 每个挂载点一个 VFS 结构                  |   NFS 协议         |                               |
|- 每个打开文件一个 v-node                  |   (mount 协议     |        |                      |
|- 远程文件:v-node 里存服务器地址          |     独立)         |        v                      |
|  + file handle                            |                    |Unix File System / 本地磁盘    |
|        |                    |             |                    |[服务端缓存:delayed write     |
|        v                    v             |                    | 或 write-through(见 22.2.7)]|
|Unix 本地 FS         NFS 客户端模块        |                    |                               |
|(本地磁盘)           + 块级缓存 8 KB       |                    |(服务器导出目录 / export)    |
|                     + 属性缓存(新鲜期 t)|                    +-------------------------------+
|                                           |                                                     
|[客户端缓存:块 8 KB,标签 Tc / Tm]        |                                                     
+-------------------------------------------+                                                     
  • 客户端系统(NFS client):相当于 Vanilla DFS 里的 Client service,但它集成在内核里(因此对应用完全透明),通过 RPC 向服务器发请求。它维护 VFS 层
    • 为每个挂载的文件系统保存一个 VFS 结构
    • 为每个打开的文件保存一个 v-node:若文件在本地,v-node 指向本地 inode(磁盘块);若是远程文件,v-node 里存的是远端 NFS 服务器的地址 + 文件句柄(file handle)
    • 对每次文件访问决定路由:走本地文件系统还是走 NFS 客户端模块。
  • 服务器系统(NFS server):同时扮演 Vanilla DFS 里的 Flat file service + Directory service 两个角色,并允许挂载(mount)文件与目录。讲义的例子:

    /usr/tesla/inventions 挂载到 /usr/edison/my_competitors,于是 /usr/edison/my_competitors/foo 就是指 /usr/tesla/inventions/foo挂载不克隆(不拷贝)文件,只是”让这个目录现在指向那个目录”。

  • 两套协议分离mount protocol(客户端问”我能挂载哪个目录?给我它的根 file handle”)与 NFS protocol(真正的 read/write/lookup/getattr)是分开的。这个分离很实用:挂载是低频的管理动作,可以走独立的认证与访问检查(/etc/exports 的 export 列表就是服务器端的访问控制),而 NFS 协议则要在数据路径上尽可能轻。

  • 讲义的思考题:”既然进程用 file descriptor 访问文件,NFS 服务器岂不是有状态了?” 答案:没有。 fd客户端内核本地的概念——VFS 把 (fd, offset) 翻译成 (file handle, absolute offset, length) 再发出去。服务器只看到 file handle,它不保存”谁打开了哪个文件”“读到哪儿了”“谁持有锁”这类跨请求状态。唯一的例外正是这套设计的代价:锁与缓存一致性这类”必须记住状态”的功能只能外挂(NLM/NSM),并且在服务器崩溃后必须重建(22.2.7)。

  • 关键假设与系统模型:NFS 假设客户端与服务器都支持同一套 VFS 抽象(所以它可以跨 Unix 变体、甚至跨操作系统工作);假设网络 RPC 语义是 at-least-once(Sun RPC 默认),因此所有 NFS 操作必须幂等或至少可容忍重复。

22.2.6 NFS 的文件句柄与路径解析(无状态设计的核心)

  • 定义与目的文件句柄(file handle, fh) 是 NFS 无状态设计的支点:讲义给出它是 32 字节的 token,包含 (volume, inode, generation) 三部分。
字段含义为什么需要
volume / file system ID文件位于服务器上的哪个导出文件系统服务器可以导出多个文件系统,需要区分
inode 号文件在该文件系统内的 inode 编号定位文件本体(但它是”可以指向同一 inode 的多个文件”之一)
generation number(生成号)inode 被复用的代数(每次 inode 被释放后重新分配就 +1)防止”陈旧 fh 指向新文件”:客户端拿着老 fh 去访问,服务器发现 generation 不匹配就返回 ESTALE
  • 机制图解:一次路径遍历的 RPC 开销。客户端要打开 /usr/edison/my_competitors/foo,而它只知道挂载点的根 fh:
   CLIENT                                            SERVER
     |  LOOKUP(fh_root, "usr")                        |
     |----------------------------------------------->|  解析一级目录
     |<-----------------------------------------------|  fh_usr
     |  LOOKUP(fh_usr, "edison")                      |
     |----------------------------------------------->|
     |<-----------------------------------------------|  fh_edison
     |  LOOKUP(fh_edison, "my_competitors")           |
     |----------------------------------------------->|
     |<-----------------------------------------------|  fh_comp
     |  LOOKUP(fh_comp, "foo")                        |
     |----------------------------------------------->|
     |<-----------------------------------------------|  fh_foo(终于拿到)
     |  OPEN / GETATTR(fh_foo) / READ(fh_foo, ...)    |
     |----------------------------------------------->|  之后都直接用 fh
     时间:4 个 RTT 起的“路径解析税”,每一级都是一次往返

这是无状态 + 路径依赖的直接代价。NFSv4 用 COMPOUND(复合操作) 把多个 LOOKUP + OPEN 打包成一个 RPC 来缓解(22.2.7)。

  • 无状态带来的三个好处与三个代价

     内容
    好处 1服务器崩溃后重启不需要恢复任何客户端状态——客户端只要重试即可(配合幂等 + at-least-once)
    好处 2服务器不需要为每个客户端分配内存,天然可支持大量客户端(没有”连接状态”的规模上限)
    好处 3实现简单、健壮,容易跨平台(这也是 NFS 能活到今天的原因)
    代价 1路径依赖 + 重复解析:每次请求都要带 fh;服务器无法”记住这个 fh 对应哪个文件”(多数实现靠一个 fh→vnode 的缓存来加速,但那是优化,不是语义)
    代价 2无法维护跨请求状态文件锁(需 NLM,Network Lock Manager)、状态监视(需 NSM,Network Status Monitor,用于崩溃后通知持锁方)、缓存一致性检查(每次都得客户端自己验证)都只能外挂
    代价 3fh 复用风险:文件被删除后 inode 被回收再利用,老 fh 会指向一个”新文件” ⇒ 用 generation number 使老 fh 立即失效(返回 ESTALE),把”读到错误数据”降级为”报错”
  • 关键假设与系统模型:NFS 假设客户端愿意承担路径遍历的延迟(用缓存摊薄);假设服务器不主动通知客户端(因此一致性只能是”客户端主动验证”式的);假设 at-least-once RPC(因此所有操作幂等)。

22.2.7 NFS 的缓存与 close-to-open 一致性

  • 定义与目的:NFS 的性能几乎全部来自缓存,而它的一致性弱点也几乎全部来自缓存。缓存存在于两侧:服务器侧缓存用于摊薄磁盘 I/O,客户端侧缓存用于避免网络往返

  • 服务器侧缓存(Server caching):把最近访问的文件块与目录块放在内存里。讲义强调这是 NFS 读取性能好的主要原因之一——因为”人写的程序有访问局部性(locality of access)”:最近访问过的块,很可能马上还会被访问。写入有两种风格:

    服务器写策略行为优点缺点
    Delayed write(延迟写)先写内存,每约 30 秒(例如 Unix sync)批量刷盘(客户端很快收到 ack)不一致:服务器崩溃可能丢掉最近 30 秒已 ack 的写
    Write-through(写穿)落盘后才 ack 客户端一致/持久(每次写都吃磁盘延迟)
  • 客户端侧缓存(Client caching)与验证规则:客户端缓存最近访问的块,每个缓存块打两个时间戳标签

    标签含义
    $T_c$该缓存项最近一次被验证(validated) 的时间
    $T_m$该块在服务器上最近一次被修改的时间

    于是服务器返回的属性里带着 mtime,客户端可以判定:在时刻 $T$,一个缓存项有效当且仅当

\[T - T_c < t \quad \text{(还在"新鲜期"内)} \qquad \text{或者} \qquad T_m^{\text{client}} = T_m^{\text{server}} \quad \text{(服务器上的修改时间没变)}\]

其中 $t$ 是新鲜期(freshness interval),是一致性与效率之间的折中:讲义给出 Sun Solaris 的实现是文件自适应取 3-30 秒、目录取 30-60 秒(现代 Linux 对应 acregmin/acregmaxacdirmin/acdirmax)。块被写时,客户端做 delayed write:先改本地缓存,稍后(或关闭时)批量发回服务器。

  • close-to-open consistency(关闭-打开一致性)——NFS 一致性机制的核心

    机制(两句话):open 时验证、close 时写回

    • open:客户端对文件做一次 GETATTR,把服务器上的 mtime/size 与自己缓存中的属性比较;若失效,则丢弃该文件的缓存数据块(之后的读都走服务器,或重新拉取)。
    • close:客户端把该文件所有脏块写回服务器(并在服务器上更新 mtime)。

    它保证什么:如果客户端 A 写完文件 F 并 close,那么之后任何客户端 B 打开 F,都能看到 A 的写。这就是”close-to-open”这个名字的来源——写方的”关闭”与读方的”打开”之间建立了一个一致性同步点

    机制图解(时序)

   客户端 A            服务器            客户端 B
      |                  |                  |
      | open(F)          |                  |
      |---GETATTR------->|                  |
      |<--attrs(mtime=10)|                  |
      | write(F, 0, 8KB) |                  |
      |  (只改本地缓存)    |                  |
      | close(F)         |                  |
      |---WRITE×N------->|  mtime := 11     |
      |<----ok-----------|                  |
      |                  |<---GETATTR-------|  open(F):验证属性
      |                  |----mtime=11----->|  与本地缓存 mtime=10 不一致
      |                  |                  |  ⇒ 丢弃旧缓存
      |                  |<----READ---------|  重新从服务器读
      |                  |-----A 的数据---->|
      |                  |                  |  ✅ B 看到了 A 已关闭的写

它不保证什么(关键!)对已经打开的文件,多个客户端可能看到不一致的数据——因为 NFS 不会主动失效别人缓存里的块,A 在自己的 open 会话里会一直用自己缓存的旧数据。所以 NFS 提供的是弱一致性:没有原子性保证,不适合多写者共享同一个文件(尤其是并发写同一区域)。

  • 具体的反例(本章会反复用到):A 与 B 都打开 F 的同一区域:
    1. A open(F),读到 block 3 = AAAA(缓存住);
    2. B open(F),写 block 3 = ZZZZclose(F)(服务器 mtime 变了);
    3. A 仍在同一个 open 会话里读 block 3 ⇒ 仍返回 AAAA(陈旧!)——A 没有任何机制知道别人改了;
    4. A close(F) 后重新 open(F) ⇒ 这次才看到 ZZZZ。 (22.4.2 的代码把这个反例跑了出来的:A.stale = 1,重新打开后才读到 B 的写。)
  • NFS 的版本演进(把”什么时候被放松的”讲清楚)

    版本关键变化对一致性的影响
    NFSv2(1985)32 位偏移(文件 ≤ 2 GB)、只支持 UDP、读写 8 KB 上限close-to-open + 属性缓存新鲜期(默认 3 s 起)
    NFSv3(1995)64 位偏移、TCP、异步写(async write)write-through 两种服务器策略、更细的属性验证、可调的 acregmin/acregmax/acdirmin/acdirmaxCOMMIT 操作属性缓存时间越长 ⇒ 请求越少但一致性窗口越大(可调旋钮)
    NFSv4(2000/2003)有状态:引入 OPEN/CLOSE 状态与租约(lease)回调(callback) 用于缓存失效;COMPOUND 复合操作减少 RPC 次数;集成锁(不再需要 NLM/NSM);集成安全(Kerberos);单一端口 2049 穿越防火墙从”客户端拉取验证”走向”服务器推送失效”——这正是 AFS 的路线(22.2.9),说明 NFSv4 承认了回调更优
  • 关键假设与系统模型:close-to-open 的保证依赖于”open 时确实做了一次验证”:如果属性缓存还在新鲜期内(默认最小值 3 秒),某些实现会跳过 GETATTR,此时”B 在 A 关闭后 3 秒内打开”就可能读不到 A 的写。所以严格来说,NFS 的保证是”close-to-open + 属性缓存新鲜期”的联合条件,这也是生产上把 actimeo 调成 0 来换强一致的做法(代价是请求数暴涨)。此外,mtime 只有秒级分辨率(NFSv3 的 mtime 是秒 + 纳秒字段,但很多实现只填秒),同一秒内的两次修改可能无法区分。

22.2.8 AFS:设计目标与”整文件”哲学

  • 定义与目的AFS(Andrew File System)CMU 设计(名字取自 Andrew Carnegie 与 Andrew Mellon——CMU 里的 “C” 与 “M”),目标是可扩展性(scalability):支持数以万计的客户端。今天它仍用于一些集群(尤其是大学集群),开源继承者是 OpenAFS

  • 两个”不寻常”的设计决定(讲义原文:Two unusual design principles)
    1. Whole file serving(整文件服务)不按块传,而是整个文件传
    2. Whole file caching(整文件缓存):客户端把整个文件缓存到本地磁盘,而且这个缓存是永久的(permanent,重启后仍在)(AFS 的典型缓存大小是 100 MB 量级)。
  • 它建立在一组”(经过验证的)假设”上——这组假设正是 AFS 与 NFS 分野的根源:

    假设依据推论
    多数文件访问是单一用户的实测:绝大多数文件只有一个用户读写缓存冲突少,整文件缓存值得
    多数文件很小实测文件大小分布长尾在几 KB~几十 KB整文件传输代价可接受
    100 MB 的客户端缓存是可承受的当时的机器内存/磁盘已能容纳大部分读能在本地命中
    读远多于写,且通常是顺序读实测读写比大约 6:1 以上打开时一次取回、之后零网络流量,收益极大
  • 直观解释(”它是什么?”):如果把 NFS 想成”每次去图书馆抄一页“,AFS 就是”把整本书借回家看,图书馆如果有人要改这本书,就打电话让你把手里的版本作废“。对于”读多写少、文件不大”的典型工作负载,这个策略的服务器负载与客户端规模几乎无关——因为客户端一旦拿到文件,之后无论读多少遍都不产生任何网络流量

  • 关键假设与系统模型:AFS 假设客户端有可用的本地磁盘(缓存要持久化);假设文件不太大(大文件整传代价高);假设写不频繁(回调流量可控)。

22.2.9 AFS 架构:Venus / Vice / Cell / Volume 与回调承诺

  • 定义与目的:AFS 把系统分成客户端进程 Venus服务器端 Vice,并用 Cell(单元)Volume(卷) 两个概念来组织管理域。

  • 机制图解:AFS 的整体结构与回调流程。

   CLIENT 机器(Venus)                              SERVER 机器(Vice)

+--------------------------------+                +---------------------------------+
|Process                Process  |                |File Server(文件服务器)        |
|   |                      |     |                |- 文件本体 + 版本号(持久化)    |
|   v                      v     |   Fetch        |- callback 表:file -> {哪些客   |
|Venus(客户端进程)             |   Store        |  户端缓存了它}(内存,崩溃即丢)|
|- 整文件缓存(本地磁盘、永久,  | <========>     |                                 |
|  约 100 MB)                   |  Callback      |Location Database (LDB)          |
|- callback promise 表,只有     |  (反向)      |- 卷 -> 服务器 的映射(位置透明)|
|  valid / canceled 两个状态     |                |                                 |
|- 与 Vice 之间只有 Fetch / Store|                |Volume(卷)                     |
+--------------------------------+                |- 文件树的子树,可整体迁移 / 复制|
                                                  +---------------------------------+

   Cell(单元)= 一个自治的管理域;全局根 /afs,路径形如 /afs/cell/user/...
  • 各组件的职责

    组件位置职责
    Venus客户端进程缓存管理器;把应用的 open/read/write/close 翻译成对 Vice 的 Fetch/Store RPC;维护 callback promise 状态
    Vice服务器文件服务器(文件本体 + 版本号)+ Location Database(LDB)(记录”哪个卷在哪台服务器上”)+ Volume 管理
    Cell管理域一个自治单元;全局根 /afs,路径形如 /afs/cell/user/...位置透明(服务器换了,路径不变)
    Volume(卷)逻辑单元文件树的一棵子树;可以整体迁移与复制(用于负载均衡与容灾),迁移后只需更新 LDB
  • 回调承诺(Callback promise)——AFS 的灵魂

    定义:当 Venus 打开一个文件、Vice 把整个文件发给它时,Vice 同时给出一份 callback promise“如果之后有别的客户端修改并关闭了这个文件,我会主动发一个 callback 通知你”。客户端那边的状态只有二值valid(承诺有效,我的副本是最新的)canceled(承诺已被打破,我的副本已失效)

    机制图解(回调流程)

   客户端 A(Venus)           服务器 Vice              客户端 B(Venus)
        |                        |                        |
        |  Fetch(F)              |                        |
        |----------------------->|  记录:F 的 callback 表 += A
        |<== 整个文件 + 版本号 v1 ==|  (对比 NFS:服务器这里什么都不记)
        |  promise = VALID       |                        |
        |  ...本地随便读写...      |                        |  Fetch(F) ---> 整文件 v1
        |                        |                        |  promise = VALID
        |                        |<------- Store(F, v2) ---|  B 关闭时整文件写回
        |                        |  版本号 v1 -> v2         |
        |<==== BreakCallback(F) ==|  向 callback 表里除 B 以外的所有人回调
        |  promise = CANCELED    |                        |
        |  (丢弃本地副本)        |                        |
        |  下次 open(F):承诺无效  |                        |
        |  ---> 重新 Fetch 拿到 v2 |                        |
  • 与 NFS 的对照(一句话总结)NFS 是”客户端主动拉取验证”(pull-based validation),AFS 是”服务器主动推送失效”(push-based invalidation)。这决定了 AFS 的一致性更强、失效更及时;也决定了 AFS 的服务器必须是有状态的

  • 讲义思考题:”回调的优缺点是什么?”

     内容
    优点 1一致性更强:失效是被推送的,不是等客户端下次打开才发现 ⇒ 客户端手里的数据”要么是最新的,要么已知失效”
    优点 2更省流量、更精准:服务器知道缓存了这个文件,只通知需要通知的人;而 NFS 的客户端只能靠”每次打开都问一遍”或”等新鲜期过期”
    优点 3打开后零网络流量:整文件在本地,重复的、顺序的读全部本地命中
    缺点 1服务器必须有状态:要为每一个 (客户端, 文件) 维护 callback 记录 ⇒ 内存开销与客户端数成正比;也违背了 NFS “无状态最健壮”的信条
    缺点 2崩溃恢复复杂:服务器崩溃会丢掉 callback 表,重启后无法知道谁手里有旧数据 ⇒ 必须保守处理(见 22.3.2 的正确性论证)
    缺点 3客户端离线/网络分区:服务器回调不到客户端时无法确认失效(AFS 用超时与”承诺最终作废”处理,Coda 更进一步支持断开操作 disconnected operation

22.2.10 AFS 的一致性模型

  • 定义与目的:AFS 的接口仍然是 open/read/write/close远程访问式的接口),但它的实现是”下载/上传”式的打开即下载整个文件,关闭即上传整个文件。一致性协议就挂在这两个动作上。

  • 读写都是”乐观的(optimistic)”:讲义原话——读写都在客户端的本地副本上进行文件关闭时,写才传播回 Vice。也就是说,读和写在本地是零网络开销的,代价是:在文件关闭之前,别人看不到你的修改

  • 打开时的一致性检查:Venus 打开文件时,先看自己的 callback promise

    情形动作
    本地有副本且 promise = VALID直接用本地副本(0 次 RPC!这是 AFS 可扩展性的关键)
    本地无副本,或 promise = CANCELED向 Vice Fetch 整个文件,重新获得 promise
  • 关闭时的传播与失效:若文件被修改过,Venus 在 close 时把整个文件写回 Vice;Vice 做三件事:
    1. 落盘递增文件版本号
    2. 向 callback 表里的所有其他客户端推送失效(callback break)
    3. 把该文件的 promise 重新授予写回者。
  • AFS 保证的语义(这是本章必须背下来的一句话)

    若两个客户端不并发写同一个文件,则它们都能看到对方最新”已关闭”的版本。

    也就是说,AFS 提供的不是 NFS 那种”要等到下次打开才检查”的弱保证,而是“关闭即广播失效”的强保证:写方一 close,所有持有旧副本的客户端立刻知道自己的副本作废了,下次访问必定重新取回最新版本。

  • AFS 不保证什么并发写的正确性。AFS 没有锁、没有原子性——如果两个客户端同时打开同一个文件并各自修改、然后相继关闭,后关闭者的整份副本会完整覆盖前者的整份副本(lost update,丢失更新)。这不是 bug 而是设计取舍:AFS 假设”多数文件是单一用户访问的”,因此把共享写的正确性交给应用层锁(应用自己用文件锁/租约来串行化写者;AFS 提供锁服务但默认不用,因为锁服务会引入服务器的额外状态与调用)。

  • 直观解释(”它是什么?”):AFS 的一致性像共享文档的”签出/签入”:你把整份文档借回家改(打开即签出),改完交回去(关闭即签入);图书馆一旦收到新版本,就会挨个打电话通知所有借了旧版本的人”你手里的作废了”。但如果两个人同时借走同一份并都交了回来,后交的那份会整个覆盖前一份——图书馆没有”逐段合并”的能力。

  • 关键假设与系统模型:AFS 假设 close 是应用表达”我的修改完成了”的信号(因此不 close 的写永远不保证被看到);假设 Vice 能可靠地把 callback 送达(否则要靠超时保守失效);假设 文件大小适合整传

22.2.11 AFS 的可扩展性分析(定量)

这一节回答:为什么 AFS 能扩展到数万个客户端,而 NFS 在那个规模下服务器会先垮?

  • 服务器负载模型(本章最重要的一条定量对比)
 NFSAFS
服务器要处理的请求每次块级访问 + 每次属性验证LOOKUP/GETATTR/READ/WRITE…)每次打开(可能 fetch)+ 每次关闭(可能 store),以及回调失效
负载与什么成正比$\propto$ 客户端数 × 单位时间访问次数与读的数据量、文件大小、会话长度都成正比$\propto$ 打开/关闭次数(与”打开后读多少遍”“读了多少块”无关
会话内重复访问每次缓存未命中就是一次 RPC;块级粒度 ⇒ 大文件顺序读 = 大量 RPC0 次 RPC(整文件已在本地)
形式化$L_{\text{NFS}} \approx N \cdot r_{\text{block}} \cdot B$($B$ = 每会话读的块数)$L_{\text{AFS}} \approx N \cdot r_{\text{open}} \cdot (1 + p_{\text{miss}})$($p_{\text{miss}}$ 主要来自回调失效)
结论客户端数 $N$ 增长时,服务器负载线性增长且系数很大只要”少量文件被打开一次、长期反复读“,服务器负载与客户端规模近似无关
  • 一个直接推论(也是 AFS 的设计出发点)整文件缓存把”块级访问次数”从服务器负载公式里彻底消掉了。讲义给出的三条假设(文件小、单用户、读多写少)恰好保证:每个文件每次被”取一次”就能服务大量的读。这就是 AFS 能扩展到数万客户端的根本原因。

  • 22.4.2 的实测(本笔记的模拟代码,random.seed(425) 固定)

负载客户端数 $N$NFS 服务器请求数AFS 服务器请求数NFS 陈旧读AFS 陈旧读
只读(每会话顺序读两遍)2 / 10 / 5084 / 420 / 2036(每客户端 ≈ 40.7)8 / 40 / 192(每客户端 ≈ 3.8)00
读多写少(5% 会话写一块)2 / 10 / 5094 / 585 / 3764(≈ 75.3/客户端)11 / 70 / 580(≈ 11.6/客户端)0 / 70 / 12730 / 0 / 0
  • 两个数字要读懂
    1. 只读场景:AFS 的请求数 ≈ $N \times$ 文件数(每个客户端把 4 个文件各取一次 = 4 次请求,之后所有会话都是 0 请求);NFS 的请求数 $\approx N \times$ 会话数 $\times$ (1 次 GETATTR + 冷块读)AFS 与”读了多少遍”无关,NFS 与它强相关
    2. 混合场景:AFS 的请求数也涨了(580),但涨的原因是回调失效导致的重新取回(536 次回调);同时 AFS 传输的字节数反而更多(4.75 MB vs 3.34 MB)——整文件传输的粒度粗,这是它的代价。而 NFS 出现了 1273 次陈旧读,AFS 是 0:这正是”用服务器状态换一致性”的收益。
  • 关键假设与系统模型:AFS 的可扩展性结论依赖”文件小 + 读多写少”。如果工作负载是”大文件 + 频繁随机小块访问”,AFS 会崩溃得很惨(每次打开都要传整个文件),而 NFS 反而更合适——假设变了,结论就反转

22.2.12 NFS vs AFS 全面对比

维度NFS(Sun,1980s,今天仍在用)AFS(CMU,OpenAFS 继承)
缓存粒度块级(默认 8 KB)+ 属性缓存整文件(whole-file),缓存在本地磁盘跨重启持久
服务器是否有状态无状态(不记打开文件表;锁/状态靠 NLM/NSM 外挂)有状态:维护 callback 表(谁缓存了哪个文件)
一致性机制客户端主动验证(pull):open 时 GETATTR 验证 + 属性新鲜期 $t$;close 时写回服务器主动失效(push):close 时版本号递增 + 回调所有持有者
一致性强度:只保证 close-to-open;并发打开时可能互相看不见(22.2.7 的反例)较强:只要不并发写同一文件,关闭后立刻全可见
原子性无(write 是覆盖,没有跨客户端原子性)无(整文件覆盖,甚至更容易丢失更新)
接口模型远程访问(细粒度 (fh, offset, len)接口是远程访问,实现是下载/上传
可扩展性服务器负载 $\propto$ 客户端数 × 块访问频率 ⇒ 服务器是瓶颈服务器负载 $\propto$ 打开/关闭次数 ⇒ 可扩到数万客户端
适合的工作负载通用(共享小文件、随机访问、Unix 全家桶);本地缓存 + 服务器缓存都有用读多写少、中小文件、顺序读、长会话;不适合大文件与随机小块
大文件友好(只传需要的块)不友好:打开一次传整份
崩溃恢复难度容易:服务器无状态,客户端重试即可较难:callback 状态丢失后必须保守失效(22.3.2)
命名透明性位置透明(挂载点后无服务器名),但 fh 绑定服务器 ⇒ 位置不独立位置透明 + 位置独立/afs/cell/...,卷可整体迁移/复制
代表部署极广:Unix/Linux 默认文件共享、企业 NAS、HPC、K8s 的 nfs volumeCMU 与多所大学的校园计算;OpenAFS 在部分大学/科研集群仍在使用

22.2.13 GFS 的设计背景与假设(所有设计决策的源头)

  • 定义与目的GFS(Google File System) 是 Google 在 2003 年公开的分布式文件系统(SOSP 2003 论文),是 HDFS 的直接蓝本,也是 MapReduce 的存储底座(详见 Lecture 4)。理解 GFS 的唯一正确姿势是:先接受它的六条假设,然后你会发现每一个”奇怪”的设计决策都是这些假设的必然推论。
#设计假设(GFS 论文口径)由此产生的设计决策
1组件故障是常态,不是例外:系统由数千台商品化(commodity)机器组成,磁盘/内存/网络/软件每天都有东西在坏系统必须自动检测、自动容错、自动恢复复制(默认 3 副本)+ 校验和 + 自动修复是默认配置,不是可选项;不能依赖任何单点长期存活
2文件巨大:常见的是 GB 到 TB 级;几百万个小文件不是优化目标(但必须支持)chunk 取 64 MB(大块减少元数据量);元数据放内存小文件性能要求被主动放弃
3工作负载:大文件的顺序读为主(多为流式读),追加写远多于随机写客户端不缓存数据(顺序流式读没有复用价值);Record Append 成为一等公民
4两种读模式:大顺序读 + 小随机读(后者不优化)不做块级缓存优化、不做低延迟优化
5高持续带宽比低延迟更重要:批量数据处理关心”每小时能过多少 TB”pipeline 推数据(吞吐优先);用大 chunk 摊薄开销;接受”单次操作延迟高一点”
6应用可以配合:客户端与文件系统可以协同设计,应用愿意接受放松的一致性模型松弛一致性(relaxed consistency) + 应用层去重/校验/检查点(22.2.18-22.2.19)
  • 直观解释(”它是什么?”):GFS 不是”给程序员用的通用文件系统”,而是“给批量数据处理程序用的传送带”。通用文件系统追求”随时能改、谁都能用、语义精确”;GFS 追求”一次写进去、很多人顺序读出来、坏几台机器也不停机“。把这点想清楚,就不会再问”GFS 为什么不支持并发修改同一区域”——因为它假设的应用根本不需要这个功能。

  • 关键假设与系统模型:崩溃-恢复(crash-recovery)故障模型;单个 master(元数据单点,靠日志+检查点+影子 master 提升可用性);网络不可靠(数据靠 checksum 而不是靠网络保证正确);没有全局时钟(租约用本地时钟近似,要求时钟漂移有界);假设应用会配合

22.2.14 GFS 架构:单 master + 多 chunkserver + 客户端(本章最重要的图)

  • 定义与目的:GFS 集群由三类角色组成:一个 master(管元数据)、多个 chunkserver(存数据)、多个 client(读写)。整个设计最关键的一句话是:元数据走 master,数据不走 master。

  • 机制图解:GFS 架构,以及”元数据路径”与”数据路径”的分离。

+----------------------------------------------------------------------------+
| MASTER(单节点,元数据全部在内存)                                         |
+----------------------------------------------------------------------------+
| 命名空间(目录树)   |  文件 -> chunk 映射   |  chunk -> 副本位置            |
| 租约(lease)管理     |  垃圾回收(GC)         |  chunk 迁移/再平衡           |
| 操作日志(operation log) + 检查点(checkpoint)(持久化)                     |
+----------------------------------------------------------------------------+
        ^                                  ^
   (1) 元数据请求(控制流)        (1) 元数据请求
     "这个 chunk 在哪几个 cs 上?"   "谁是 primary?"
        |                                  |
+----------------------------------------------------------------------------+
| CLIENT                                                                     |
| 只缓存元数据(chunk 位置 / 谁是 primary),不缓存文件数据                  |
+----------------------------------------------------------------------------+
        |                                  |
   (2) 数据流:client <=================> chunkserver(不经过 master)
            v                          v                          v
+-----------------------+   +-----------------------+   +-----------------------+
| ChunkServer cs1       |   | ChunkServer cs2       |   | ChunkServer cs3       |
| chunk 数据(64 MB/个)|   | chunk 数据(64 MB/个)|   | chunk 数据(64 MB/个)|
| 每 chunk 3 副本之一   |   | 每 chunk 3 副本之一   |   | 每 chunk 3 副本之一   |
| 64 KB 一块 + 32 位 CRC|   | 64 KB 一块 + 32 位 CRC|   | 64 KB 一块 + 32 位 CRC|
+-----------------------+   +-----------------------+   +-----------------------+

元数据路径:client <-> master(只传位置信息)    数据路径:client <-> chunkserver(传字节)
  • 单 master 存什么(全部是元数据,没有文件数据)
    1. 命名空间(namespace):目录树、文件名、属性;
    2. 文件 → chunk 的映射(哪个文件的第几个 64 MB 是哪个 chunk handle);
    3. chunk → 副本位置(哪个 chunk 在哪几个 chunkserver 上);
    4. 租约(lease):哪个 chunk 当前把 primary 租给了谁、何时到期;
    5. 垃圾回收(GC):延迟删除(先标记、后回收),以及 chunk 迁移/再平衡(rebalancing)

    注意第 3 项的特殊性master 并不持久化保存”副本位置”——它在启动时通过 chunkserver 的心跳(heartbeat)汇报来重建位置信息。这是有意的:位置信息随机器上下线高频变化,持久化只会带来不一致;而”哪些 chunk 属于哪个文件”才是必须持久化的权威事实

  • 为什么 master 能很快它不存文件数据,所有元数据都在内存里;论文给出的量级是每个 64 MB chunk 只需不到 64 字节的元数据(22.5 会把这个数字算给你看:100 TB 数据 ≈ 156 万个 chunk ≈ 100 MB 元数据;即便按”每 chunk 数百字节到 1 KB”的保守工程估计,也只有约 1.6 GB,依然放得下)。

  • 为什么 chunk 是 64 MB(必须记住的三个理由 + 一个代价)

    理由推导
    减少 master 的元数据量chunk 越大,同样数据量需要的 chunk 数越少 ⇒ 内存中的映射表越小 ⇒ 单 master 能管的数据总量越大
    减少客户端与 master 的交互一次租约覆盖更多数据:客户端问一次 master 就能对 64 MB 连续读写(若 chunk 是 1 MB,同样的写要问 64 次)
    摊薄 TCP 连接与握手开销与 chunkserver 建立连接、发送请求的固定成本被分摊到更大的数据量上;长连接可以保持,避免反复建连
    代价:小文件与热点 chunk小文件占用一整个 chunk(内部碎片);同一个 chunk 被大量客户端同时访问时,持有该 chunk 的 chunkserver 会成为热点(虽然 3 副本能分摊读,但无法分摊”所有客户端都抢同一个 64 MB”的热点写)。缓解手段:提高副本数、允许客户端从别的客户端读、让应用错开启动时刻批量写
  • 客户端的两条铁律
    1. 客户端只缓存元数据(chunk 位置、primary 是谁),不缓存文件数据——因为工作负载是流式顺序读,读过的数据不会再读,缓存毫无收益(这一点和 NFS/AFS 完全相反!);
    2. 客户端直接和 chunkserver 传数据,数据不经过 master——否则 master 会成为整个集群的带宽瓶颈,架构也就无法扩展到数百个 chunkserver。
  • 关键假设与系统模型:GFS 假设 chunk 数量(元数据量)能放进单机内存;假设元数据操作频率远低于数据操作频率(否则单 master 会成为瓶颈);假设每个 chunk 有 3 个副本,且分布在不同的机架(讲义 Lecture 4 的 2+1 布局:同一机架 2 份、另一机架 1 份,既容机架故障又限制跨机架写带宽)。

22.2.15 GFS 的读流程(客户端 ↔ master ↔ chunkserver)

  • 定义与目的:读流程是 GFS 里最简单也最能体现”元数据/数据路径分离”的流程。

  • 逐步流程(五步,必须记住顺序)

    1. 客户端把文件偏移换算成 chunk 索引:GFS 的文件被切成固定大小 64 MB 的 chunk(最后一个 chunk 可以不满),因此 chunk index = offset / 64 MBoffset in chunk = offset % 64 MB这个换算完全在客户端本地完成,不需要任何网络交互
    2. 客户端向 master 请求该 chunk 的副本位置:发 (文件名, chunk index);master 回复 chunk handle(一个 64 位全局唯一 ID)+ 每个副本所在的 chunkserver 位置
    3. 客户端缓存这份元数据(带超时):下次访问同一个 chunk 不再问 master。这一步就是”元数据可缓存”——单 master 不至于成为瓶颈的第二道保险。
    4. 客户端直接向其中”最近”的一个 chunkserver 发读请求(chunk handle, byte range)。最近的判定标准是网络拓扑距离(同一台机器 / 同一机架 / 同一边界内的机房),客户端会优先选最近的副本——这正是 Lecture 4 讲的数据本地性(data locality)思想。
    5. chunkserver 返回数据;整个过程master 只参与了第 2 步
   CLIENT                    MASTER                 ChunkServer cs1     cs2     cs3
     |                         |                          |              |       |
     |--(1) 本地换算 offset -> chunk index                  |              |       |
     |--(2) get chunk H 位置-->|                          |              |       |
     |<-- H 在 {cs1,cs2,cs3} --|  (master 只做元数据)      |              |       |
     |--(3) 缓存 (H -> 位置)                                   |              |       |
     |                                                       |              |       |
     |--(4) 读 (H, [0,1MB)) ------------------------------>|  (最近副本)   |       |
     |<--------------- 数据 ---------------------------------|              |       |
     |                                                                      |       |
     |  (若 cs1 超时/不可用,客户端改向 cs2 读 —— 副本的容错价值在这里体现)  |       |
  • 正确性直观论证:因为 chunk 的内容由主副本按序列号串行化(22.2.16)后一致地写到所有副本,所以从任意一个副本读都能得到相同的数据(”consistency”的定义);副本之间唯一的差别是位置与负载,不是内容。除非(a)该 chunk 上一次写只部分成功(22.2.18 的”不一致”情形),或(b)某个副本发生了静默位腐败——这时靠 checksum 检测并改读其他副本(22.2.19)。

  • 复杂度:一次读的网络往返最少 1 次(若元数据已缓存)或 2 次(先问 master 再读数据);客户端缓存元数据后,连续顺序读的额外开销趋于 0

  • 关键假设与系统模型:读操作幂等(同一 (handle, offset, len) 重复读结果不变),因此可以安全重试;读不需要租约(只有写才需要 primary);客户端假设master 回复的位置可能过期(chunkserver 可能已下线)⇒ 必须能”换一个副本重试”。

22.2.16 GFS 的写流程与租约(本章最精巧的部分)

  • 定义与目的:GFS 的写要同时满足两件互相拉扯的事:(a)高吞吐(不能让 master 参与每一次写),(b)副本内容一致(不能让三个副本的写顺序不同)。它的解法是:master 只授予”主副本资格”(租约),写顺序由主副本自己串行化。

  • 为什么必须有一个 primary(主副本):如果三个副本各自接受客户端的写请求并各自决定顺序,那么当两个客户端并发写同一 chunk 时,三个副本可能以不同顺序应用两次写 ⇒ 副本内容不同(inconsistent)。GFS 的解法是:对每个 chunk,在任一时刻最多只有一个 chunkserver 是 primary;所有写请求都先到 primary,由 primary 分配序列号(serial number)并决定顺序,secondary 严格按 primary 给的顺序执行。于是”多个客户端并发写”被化简成”primary 上的一个串行序列”。

  • 为什么用租约(lease)而不是问 master:如果每次写都要先问 master”谁是 primary”,那么 master 就得处理每一次写的元数据请求,单 master 立刻成为瓶颈。租约是”限期授权”:master 授予某个 chunkserver 该 chunk 的 primary 资格,默认 60 秒(可续约);在租约有效期内,客户端可以直接使用缓存的 primary 位置,master 完全不参与。租约把 master 的交互次数从”每次写一次”降到”每 60 秒每 chunk 一次”(22.3.4 会给出定量证明)。

  • 完整写流程(七个步骤 + 时序图)——注意数据流与控制流是分开的两条往返,这是 GFS 设计中最容易考也最容易答错的地方:

   CLIENT                 MASTER              PRIMARY(cs1)      SECONDARY(cs2)   SECONDARY(cs3)
     |                      |                      |                 |                |
     | (1) 谁是这个 chunk 的 primary?              |                 |                |
     |--------------------->|                      |                 |                |
     |                      | 无租约/租约已到期 => 选 cs1 为 primary,授予 60s 租约     |
     |                      |--------------------->| lease(H, +60s)  |                |
     |<-- primary=cs1, secondaries=[cs2,cs3], expire=|                 |                |
     |                      |                      |                 |                |
     | (2) 数据流:先把数据推给所有副本(此时还【不发】写指令)        |                |
     |===== data chunk ============================>|                 |                |
     |                      |                      |=== data =======>|                |
     |                      |                      |                 |=== data ======>|
     |<==================== ack ====================|<==== ack =======|<==== ack ======|
     |    (链式 pipeline:数据沿 cs1 -> cs2 -> cs3 逐跳转发,ack 反向回传)              |
     |    (cs1/cs2/cs3 都把数据放进 LRU 缓冲,此时【还不知道】它要写到哪个偏移)        |
     |                      |                      |                 |                |
     | (3) 控制流:把写请求发给 primary(只带 data_id + 偏移 + 长度)  |                |
     |-------------------------------------------->|                 |                |
     |                      |   (4) primary 分配序列号 s=7,先写到本地   |                |
     |                      |   (5) ---- apply(s=7, offset) ---->     |                |
     |                      |                      |---------------->|                |
     |                      |                      |------------------------------- ->|
     |                      |   (6) <--- ack ----  |<---- ack -------|<---- ack ------|
     |<-- (7) ok(serial=7) -|                      |                 |                |
  • 每一步的动机(为什么这样设计)

    步骤设计动机
    (1) 问 primary只有这一步需要 master;客户端会缓存 primary 位置直到租约到期
    (2) 数据流:pipeline 推送到所有副本关键设计:数据流与控制流分离。 数据沿 chunkserver 链式(pipelines) 传输:client -> 最近的副本 -> 下一个副本 -> …动机:若客户端自己给 3 个副本各发一份完整数据,则客户端的上行带宽要承载 3 份数据,客户端成为瓶颈;而链式转发时客户端上行的数据量只有 1 份,其余由服务器之间并行转发。22.4.3 的实测给出精确结论:只要客户端上行 $U$ 小于副本数 $k$ 与服务器链路 $B$ 的乘积($U < kB$),pipeline 就更快,收益最高可达 $k$ 倍
    (2) 数据”先到、偏移未定”副本先把数据放进 LRU 缓冲区,此时还不知道写到哪里 ⇒ 数据到达可以流水线并行,不必等 primary 定序;这也是”控制流与数据流分离”的直接体现
    (3) 写指令只发给 primary一个实体定序,避免多主冲突
    (4) primary 分配连续的序列号这就是”串行化”(serialization):primary 为落在同一 chunk 上的所有并发写分配单调递增的序列号,并按序列号顺序在本地应用
    (5) primary 按序列号顺序转发给 secondarysecondary 按收到顺序执行;由于 primary→secondary 之间是可靠 FIFO 通道(TCP),且 primary 是唯一发送者,所以每个 secondary 上的应用顺序 = primary 的序列号顺序
    (6) 收集所有 secondary 的 ack只要有任何一个 secondary 失败,primary 就向客户端报告失败
    (7) 返回客户端客户端可以重试整个写(回到 (1) 重新确认 primary)
  • 失败处理与”重复数据”的根源:若某个 secondary 失败,primary 向客户端报错,客户端重试整个写(可能重试多次)。这带来两个后果:
    1. 重试的写会拿到一个”新的、更靠后”的序列号,因此可能在同一 chunk 上产生重复数据(GFS 论文明确承认这一点);
    2. 因此 GFS 的写语义不是”精确一次”,一致性模型只能是”松弛的”(22.2.18)。
  • “客户端如何知道 chunk 里有什么”普通写由客户端自己指定偏移,所以写入方自己知道写在哪;但并发写时它无法知道别人的写是否与它交错(22.2.18 的 “undefined”)。这正是 Record Append 要解决的问题。

  • 关键假设与系统模型:租约的唯一性依赖 master 是唯一授权者且时钟漂移有界(否则旧 primary 可能在”自己以为租约还有效”时继续服务,出现脑裂式的双主);副本位置的新鲜度依赖心跳;数据完整性依赖 checksum(因为数据流是”先推后写”,缓冲区中的数据必须能被校验)。

22.2.17 Record Append:GFS 的杀手级特性

  • 定义与目的Record Append(记录追加) 是 GFS 为”多生产者、单消费者”场景(日志收集、MapReduce 输出、多客户端并发写同一个结果文件)专门设计的操作:客户端只给数据,偏移由 GFS 决定

  • 语义(必须逐字理解)

    客户端提供数据,GFS 原子地把数据追加到文件末尾至少一次(at-least-once),并把数据被写入的实际偏移量返回给客户端。

    三件事同时成立:(a)GFS 选择偏移;(b)追加是原子的(不会与别人的追加交错);(c)至少一次(可能重复、可能出现填充)。

  • 与普通写的差别(这是它的全部价值)
 普通写(Write)Record Append
谁决定偏移客户端write(offset, data)GFS(primary)
并发写同一区域允许发生 ⇒ 各客户端的写互相覆盖/交错不可能发生(偏移由 primary 串行分配,两条记录不会重叠)
偏移是否可知客户端本来就知道客户端从返回值得知(因此是 “defined” 的)
典型用法覆盖写固定位置多个客户端/多台机器并发往同一个日志文件追加
失败后果副本可能出现不同内容可能出现重复记录填充,但不会有交错
  • 实现要点(四步)
    1. 客户端把数据推送到 primary 与所有 secondary(同普通写的 pipeline);
    2. 客户端把追加请求发给 primary;
    3. primary 决定追加位置(= 当前 chunk 的逻辑末尾),串行化所有副本的追加位置,然后按序列号顺序让副本在同一个偏移写入;
    4. primary 把实际偏移返回给客户端。
  • 跨 chunk 边界的 padding:如果当前 chunk 剩余空间装不下这条记录,primary 会把当前 chunk 填充(pad)到 64 MB(写零),并告诉客户端”请到下一个 chunk 重试“。
   chunk 0(64 MB,逻辑末尾 = 64 MB - 200 B)        chunk 1(新分配)
   +--------------------------------------------+   +--------------------------+
   | ...... 已有记录 ......        | 剩余 200 B  |   |                          |
   +--------------------------------------------+   +--------------------------+
                            |                                    ^
       客户端要追加 300 B 的记录(200 B 装不下)                    |
                            v                                    |
   +--------------------------------------------+   +--------------------------+
   | ...... 已有记录 ......        | 000...(padding) |  | 该记录最终写在这里(offset = 64 MB)|
   +--------------------------------------------+   +--------------------------+
        padding 浪费最多 64 MB - 1 的空间,但换来了“一条记录不会被切成两半”的原子性
  • at-least-once + 可能重复 + 可能填充,这套语义为什么”够用”
    • 应用接受重复:日志收集、MapReduce 输出这类应用天然可以容忍重复记录——只要每条记录带唯一 ID,读者在消费时去重即可(这正是 MapReduce 的做法:Reduce 输出写临时文件,见 22.2.19)。
    • 应用能识别填充:填充是已知的零字节区域;读者可以用记录长度/校验和/魔数把填充与真实记录区分开(GFS 论文建议记录里带 checksum,读到损坏或填充就跳过)。
    • 换来的是吞吐多生产者可以无协调地并发追加到同一个文件——这在普通写模型下需要应用自己做分布式锁,代价远高于”偶尔重复一条记录”。
  • 关键假设与系统模型:应用愿意在应用层处理重复与填充(这是”松弛一致性”的最低门槛);primary 在租约期内必须唯一(否则两条记录可能被分配同一个偏移);至少一次意味着不保证每个副本都有(失败的副本可能缺这条记录,它的数据区域会留下空洞,直到被修复)。

22.2.18 GFS 的一致性模型(最容易考、也最容易被误解的部分)

  • 定义与目的:GFS 明确放弃了 POSIX 式的强一致性,采用的是一套松弛的一致性模型(relaxed consistency model),理由是:在 GFS 的目标工作负载下,要求强一致的性能与复杂度代价远大于收益。但”松弛”不等于”随便”——GFS 精确定义了三种状态,应用开发者必须按这套定义写程序。

  • 三个术语的严格定义(论文口径,必须一字不差地理解)

    术语定义直观含义
    一致的(consistent)无论从哪个副本读,所有客户端看到相同的数据“三个副本内容一样”
    确定的(defined)在一致的基础上,客户端还能看到”完整的”写内容(因为写在副本上是被串行化执行的)“我不但知道大家看到的一样,还知道这段数据就是我写的那段”
    不一致的(inconsistent)不同副本返回不同的数据“副本之间已经对不上了”

    注意”一致”与”确定”是两个独立的维度:”一致”说的是副本之间的关系;”确定”说的是客户端能否预测内容。可能有”一致但未定义”的组合,但不可能有”确定但不一致”(确定以一致为前提)。

  • 完整状态矩阵(必考)

写类型串行成功(serial success)并发成功(concurrent success)失败(failure)
普通写(Write)确定(defined):所有副本内容 = 这一次写的完整内容(因为只有一个写者,序列号顺序就是它)一致但未定义(consistent but undefined)所有副本相同,但内容是若干次并发写的任意交错片段,客户端无法预测看到什么不一致(inconsistent):不同副本内容不同(某些副本写了、某些没写)
记录追加(Record Append)确定(defined):所有副本在同一偏移同一条完整记录(定义良好的”至少一次”语义)确定的部分 + 交错区域:在定义好的偏移处,所有副本数据相同;但某些副本可能有重复记录或填充(因此不同副本可能在尾部长度上不同)不一致(inconsistent)
  • ASCII 版本的状态矩阵(考试时画这个):
+------------------------+----------------------------------+----------------------------------+
|                        | Write(客户端定偏移)            | Record Append(GFS 定偏移)      |
+------------------------+----------------------------------+----------------------------------+
| 串行成功 (serial)      | DEFINED:三副本内容相同,且就是  | DEFINED:同一偏移处有同一条完整  |
|                        | 这一次写的完整内容               | 记录(定义良好的至少一次语义)   |
+------------------------+----------------------------------+----------------------------------+
| 并发成功 (concurrent)  | CONSISTENT but UNDEFINED:三副本 | DEFINED 的部分 + 交错区域:定义  |
|                        | 相同,但内容是若干次并发写的     | 偏移处数据相同,但某些副本可能   |
|                        | 任意交错片段,客户端无法预测     | 有重复记录或 padding(尾部不同) |
+------------------------+----------------------------------+----------------------------------+
| 失败 (failure)         | INCONSISTENT:不同副本内容不同   | INCONSISTENT:不同副本内容不同   |
|                        | (有的写了、有的没写 ⇒ 空洞)    | (有的写了、有的没写 ⇒ 空洞)    |
+------------------------+----------------------------------+----------------------------------+
一致(consistent) = 副本之间相同;确定(defined) = 客户端能预测内容;确定 ⇒ 一致
  • “一致但未定义”到底是什么意思(必须用一个具体例子讲透): 设客户端 A 在偏移 0 写 600 字节的 'A',客户端 B 在偏移 300 写 600 字节的 'B',两者并发
    • primary 把它们串行化,假设顺序是 B 然后 A。于是三个副本最终都是:[0,300)='A'(A 写的后半段覆盖了 B 的前 300 字节?——不对,让我们仔细算):
      • 先应用 B:[300,900) = 'B'
      • 再应用 A:[0,600) = 'A',覆盖 [300,600)
      • 结果:[0,600)='A'[600,900)='B' ⇒ 布局 AAABBBBBB
    • 若顺序相反(A 然后 B):[0,300)='A'[300,900)='B' ⇒ 布局 AAAAAABBB
    • 两种情况下三个副本都完全相同(consistent),但结果取决于客户端无法观测的序列顺序(undefined):A 完全无法知道”我写的 600 字节是不是还完整”(在第二种情形下它被 B 覆盖了一半)。 22.4.1 的实验 1 就是这样演示的:同样的两次写,只因网络时序变化,布局就在 AAABBBBBBAAAAAABBB 之间摆动,而三个副本始终一致
  • GFS 如何保证”基本可用”(松弛一致性不等于不保护数据)
    1. chunk 完整性靠 checksum(22.2.19):把 chunk 分成 64 KB 的块每块一个 32 位校验和;读到不匹配 ⇒ 判定该副本损坏 ⇒ 改读其他副本
    2. 通过副本间比较发现”静默不一致”:chunkserver 周期性互相比较 checksum,发现不一致就从健康副本修复(论文称这种”影子副本/副本间校验”是应对位腐败(bit rot / silent data corruption) 的关键手段);
    3. 写路径的”至少一次”:宁可重复,不可丢失。
  • 应用如何应对松弛一致性(四条工程惯例,必须记住)
    1. 优先使用 append 而不是覆盖写(append 更高效、更容错;GFS 论文建议应用”追加、偶尔检查点”,而不是”随机写”);
    2. 用检查点(checkpoint):应用定期写检查点(一个完整、自洽的状态快照),使写入者可以从检查点恢复,而不是依赖底层文件系统的一致性
    3. 写”自验证、自标识”的记录:每条记录带 唯一 ID + 校验和,读时验证并去重(这就是”用应用层逻辑把 at-least-once 变成 effectively exactly-once”);
    4. 用”临时文件 + 原子 rename”实现原子提交(MapReduce 的做法):把结果写到临时文件,全部写完后在 master 上做一次原子的 rename ⇒ 读者要么看到完整结果,要么什么都看不到,永远看不到半成品(这也解释了为什么 MapReduce 的 reduce 输出是”先写 DFS 临时文件再改名”)。注意原子性来自 master 的元数据操作(单 master 串行处理命名空间操作),而不是来自 chunk 写入本身。
  • 关键假设与系统模型:GFS 假设应用愿意适配弱语义(这是最”社会性”的一条假设);假设批处理任务可以重跑(MapReduce 的重试模型);假设应用能容忍重复记录

22.2.19 完整性、校验、修复与 master 的高可用

  • 定义与目的:GFS 假设硬件会坏、位会翻转,因此必须在软件层检测损坏并自动修复,同时保证元数据服务不因单点故障而长期不可用

  • checksum(校验和)机制
    • 粒度:每个 chunk(64 MB)划分为 64 KB 的块每块保存一个 32 位校验和(因此一个 chunk 有 1024 个校验和,约 4 KB 的校验数据);
    • 写入时:每次写都重算受影响块的校验和;
    • 读取时先校验再返回;不匹配 ⇒ 拒绝服务并报错
    • 为什么放在 chunkserver 而不是客户端:客户端与 chunkserver 之间、以及副本之间的一致性检查都需要一个独立于数据的权威副本
    • 收益:把”读到错误数据(silent data corruption)”降级为”读到明确的错误,然后改读其他副本“——这是”检测(detection)”而非”纠正(correction)”。
  • 机制图解:位腐败的检测与修复
   正常状态: cs1 [chunk H: 数据 D | 校验和 S]      cs2 [H: D | S]      cs3 [H: D | S]
                                                       |
   发生静默位腐败(磁盘位翻转 / 内存位翻转 / 硬件 bug,不改校验和):
                                                     v
                                             cs2 [H: D' | S]   <-- D' != D,但 S 没变
   客户端读 cs2:
       1. cs2 读到数据后重算校验和 S' = crc(D') != S
       2. cs2 返回“损坏”错误(而不是把 D' 当数据返回)        <== 检测
       3. 客户端把损坏事件报告给 master
       4. 客户端改向 cs1/cs3 读取(幂等读 + 多副本 = 可用性)   <== 容错
   master:
       5. 把 cs2 上这个副本标记为无效,指令 cs2 从健康副本 cs1 重新复制(re-replicate)
       6. cs2 从 cs1 拉取整个 chunk,重算校验和            <== 修复
   结果: 3 副本恢复一致;期间客户端读服务没有中断
  • master 的高可用
    • 操作日志(operation log)所有元数据变更都先追加到日志(并落盘、复制到远端),它是元数据的权威来源,也是”逻辑时间线“——它定义了并发操作的顺序(例如”文件 F 的创建在 G 的 rename 之前”);只有日志落盘后,变更才对客户端可见。
    • 检查点(checkpoint):定期把内存中的元数据结构序列化成检查点;master 崩溃重启后,先加载最近的检查点,再重放(replay)其后的日志即可恢复到一致状态;日志与检查点都会复制到多台机器
    • 影子 master(shadow master):一个只读的 master 副本(消费同一份日志,落后主 master 一点点),提供只读服务。当主 master 故障时,客户端仍能读到(可能略旧的)元数据,而不是完全不可用。
    • 为什么单 master 不是数据瓶颈客户端不通过 master 传数据(数据走 chunkserver),元数据可以被客户端缓存(chunk 位置、primary 位置)⇒ master 只处理元数据请求。论文的实测口径是:单 master 的元数据吞吐在每秒数百次量级,而客户端与 chunkserver 的实际负载远低于此;真正的瓶颈是master 的元数据规模(内存)崩溃恢复的时间(日志大小)。(补充说明:具体 QPS 数字随论文版本与集群规模而异,此处给的是量级口径。)
    • 单 master 的风险与工程改进
      1. 命名空间分片:论文提出可以用多个 master,每个负责一部分命名空间(namespace sharding);生产上 Google 的 Colossus 就走向了多 master 的架构;
      2. HDFS 的路线NameNode HA(Active/Standby 两个 NameNode + ZKFC 健康检查 + JournalNode 共享编辑日志)解决单点;Federation(联邦) 用多个命名空间分担元数据规模。
  • 关键假设与系统模型:checksum 只保证检测,修复依赖至少一个健康副本(3 副本容忍 1 个损坏,但若同一 chunk 的两个副本同时损坏就需要更快的检测与更多的副本);master 的可用性依赖日志/检查点的复制故障切换机制

22.2.20 HDFS 与 GFS 的关系与差异

  • 定义与目的HDFS(Hadoop Distributed File System) 是 GFS 的开源实现(Yahoo!/Apache,2006 起),架构上几乎是 GFS 的翻版(NameNode 对应 master、DataNode 对应 chunkserver、Block 对应 chunk),但在工程细节与演进路线上有几处重要差异
维度GFS(Google,2003)HDFS(Apache,2006→至今)
块大小64 MB(论文默认)早期 64 MB,Hadoop 2.x 起默认 128 MB(讲义 Lecture 4 提到的”GFS/HDFS 的 chunk 是 64 MB”对应 HDFS 早期的默认值)
写模型支持对已打开文件的并发 Record Append(多生产者追加同一个文件)单写者模型(single writer):一个文件同时只能有一个写者;文件一旦被 close 就不能再写(除非以 append 模式重新打开,且此时也仍是单写者)⇒ 多客户端”并发追加同一个文件”在 HDFS 上不被支持,需要应用自己拆文件或走别的路径
元数据高可用单 master + 日志/检查点 + 影子 master(只读早期 NameNode 是单点(SPOF) ⇒ 后期引入 NameNode HA(Active/Standby + ZKFC 故障检测 + JournalNode 共享 EditLog)与 Federation(多个命名空间,分担元数据规模)
副本位置管理master 通过心跳重建位置(不持久化位置)类似:DataNode 周期性发 block report 给 NameNode,NameNode 据此维护块位置映射
写管道与租约master 授予 60 s 租约 → primary 串行化同样用租约(HDFS 的 lease 由 NameNode 授予:软上限默认 60 秒、写者需要周期续租,硬上限默认 60 分钟,超时则强制回收,用于防止文件被多个写者打开)与 写管道(pipeline),管道成员失败时做 pipeline recovery(替换失败的 DataNode 并重建管道)
数据可见性控制应用靠 Record Append 语义与论文建议明确提供 hflush()(让数据对读可见,但不保证落盘)hsync()(落盘 + 对读可见),把”何时可见”变成显式 API
副本数默认 3,机架感知 2+1默认 3,机架感知策略相同(讲义:2 份在同一机架、1 份在另一机架,既容机架故障又限制跨机架写带宽)
生态角色MapReduce 的输入/输出存储;GFS 之上还有 BigTableMapReduce / Spark / HBase / Hive / Flink 的存储底座;Hadoop 3.x 进一步引入纠删码(erasure coding)替代 3 副本、多 NameNode、GPU 支持等
  • 直观解释(”它是什么?”):HDFS 就是”开源版的 GFS + 十年工程补丁“:把 GFS 论文里”我们假设 master 不会长期挂掉”的乐观,补成了 HA 与 Federation;把 GFS 论文里”应用自己处理重复记录”的自由,收紧成”单写者 + 显式 hflush/hsync”的规则。理解 GFS 就理解了 HDFS 的 80%;剩下的 20% 是单写者约束、可见性 API 与高可用工程。

22.2.21 其他分布式文件系统(简述,用于建立坐标系)

系统一句话定位关键机制
CodaCMU 的 AFS 后继,面向移动/不可靠网络断开操作(disconnected operation):客户端可以在与服务器失联时继续读甚至写(记录日志,重连后回放);hoarding(预取)与回放;一致性用 AVSG(可用卷存储组) 判定
Sprite80 年代末 Berkeley 的日志结构网络文件系统全局共享的 prefix 缓存(按路径前缀路由)、客户端/服务器缓存统一管理
xFSBerkeley 的无服务器(serverless) 文件系统把文件系统的功能分布到所有节点(条带化 + 协作缓存),没有专用服务器
Lustre / GPFS / PanFSHPC 并行文件系统元数据服务器(MDS)与对象存储服务器(OSS)分离、条带化、客户端并行访问;面向”大文件 + 高聚合带宽”
Ceph统一存储(对象/块/文件),CRUSH 算法无中心元数据表:用 CRUSH 把对象确定性映射到 OSD(object -> PG -> OSD),元数据(MDS)只服务目录树,可动态扩缩
S3 / GCS 等对象存储不是文件系统,是键值/对象语义(HTTP PUT/GET,扁平桶)无目录(只有 key 前缀)、不支持随机写(对象整体替换)、历史上是最终一致(AWS 2020 起对新建/覆盖 PUT 提供强一致读后写)——接口约束换来了无限扩展性
Alluxio / JuiceFS“内存层/元数据与数据分离”的新型文件系统Alluxio:内存为中心的虚拟分布式存储,为 Spark/Presto 做缓存层;JuiceFS:元数据入 Redis/数据库,数据入对象存储,把 S3 包装成 POSIX 文件系统
  • 一句话总结这一节的用意“分布式文件系统”不是一种东西,而是一个被接口(对象/文件/块)、一致性(强/弱)与工作负载(顺序追加/随机读写)划出的设计空间。 本章的三位主角占据其中三个极端位置。

22.2.22 图书馆类比与”下载/上传 vs 远程访问”的根本分野

  • 直观解释(”它是什么?”)——用图书馆一次性记住三者的差别
系统类比带宽/延迟特征一致性特征
NFS每次去图书馆抄一页:你只借你需要的页(块),但每次都要跑一趟(RPC);如果别人改了书,你要重新来一趟才发现(open 时验证)请求数多、每次数据少、延迟敏感弱(close-to-open)
AFS把整本书借回家:一次借走整本(整文件),在家随便翻(本地读,0 网络流量);有人要改这本书,图书馆打电话让你作废手里的版本(callback);还书时整本交回(close 时写回)请求数少、每次数据多、大文件吃亏较强(回调失效,关闭后立即可见)
GFS把书切成分册存三个库房,只有一个库房负责登记改动:读者去最近的库房取(数据不经 master),改动要先问登记处”谁是主库房”(租约),主库房发号(序列号)后各库房按号执行高聚合吞吐、延迟不敏感、追加友好松弛(一致但可能未定义)
  • “下载/上传模型 vs 远程访问模型”——本章的根本分野
 下载/上传模型(upload/download)远程访问模型(remote access)
接口形态整个文件搬来搬去(Fetch/Store细粒度读写(read(fh, off, len)
优点本地访问零延迟;重复读无网络流量;服务器负载与访问次数解耦只传需要的数据;大文件友好;随机访问友好
缺点大文件代价高;首次打开延迟高;写回粒度粗(整文件覆盖 ⇒ 丢失更新风险大)每次访问都要 RPC(延迟敏感);服务器负载与访问频率成正比
一致性难点需要”版本/回调”来知道自己手里的副本是否过期需要”验证/租约”来决定缓存是否可用
代表AFS(True 下载/上传);S3(对象整体 PUT/GET)NFS(块级);GFS/HDFS(handle, offset, len)
  • GFS/HDFS 的”第三极”:用接口约束换一致性简化。GFS/HDFS 选择 “一次写、多次读(write-once-read-many)”不可变/追加模型,这带来一连串简化:
    1. 消除了大部分一致性问题:如果文件写完就不再改,那么”多个副本的内容一致”只需要在写入那一次保证(primary 序列号 + 3 副本同步写),之后所有读天然一致(这就是”serial success ⇒ defined”为什么是 GFS 的常态);
    2. 简化了缓存:客户端不缓存数据也不会丢失什么(顺序流式读没有复用);元数据缓存可以用超时简单处理;
    3. 简化了元数据:只追加 ⇒ 块的位置信息简单(追加在末尾),无需复杂的块索引更新;
    4. 换取到高吞吐:大块(128 MB)+ 顺序追加 + 数据与元数据路径分离 ⇒ 单集群数百 PB、数千节点。 这是一次极其经典的取舍:把”接口的表达能力”(不允许随机改)压缩,换来”一致性与实现复杂度”的大幅下降。 代价是:不适合随机写、不适合小文件、不适合多写者——所以 HDFS 至今不适合做数据库的存储引擎(HBase 用它在文件层之上重新组织成 LSM 结构)。
  • 黄金法则(本章总纲,务必背下来)

分布式文件系统的设计由工作负载假设决定:NFS 假设通用工作负载 ⇒ 选择无状态 + 块级远程访问;AFS 假设读多写少的中小文件 ⇒ 选择整文件缓存 + 回调;GFS 假设少量巨型文件的追加写 ⇒ 选择大 chunk + 单 master 元数据 + 松弛一致性。三者没有优劣,只有假设是否匹配。

22.3 算法伪代码与正确性分析

本节给出六个核心算法的假设 → 伪代码 → 逻辑解说 → 正确性论证(安全性/活性)→ 复杂度闭环。前两个算法刻画 NFS 与 AFS 的一致性机制,后四个算法刻画 GFS 的读、写、追加与修复。

算法 22.3.1:NFS 的无状态文件句柄与路径解析(LOOKUP)+ close-to-open 缓存验证

假设与系统模型

  • 进程:客户端进程与服务器进程都可能崩溃-恢复(crash-recovery);服务器无状态(不保存打开文件表、锁、缓存一致性记录),因此崩溃重启后无需恢复任何东西。
  • 通道:RPC 采用 at-least-once 语义(超时重发,Lecture 19-20);因此所有操作必须幂等
  • 时间:客户端与服务器时钟不同步,但服务器的时间戳(mtime/size)是唯一的”版本”权威;假设(理想化)mtime 精度足以区分相邻两次修改(现实中的秒级分辨率是这一保证的已知破口,见正确性论证)。
  • 文件句柄fh = (volume, inode, generation),客户端可以自由复制、缓存它;服务器不记得任何 fh 与客户端的关联
  • 客户端状态cache[path] = {fh, attrs, blocks, dirty};属性缓存有一个新鲜期 $t$(Solaris 实测:文件 3-30 s,目录 30-60 s)。

伪代码

# ================= 服务器(无状态,每次请求自带 fh)=================
upon request LOOKUP(fh_dir, name):
    d := resolve(fh_dir)                      # 校验 volume/inode/generation
    if d == NONE: return ESTALE               # generation 不匹配 => 老句柄
    if name not in d.entries: return ENOENT
    return fh_of(d.entries[name])             # 返回子文件的 fh

upon request GETATTR(fh):
    i := resolve(fh)
    if i == NONE: return ESTALE
    return (type, mode, owner, size, atime, mtime, ctime)

upon request READ(fh, offset, count):         # 幂等
    i := resolve(fh); if i == NONE: return ESTALE
    return i.data[offset : offset + count]

upon request WRITE(fh, offset, data):         # 幂等(同偏移写同数据 => 结果相同)
    i := resolve(fh); if i == NONE: return ESTALE
    i.data[offset : offset+len(data)] := data
    i.size := max(i.size, offset + len(data))
    i.mtime := server_clock()                 # 服务器时间戳 = 版本
    return OK

# ================= 客户端 =================
state:
    cache := { }        # path -> {fh, attrs, blocks, dirty};blocks[i] = (data, tag)
    attr_fresh_until := { }   # path -> 本地时间,属性缓存的新鲜期截止

function LOOKUP_PATH(path):                   # 逐级解析:每一级一个 RTT
    fh := MOUNT_ROOT_FH
    for name in components(path):
        fh := RPC LOOKUP(fh, name)            # NFSv4 用 COMPOUND 把多级打包成一个 RPC
    attrs := RPC GETATTR(fh)
    return fh, attrs

function OPEN(path):
    fh, attrs := LOOKUP_PATH(path)
    c := cache.get(path)
    if c == NONE or c.attrs.mtime != attrs.mtime or c.attrs.size != attrs.size:
        if c != NONE: drop(c.blocks)          # ★ 关键:验证失败 => 丢弃数据缓存
        cache[path] := {fh, attrs, blocks: {}, dirty: {}}
    attr_fresh_until[path] := now() + t
    return fd(path)                           # fd 是客户端本地的,服务器不知道

function READ(path, i):
    c := cache[path]
    if c.blocks.has(i):                       # ★ 会话内不再验证 => 可能读到陈旧数据
        return c.blocks[i].data
    d := RPC READ(c.fh, i*BS, BS)
    c.blocks[i] := (d, now())
    return d

function WRITE(path, i, data):                # delayed write:先改本地
    cache[path].blocks[i] := (data, now())
    cache[path].dirty.add(i)

function CLOSE(path):                         # ★ 关闭时写回所有脏块
    for i in cache[path].dirty:
        RPC WRITE(cache[path].fh, i*BS, cache[path].blocks[i].data)
    cache[path].dirty := {}

算法逻辑解说(含一个具体数值例子)

BS = 8 KB,客户端 A、B 的本地时钟相同,服务器文件 F 的 mtime 初始为 100。

时刻A 的动作B 的动作服务器状态
$t_1$OPEN(F):4 级路径 ⇒ 4 次 LOOKUP + 1 次 GETATTR(5 RTT)⇒ 拿到 mtime=100,缓存块mtime=100
$t_2$READ(F, 3):未命中 ⇒ 1 RTT,缓存 block3(标签 mtime=100不变
$t_3$OPEN(F)WRITE(F,3,'ZZZZ')CLOSE(F) ⇒ 1 次 WRITE 落服务器mtime=101
$t_4$READ(F, 3)命中缓存,直接返回旧的 AAAA陈旧读mtime=101
$t_5$CLOSE(F)OPEN(F):GETATTR 得 mtime=101 != 100 ⇒ 丢弃缓存mtime=101
$t_6$READ(F,3):未命中 ⇒ 1 RTT ⇒ 读到 ZZZZmtime=101

谁发起、消息如何流转:所有请求由客户端发起;服务器从不主动发消息(这是 NFS 与 AFS 最本质的区别)。路径解析是逐级的,因此路径越长,打开越贵;一旦拿到 fh,后续读写都是”一次请求一次响应”。

正确性论证

安全性(Safety,分两条:一条保证,一条明确不保证)

(S1) close-to-open 保证成立若客户端 A 在 $t_c$ 关闭对 F 的写,客户端 B 在 $t_o > t_c$ 打开 F,则 B 在此之后对 F 的读能看到 A 的写。

证明:A 的 CLOSE 会把它所有的脏块用 WRITE 送到服务器(伪代码 CLOSE),且每个 WRITE 都把服务器上的 mtime 更新为服务器当前时间,因此在 $t_c$ 之后,服务器上 F 的 mtime 严格大于 B 缓存里的 mtime(或 B 根本没有缓存)。B 的 OPEN 必然调用 GETATTR(伪代码 OPEN 的第一步),得到的 attrs.mtime != c.attrs.mtime,于是执行 drop(c.blocks);此后 B 的 READ 未命中缓存,只能向服务器要,于是返回的是包含 A 写入的最新数据。$\blacksquare$

这条证明依赖三个假设,缺一不可(这正是”课堂结论”与”生产真相”的距离):

  1. 写回发生在 close 时刻(delayed-write-at-close)。若实现把脏块的写回继续推迟到 close 之后(例如某些缓存超时策略),$t_c$ 时刻服务器还没有数据,保证失效。
  2. open 时确实执行了验证。若属性缓存在新鲜期 $t$ 内(默认最小值 acregmin = 3 s)而实现跳过了 GETATTR,那么 $t_o - t_c < 3$ s 时就可能读不到 A 的写。这就是”close-to-open 保证”必须与”属性缓存新鲜期”联合陈述的原因。
  3. mtime 能区分相邻两次修改mtime 若只有秒级分辨率,A 在同一秒内对 F 的两次修改、或”A 修改后 B 又改回同一秒”,都可能让 B 认为缓存仍然有效。工程上因此同时比较 sizectime

(S2) 明确不提供的保证对已打开的文件,NFS 不保证多客户端一致,也不保证任何原子性。

给出具体反例(对应上表的 $t_4$):设 A 在 $t_1$ 打开 F 并在 $t_2$ 读取并缓存 block 3;B 在 $t_3 \in (t_1, t_4)$ 写入 block 3 并关闭;A 在 $t_4$ 再读 block 3。A 的 READ(伪代码)在缓存命中时既不验证属性也不联系服务器,因此返回的是 $t_2$ 时刻的数据,未反映 B 在 $t_3$ 已提交到服务器的写。所以 close-to-open 不蕴含“并发打开者互相可见”,甚至不蕴含 read-your-writes-across-clients。此外,A 与 B 各自做 read(block3) → 修改 → write(block3) 时,后写者会完整覆盖前写者的结果(lost update)——NFS 的 WRITE 是”整块覆盖”,没有跨客户端的原子性。$\blacksquare$

活性(Liveness)

  • 请求最终被满足:RPC 层保证”重发直到收到响应”,且所有操作幂等READ/WRITE/LOOKUP/GETATTR 重复执行结果相同),因此重发不会破坏状态。
  • 服务器崩溃后恢复时间为 0:服务器无状态,重启即可服务(对比 AFS 需要重建回调表,见 22.3.2)。
  • 锁与一致性检查的活性不在本算法内:锁由 NLM 提供,崩溃后需要 NSM 通知持锁者(否则锁会被永久持有)——这是”无状态”外挂出来的一整套协议。

复杂度

操作消息复杂度说明
LOOKUP_PATH(冷)$d + 1$ 个 RTT($d$ = 路径深度)NFSv4 的 COMPOUND 合并为 1 个 RTT
OPEN$d+1$(冷)或 0(属性缓存新鲜期内,代价是一致性) 
READ(缓存命中)0但可能陈旧
READ(未命中)1 个 RTT稳态下服务器请求数 $\propto$ 块访问次数
WRITE + CLOSE脏块数(每个块一次 RPC)关闭时批量写回
服务器空间$O(1)$(不保存客户端状态)可支撑大量客户端
客户端空间$O(\text{缓存大小})$块级 + 属性

这一条复杂度公式就是 NFS 的软肋:服务器请求数 $L_{\text{NFS}} \propto N \cdot r_{\text{block}}$,与”读了多少块”成正比 ⇒ 客户端规模一上来,服务器先垮(22.2.11)。

算法 22.3.2:AFS 的整文件缓存 + 回调承诺(callback promise)失效协议

假设与系统模型

  • 进程:Venus(客户端)与 Vice(服务器);服务器有状态,崩溃后状态丢失(crash-recovery,且”恢复”指的是服务器进程重启,而不是恢复到崩溃前)。
  • 通道:Fetch/Store/Callback 都是 RPC;回调用可靠通道传输(TCP),若长时间无法送达,客户端最终把承诺视为作废(保守)。
  • 文件模型:整文件读写;每次 Store 使版本号单调递增Store 的内容是整个文件(不做差分)。
  • 关键不变量(本算法要证明的核心):客户端持有的 promise = VALID 蕴含其本地副本等于服务器当前版本。

伪代码

# ================= 服务器 Vice(有状态)=================
state:
    files[f]     := {data, version}          # 持久化
    callbacks[f] := set of client ids        # ★ 谁缓存了 f(内存,崩溃即丢)
    epoch        := 递增的“状态世代号”

upon receive Fetch(f, cid) from Venus_c:
    callbacks[f].add(cid)                    # 登记:以后 f 改变要通知 c
    send (files[f].data, files[f].version) to Venus_c

upon receive Store(f, data, cid) from Venus_c:
    files[f].data    := data
    files[f].version := files[f].version + 1        # 版本号递增
    for c in callbacks[f] \ {cid}:                  # ★ 推送失效(push)
        send BreakCallback(f) to Venus_c
    callbacks[f] := {cid}
    send files[f].version to Venus_c

upon restart:                                       # 崩溃恢复:callback 表已丢失
    callbacks := {}                                 # 忘记谁缓存了什么
    epoch     := epoch + 1                          # 世代号变化 = 告诉客户端“承诺全废”
    # 客户端在下一次接触服务器时发现 epoch 变化,把所有 promise 置为 CANCELED

# ================= 客户端 Venus =================
state: cache[f] := {data, version, promise ∈ {VALID, CANCELED}, dirty}

upon Open(f):
    if cache[f] exists and cache[f].promise == VALID:
        return local_copy                        # ★ 0 次 RPC:可扩展性的根源
    (data, v) := RPC Fetch(f, myid)
    cache[f] := {data, version: v, promise: VALID, dirty: false}
    return local_copy

upon Read(f, i):  return cache[f].data[i]         # 0 次 RPC,且必定是最新版本

upon Write(f, i, d):                              # 乐观写:只改本地
    cache[f].data[i] := d; cache[f].dirty := true

upon Close(f):
    if cache[f].dirty:
        cache[f].version := RPC Store(f, cache[f].data, myid)
        cache[f].dirty := false
    cache[f].promise := VALID                     # 写回者持有最新版本

upon receive BreakCallback(f) from Vice:
    cache[f].promise := CANCELED                  # 二值状态,无需版本比较
    drop(cache[f])                                # 丢弃整份本地副本

upon detect epoch_changed:                        # 服务器重启过
    for all f: cache[f].promise := CANCELED; drop(cache[f])

算法逻辑解说

以讲义的口径走一遍:Venus 打开文件 → Vice 发送整个文件并给出 callback promise;Venus 在本地读/写;文件关闭时把写传播回 Vice;Vice 递增版本号并回调所有其他缓存者。客户端的状态只有二值validcanceled——这是 AFS 的一个优雅之处:客户端不需要保存版本号、也不需要做时间戳比较,只需要回答”我的承诺还有效吗”。

一个数值例子:文件 F 在服务器上 version = 7

  1. Venus_A Open(F):本地无副本 ⇒ Fetch ⇒ 得到 (data₇, 7)promise = VALID,Vice 的 callbacks[F] = {A}
  2. Venus_B Open(F):同样 Fetch(data₇, 7)callbacks[F] = {A, B}
  3. A 修改 F 并 CloseStore(F, data₈, A) ⇒ 服务器 version = 8,向 BBreakCallback(F)callbacks[F] = {A}
  4. B 的 promise 变为 CANCELED 并丢弃本地副本;B 下一次 Open(F) 时承诺无效 ⇒ 重新 Fetch ⇒ 拿到 data₈ ✅(注意:B 不需要 close 再 open 也能看到,只要它重新 open;而且它不会”悄悄地”继续用旧数据
  5. 若此时 Vice 崩溃重启:callbacks 清空、epoch++;A 与 B 在下次接触时把各自的 promise 置为 CANCELED,各自重新 Fetch一次”惊群”),随后重新登记。

正确性论证

安全性(Safety)

(S1) 核心不变量若客户端 c 的 cache[f].promise = VALID,则 cache[f].data 等于服务器上 files[f].data 的当前内容。

证明(对”版本号变化事件”做归纳):令 $E$ 表示任意一次 Store(f, ·, ·) 事件(它是唯一能改变 files[f].version 的事件;Fetch 只读)。

  • Fetch 返回 (data, v) 时,服务器此刻的 version 就是 $v$,且 data 就是那一刻的 files[f].data;客户端据此把 promise 置为 VALID。不变量成立。
  • 归纳步:假设不变量在 $E$ 之前成立,且 c 的 promise = VALID(即 c 是在某次 Fetch 后未被失效的缓存者)。事件 $E$ 发生时,代码执行 for c' in callbacks[f] \ {cid}: send BreakCallback。而 c 在 Fetch 时已被 callbacks[f].add(cid) 登记,此后只有在两种情况下会离开 callbacks[f]:(i) 另一次 Store 把它重置掉——但那次 Store 同样会给 c 发 BreakCallback;(ii) 服务器重启——此时客户端会因 epoch 变化而把 promise 置为 CANCELED(代码最后一段)。因此在 $E$ 之后,c 要么已经收到 BreakCallbackpromise 变为 CANCELED,不变量前提不成立),要么它自己就是写者(promise 保持 VALIDversion 被更新为 $v+1$,其数据就是新的 files[f].data)。两种情况下”promise = VALID ⇒ 数据最新”都成立。 $\blacksquare$
  • 依赖的假设:(a) 回调最终送达或客户端超时作废(若回调静默丢失且不超时,不变量就破了——这正是”崩溃恢复必须保守”的原因);(b) 服务器重启后不会假装自己没重启(必须递增 epoch 或让所有客户端失效)。

(S2) AFS 的可见性保证若客户端 A 对 F 完成一次 StoreClose 返回),则此后任何客户端 B 对 F 的 Open 都能看到 A 的修改。

证明:由 (S1),B 在 Open 时若 promise = VALID,则它手里的副本本来就等于服务器当前版本(由于 A 的 Store 发生在 B 的这次 Open 之前,B 若在这次 Store 之前 Fetch 过,则它已被 BreakCallback 置为 CANCELED,因此 promise != VALID);于是无非两种情况:promise 已 CANCELED ⇒ 重新 Fetch(看到 A 的写);promise 仍 VALID ⇒ 说明自 B 上次 Fetch 以来没有任何 Store 发生,但 A 的 Store 就在其之前,所以 B 的 Fetch 必然发生在 A 的 Store 之后——仍然看到 A 的写。$\blacksquare$

与 NFS 的对比(这才是重点):NFS 的 (S1) 依赖”open 时验证 + 属性新鲜期”,是拉取式(pull)的,存在一个”最长到新鲜期 $t$”的窗口;AFS 的 (S2) 是推送式(push)的:Store 一完成,失效立刻被推给所有缓存者。因此 AFS 的一致性窗口 ≈ 一次回调的传播延迟,而 NFS 的窗口 ≈ 属性缓存新鲜期

(S3) 明确不提供的保证:并发写会丢失更新。 反例:F 在服务器上 version = 7;A、B 同时 Open(F)(都拿到 data₇,都 promise = VALID);A 在本地把第 3 块改成 XClose ⇒ 服务器 version = 8data₈ = A 的整份副本;B 随后把第 5 块改成 YClose ⇒ 服务器接受 B 的 Store(没有任何版本冲突检查,也没有 if version == 8 的条件写)⇒ version = 9data₉ = B 的整份副本A 对第 3 块的修改被整体覆盖(lost update)。AFS 的 Store 是”整文件覆盖 + 无版本校验”,因此共享写的正确性必须由应用层锁保证(AFS 提供锁服务,但它会引入额外状态与调用,因此默认依赖”多数文件单用户”这一假设)。$\blacksquare$

(S4) 崩溃后的保守失效是安全的服务器崩溃丢失 callbacks 后,让所有客户端把 promise 置为 CANCELED,只会造成”多余的重新取回”,不会造成”使用过期数据”。

证明:CANCELED 是一个保守的标记——它把”可能过期”一律当作”确定过期”。由 Open 的代码,promise != VALID 一律触发一次 Fetch,于是客户端在此之后拿到的必定是服务器当前版本(Fetch 返回的是 files[f].data 的实时快照)。错误的方向只会是”把最新数据当成过期”(性能损失),永远不会是”把过期数据当成最新”(安全性损失)。 $\blacksquare$

活性(Liveness)

  • 每次 Open 在有限时间内返回:最多一次 Fetch(1 个 RTT,整文件传输时间)。
  • 失效后承诺可被重建:下一次 Fetch 成功即重新建立 promise = VALID
  • 活性风险:服务器重启会导致全部客户端同时重新 Fetch(惊群 / thundering herd);服务器突发负载会短暂飙升——这是”有状态设计 + 保守失效”的代价。生产 AFS 用回调超时分级客户端重试抖动来缓解。

复杂度

维度复杂度
消息Open:0(承诺有效)或 1 次整文件传输;Close:0(未改)或 1 次整文件传输;每次 Store 触发 $O(\lvert\text{callbacks}[f]\rvert)$ 次回调
服务器空间$O(\sum_f \lvert\text{callbacks}[f]\rvert)$ = $O(\text{客户端数} \times \text{被缓存文件数})$(这就是”用状态换一致性”的账单)
客户端空间$O(\text{缓存文件字节数})$(约 100 MB 量级,持久化)
服务器负载$\propto$ 打开/关闭次数 + 失效次数与”打开后读了多少遍、读了多少块”无关 ⇒ 这是 AFS 可扩展性的形式化表述

算法 22.3.3:GFS 的读流程(客户端 ↔ master ↔ chunkserver)

假设与系统模型

  • 进程:1 个 master(元数据,在内存中)、$m$ 个 chunkserver(数据 + 校验和)、任意多个客户端;crash-recovery 故障模型。
  • 数据布局:文件按 64 MB 切成 chunk(最后一个可不满),每个 chunk 有 3 个副本,分布在不同机架(2 + 1)。
  • 通道:客户端-master 与客户端-chunkserver 都是 RPC;读操作幂等,可安全重试并切换副本。
  • 元数据缓存:客户端缓存 (文件, chunk index) -> (handle, replicas),带超时;缓存文件数据被明确排除(工作负载是流式顺序读)。

伪代码

# ================= 客户端 =================
state: meta_cache := { }        # (file, chunk_index) -> (handle, [replica...], expire_at)

procedure READ(file, offset, length):
    ci := offset / CHUNK_SIZE                     # ① 本地换算,无需网络
    (handle, replicas) := GET_CHUNK_LOCATION(file, ci)
    best := argmin_{r in replicas} distance(client, r)     # ④ 选“最近”的副本
    for r in order_by_distance(replicas) with best first:
        reply := RPC READ(r, handle, offset % CHUNK_SIZE, length)
        if reply.status == OK:
            return reply.data                     # ⑤ 成功返回
        if reply.status == CORRUPT:
            RPC REPORT_CORRUPT(master, handle, r)  # 报告损坏,换副本再试
    return ERROR

function GET_CHUNK_LOCATION(file, ci):            # ② ③ 元数据:问 master 并缓存
    e := meta_cache.get((file, ci))
    if e != NONE and e.expire_at > now(): return e.handle, e.replicas
    reply := RPC GET_LOCATION(master, file, ci)
    meta_cache[(file, ci)] := (reply.handle, reply.replicas, now() + TTL)
    return reply.handle, reply.replicas

# ================= master =================
upon receive GET_LOCATION(file, ci) from client:
    if (file, ci) not in chunk_map:               # 必要时分配新 chunk(写路径才会发生)
        h := allocate_chunk(); chunk_map[(file, ci)] := h
        for r in pick_replicas(k = 3, by_rack):   # 2 个同机架 + 1 个异机架
            replica_map[h].add(r); send CREATE_CHUNK(h) to r
    send (chunk_map[(file, ci)], replica_map[h]) to client
    # master 只做元数据:不转发数据、不参与数据路径

# ================= chunkserver =================
upon receive READ(handle, off, len) from client:
    if checksum_ok(handle) == false:              # 先验校验和(64 KB 一块,32 位)
        send CORRUPT(handle) to client            # 拒绝返回损坏数据
    else:
        send data[handle][off : off+len] to client

算法逻辑解说

冷路径:客户端问 master(1 RTT)→ 客户端缓存这份元数据 → 直接向最近的 chunkserver 读(1 RTT)⇒ 共 2 个 RTT。热路径(同一 chunk 内的连续读):只有 1 个 RTT,master 完全不参与。数值例子:读文件 F 的 [0, 192 MB),chunk 大小 64 MB ⇒ 需要 chunk 0、1、2 三个 chunk ⇒ 3 次元数据请求(每个 chunk 一次,之后缓存在客户端)⇒ 之后对每个 chunk 的所有顺序读都只有数据 RTT。如果副本 cs2 返回 CORRUPT,客户端把事件报告给 master 并改读 cs1,用户看到的只是这一次读多花了一个 RTT(22.3.6 处理修复)。

正确性论证

安全性

  • (S1) 返回值来自一个副本的完整数据READ 只从一个 chunkserver 取 [off, off+len),返回的数据是该副本上该区间的字节;checksum_ok 保证了返回的数据没有静默损坏(否则返回 CORRUPT 而不是数据)。
  • (S2) 副本内容一致(在”上一次写成功”的前提下):由算法 22.3.4 的写串行化保证(primary 分配序列号 ⇒ 所有副本按同一顺序应用同一组写),因此从任意副本读都得到同样的内容。例外:上一次写部分失败(某些副本没写成功)⇒ 副本之间可能不同,此时读可能返回”不一致的内容”,这正是 GFS 一致性模型中 inconsistent 的来源,也是 22.3.6 的 checksum 对比所要检测的对象。
  • (S3) 元数据可能过期,但不会读到错误数据:客户端缓存的副本位置可能指向已下线的 chunkserver ⇒ 读会失败或超时 ⇒ 客户端换一个副本重试(读是幂等的,重试安全)。“过期”只影响可用性与延迟,不影响正确性——因为位置只是”候选地址”,不是数据本身。

活性

  • 只要至少 1 个副本可用,读就成功:3 副本容忍 2 个副本同时不可用($f \le k-1 = 2$),这是”复制换可用性”的直接体现。
  • master 不可用时:已缓存的元数据仍能让部分读成功(这也是元数据缓存的价值)。

复杂度

维度复杂度
网络 RTT冷:2(1 元数据 + 1 数据);热:1
master 消息每个 chunk 每个客户端(每 TTL)1 次;顺序读一个大文件只需 1 次/chunk
数据流经 master0 字节(这是架构可扩展到数百 chunkserver 的关键)
空间客户端 $O(\#\text{chunk} \times \text{副本列表})$ 元数据(不缓存数据);master $O(\#\text{chunk})$

算法 22.3.4:GFS 的写流程与租约(本章最重要的伪代码)

假设与系统模型

  • 进程:master、$m$ 个 chunkserver、多个客户端;crash-recovery
  • 副本:每个 chunk 有 $k=3$ 个副本;任一时刻至多一个 primary(由 master 授予租约)。
  • 通道:控制消息走可靠 FIFO 通道(TCP)⇒ 同一发送者对同一接收者的消息不会乱序、不会丢失(这是顺序正确性的关键假设);数据流沿 pipeline 逐跳转发。
  • 时间:无全局时钟;租约用本地时钟度量;假设时钟漂移率有界(存在 $\rho$,使任意两个时钟的偏差增长率不超过 $\rho$)。租约时长 $T = 60$ s。
  • 数据先于控制:数据先推送到所有副本(进 LRU 缓冲),控制指令后到(只带 data_id 与偏移)。

伪代码

# ================= master =================
state: leases[h] := (primary, expire) | NONE      # 每个 chunk 至多一条租约
       replica_map[h] := [cs1, cs2, cs3]

upon receive GET_PRIMARY(h) from client:
    if leases[h] != NONE and leases[h].expire > now():
        reply (leases[h].primary, replica_map[h] \ {primary}, leases[h].expire)   # 复用租约
    else:
        p := choose(replica_map[h] \ {旧 primary})       # 换一个副本,避免“永远同一个主”
        expire := now() + T              # ★ 只有旧租约已过期才会走到这里
        leases[h] := (p, expire); log_lease(h, p, expire)     # 记入 operation log / 审计
        send GRANT_LEASE(h, expire) to p
        reply (p, replica_map[h] \ {p}, expire)

# ================= chunkserver(作为 primary)=================
state: lease_until[h]  # 我的租约到期时刻(本地时钟)
       serial[h] := 1  # 下一个序列号
       buffer[data_id] # 已经收到的数据(LRU)

upon receive GRANT_LEASE(h, expire): lease_until[h] := expire

upon receive PUSH_DATA(data_id, data, next_hop) from prev:      # 数据流(pipeline)
    buffer[data_id] := data
    if next_hop != NONE:
        ack := RPC PUSH_DATA(data_id, data, next_hop.next) to next_hop
        if ack != OK: send ERROR; return
    send ACK to prev                                            # ack 反向回传

upon receive APPLY(h, data_id, offset, secondaries) from client:  # 控制流
    if now() > lease_until[h]:                                  # ★ 自我 fencing
        send (ERROR, no-lease) to client; return
    s := serial[h]; serial[h] := s + 1                          # ★ 分配序列号(串行化)
    write buffer[data_id] into local chunk h at offset           # 先写自己
    for sec in secondaries:                                      # 按序列号顺序转发
        ack := RPC APPLY_REPLICA(h, s, offset, data_id) to sec
        if ack != OK: send (ERROR, sec_failed) to client; return  # 任一失败 => 报告失败
    send (OK, serial = s) to client

# ================= chunkserver(作为 secondary)=================
upon receive APPLY_REPLICA(h, s, offset, data_id) from primary:
    write buffer[data_id] into local chunk h at offset           # 按 primary 给的位置写
    send OK to primary

# ================= client =================
procedure WRITE(file, offset, data):
    retry := 0
    while retry < MAX_RETRY:
        h := chunk_of(file, offset); off := offset % CHUNK_SIZE
        (p, secs, expire) := GET_PRIMARY(h)          # ① 只有这一步需要 master
        data_id := fresh_id()
        ack := RPC PUSH_DATA(data_id, data, secs) to p    # ② 数据流:pipeline 推送
        if ack != OK: retry += 1; continue
        r := RPC APPLY(h, data_id, off, secs) to p        # ③ 控制流:只发给 primary
        if r.status == OK: return SUCCESS(r.serial)
        retry += 1                                        # ④ 失败 => 重试整个写
    return ERROR

算法逻辑解说(一次成功的写,7 步)

chunk H、primary = cs1、secondaries = {cs2, cs3}、写 1 MB 数据到 chunk 内偏移 4096 为例:

消息谁做定序说明
client → master:GET_PRIMARY(H)mastermaster 发现 H 无有效租约 ⇒ 选 cs1、授予 60 s 租约、回复 (cs1, [cs2,cs3], expire)
client → cs1 → cs2 → cs3:PUSH_DATA无人数据沿链推进;每跳把数据放进 LRU 缓冲(此时还不知道偏移);ack 反向回传
client → cs1:APPLY(H, data_id, 4096, [cs2,cs3])控制指令只给 primary,数据不重传(用 data_id 引用)
cs1 本地:分配 serial = 7,写入本地primary这一步就是串行化的发生地
cs1 → cs2 → cs3:APPLY_REPLICA(H, 7, 4096, data_id)primary按序列号顺序转发;secondary 按收到顺序执行
cs2、cs3 → cs1:OK只要有一个失败,cs1 就回 ERROR
cs1 → client:OK(serial = 7)客户端知道自己这次写的序列号

若 cs3 失败:cs1 在 ⑥ 察觉 ⇒ ⑦ 返回 ERROR(sec_failed) ⇒ 客户端重试整个写(回到 ①)。重试时会拿到新的序列号 8,于是 cs1 与 cs2 上可能留下两条数据(重复),而 cs3 上一条。这正是 GFS 写语义只能”至少一次”、一致性模型只能是”松弛”的根本原因。

正确性论证

(S1) 任一时刻至多一个 primary(租约互斥)

证明:master 是唯一的授权者,且它对每个 chunk 只保存一条租约记录 leases[h]GET_PRIMARY 的代码只有在 leases[h] == NONEleases[h].expire <= now() 时才写入新租约,因此两次授权之间必然隔着一个完整的租约期。形式化:设授权事件 $G_1$ 在 $t_1$ 生效于 $p_1$(到期 $t_1 + T$),$G_2$ 在 $t_2$ 生效于 $p_2$。因为 $G_2$ 执行时必须有 leases[h].expire <= t_2(master 的本地时钟读法),即 $t_1 + T \le t_2$。而旧的 primary $p_1$ 在自己的时钟读出的时刻超过 lease_until[h] 后会拒绝一切写请求(APPLY 里的 fencing 检查)。若两个时钟的漂移率有界($\rho$),则 master 需要等待 $T(1+\rho)$ 而不是 $T$,才能保证”$p_1$ 一定已经停止服务”;GFS 依赖这个有界漂移假设(工程上 $T$ 的取值远大于典型漂移,并配合 master 主动 revoke)。$\blacksquare$

(S2) 副本的最终内容一致(consistent)

证明:设 $S$ 是在所有副本上都成功应用的写集合。对 $S$ 中的写,每个副本都收到 primary 转发的 APPLY_REPLICA(h, s, offset, data_id);由于 (a) primary 是唯一给 secondary 发控制指令的实体,(b) primary→secondary 是可靠 FIFO 通道,(c) primary 按序列号单调递增的顺序发送,故每个 secondary 收到的顺序都是 $s_1 < s_2 < \cdots$ 的同一顺序;primary 自己也按同一顺序应用(它先分配序列号再本地写)。因此所有副本都按同一个顺序同一初始状态应用同一组确定性操作(”在偏移 off 写入字节串 $d$”是确定性函数)。由复制状态机的基本引理(相同初始状态 + 相同操作顺序 + 确定性操作 ⇒ 相同状态),三个副本的最终内容相同。$\blacksquare$

注意这里依赖”所有副本都在 $S$ 中”:若某次写只在部分副本上成功(⑥ 处有人失败),那么各副本应用的操作集合不同 ⇒ 内容可能不同inconsistent),这正是状态矩阵里”失败”那一格。

(S3) 串行成功 ⇒ 确定(defined)

若某个字节区间在整个过程中只被一个客户端写了一次且成功(无并发写),则该区间的内容就是那个客户端写入的字节串(客户端自己知道偏移与长度,也拿到了 OK),因此客户端能预测内容。反之:

(S4) 并发写 ⇒ 一致但未定义(consistent but undefined)——反证与构造

构造:客户端 A 在偏移 0 写 600 字节 'A'、客户端 B 在偏移 300 写 600 字节 'B',两者并发(A 与 B 都在同一个租约期内向同一 primary 发 APPLY)。

  • 由 (S2),无论 primary 选择哪个顺序,三副本内容都相同(consistent)
  • 但最终内容取决于序列号顺序:顺序 (B, A)[0,600)='A'[600,900)='B';顺序 (A, B)[0,300)='A'[300,900)='B'
  • 客户端无法观测序列号顺序serial 只有 primary 知道,且 A 收到的只是”我这次是 7”),因此 A 无法知道自己写入的 [300,600) 是否已被 B 覆盖内容对客户端不可预测(undefined)
  • 22.4.1 的实验 1 实测到了这两种布局(AAABBBBBBAAAAAABBB),并且每次运行三个副本都完全相同。$\blacksquare$

(S5) 为什么不用锁来消除”未定义”? 因为锁会引入服务器端状态与额外的 RPC(每次写前加锁/解锁各 1 个 RTT,且要处理持锁者崩溃),与 GFS”高吞吐、松弛语义、应用配合”的设计目标相反。GFS 把”需要定序语义”的需求收敛到 Record Append(22.3.5)这一个操作上。

活性(Liveness)

  • 写最终成功:客户端循环重试(上限 MAX_RETRY);只要 (a) 存在可用的 primary(master 会在租约过期后授予新租约),(b) 至少一个 secondary 可用,下一次尝试就可能成功。注意:GFS 的写要求”所有副本都成功”才算成功,因此容忍 0 个副本故障(这是”一致性优先于可用性”的取舍;与之相对,Cassandra 的 $W=1$ 是”可用性优先”)。
  • 租约到期后总会有新 primary:master 在旧租约过期后的第一次 GET_PRIMARY 就会授予新租约 ⇒ 不存在”永久无主”的状态。
  • 风险窗口(必须知道)旧租约刚过期、新租约刚授予的瞬间,客户端可能仍拿着缓存的旧 primary 地址去写;旧 primary 会因自身 fencing 检查(now() > lease_until)而拒绝,客户端收到 ERROR清空 primary 缓存并重新问 master只是延迟增加,不会出现双主写前提仍是时钟漂移有界:若旧 primary 的时钟走得慢(或它对 lease_until 的换算有误),它可能在 master 已经授权新 primary 之后仍然接受写 ⇒ 脑裂式的双主,两个 primary 各自分配序列号 ⇒ 副本内容永久分歧。这是 GFS 租约机制最需要小心的地方。

复杂度

维度复杂度对比
master 交互次数每个 chunk 每租约期(60 s)1 次,而不是每次写 1 次若没有租约:每次写都要问 master 谁是 primary ⇒ master 消息量 $\propto$ 写速率;有租约后降为 $\propto$ 写速率 / 60s,降低约 3 个数量级(对高频写而言)
数据流pipeline 上每跳传 1 份数据,共 $k-1$ 跳相对”客户端给每个副本各发一份”(1 份上行 × $k$),客户端上行从 $kF$ 降到 $F$(22.4.3 给出精确条件 $U < kB$)
控制消息$2$ 次客户端↔primary + $2(k-1)$ 次 primary↔secondary(请求 + ack)与数据量无关(这就是”数据流与控制流分离”的收益)
空间primary/secondary 需要 buffer[data_id](LRU 缓冲,写完后可丢弃)与”一次写入的数据量 × 并发写数”成正比

算法 22.3.5:GFS 的 Record Append(原子追加 + at-least-once)

假设与系统模型

  • 同 22.3.4(单 primary + 租约 + pipeline 数据流),额外假设:记录大小可控(通常远小于 64 MB),应用容忍重复记录
  • 目标语义:偏移由 GFS 决定;追加原子(不交错);至少一次

伪代码

# ================= client =================
procedure APPEND(file, record):
    loop:
        (h, ci) := RPC GET_LAST_CHUNK(master, file)      # ① 问 master:当前末尾 chunk
        (p, secs, expire) := GET_PRIMARY(h)
        data_id := fresh_id()
        ack := RPC PUSH_DATA(data_id, record, secs) to p  # ② 数据流(同普通写)
        if ack != OK: continue
        r := RPC APPEND(h, data_id, secs) to p            # ③ 控制流:偏移不由客户端给
        if r.status == OK:
            return (ci * CHUNK_SIZE + r.offset, r.offset)  # ④ 客户端从返回值得知偏移
        if r.status == RETRY_NEXT_CHUNK:                  # ⑤ chunk 已被 padding 填满
            RPC EXTEND(master, file, r.padded)            #    让 master 记账(padding 也算文件内容)
            continue                                      #    下一次循环会拿到新的末尾 chunk
        # r.status == ERROR:某个副本失败 => 重试(★ 可能重复)

# ================= chunkserver(primary)=================
upon receive APPEND(h, data_id, secondaries) from client:
    if now() > lease_until[h]: send (ERROR, no-lease) to client; return
    data := buffer[data_id]
    cur  := logical_size(h)                       # 当前 chunk 的逻辑末尾
    if cur + len(data) > CHUNK_SIZE:              # ★ 装不下 => padding + 让客户端换 chunk
        write zeros into h at [cur, CHUNK_SIZE)   # 把当前 chunk 填满(浪费空间,换取原子性)
        for sec in secondaries: RPC APPLY_PAD(h) to sec
        send (RETRY_NEXT_CHUNK, padded = CHUNK_SIZE - cur) to client
        return
    s := serial[h]; serial[h] := s + 1            # ★ 分配序列号
    off := cur                                    # ★ 偏移由 primary 决定(不是客户端)
    write data into h at off                      # 先写自己
    for sec in secondaries:
        ack := RPC APPLY_APPEND(h, s, off, data_id) to sec
        if ack != OK: send (ERROR, sec_failed) to client; return
    send (OK, offset = off) to client

算法逻辑解说

  • 偏移的来源:普通写的偏移是客户端给的;Record Append 的偏移是 primary 用 cur := logical_size(h) 现算的,而且在分配序列号之后才确定。这消除了”两个客户端把各自的记录写到同一偏移并互相覆盖”的可能。
  • 数值例子:chunk 逻辑大小 cur = 1000、记录长度 64。前 10 条记录依次落在 0、64、…、576(cur + 64 <= 1024)。第 11 条要写时 cur = 640?—— 让我们按 1024 的 chunk 算:可以放 $\lfloor 1024/64 \rfloor = 16$ 条记录。第 17 条写时 cur = 1024 - 64 = 960960 + 64 = 1024 <= 1024 仍可放(正好装满);第 18 条cur = 10241024 + 64 > 1024padding 0 字节padded = 0)⇒ 客户端拿到 RETRY_NEXT_CHUNK,master 分配 chunk 1,记录最终落在文件偏移 1024。22.4.1 的实验 4 就跑到这个边界并打印了 chunk 0 length after padding = 1024
  • 失败重试与重复:若 cs3 在 APPLY_APPEND 处失败,primary 已经把记录写在自己和 cs2 上,并将 ERROR 返回客户端;客户端重试⇒ primary 此时 cur 已经前进了 64 ⇒ 记录被写在新偏移 ⇒ cs1 与 cs2 上出现两条相同记录,cs3 上只有一条(实测输出正是:['C1-retry','C1-retry'] × cs1/cs2,而 cs3 是 ['<zeros>','C1-retry'])。

正确性论证

(S1) 不交错(no interleaving)任何一条成功应用(在所有副本上)的记录,其字节区间不会与另一条成功应用的记录的字节区间部分重叠。

证明:设记录 $R_1$(在序列号 $s_1$ 处)与 $R_2$(在序列号 $s_2$ 处)都成功。由 primary 的唯一性(22.3.4 的 S1),它们由同一个 primary 串行分配序列号;若 $s_1 < s_2$,则 primary 在处理 $R_2$ 时读到的 cur 已经包含了 $R_1$(因为处理 APPENDwrite data into h at off 立即增加了 logical_size(h))⇒ $R_2$ 的偏移 $off_2 \ge off_1 + \lvert R_1\rvert$,即两条记录的区间首尾相接但不重叠。若 $R_1$ 跨过 chunk 边界(装不下),primary 不写半条记录,而是 pad 满当前 chunk 并让客户端到下一个 chunk 重试 ⇒ 中间不会出现”半条记录 + 下一条记录的开头”。因此记录不会被劈开,也不会交错。$\blacksquare$

(S2) 偏移一致(所有成功副本在同一位移):由 (S1) 的证明,偏移由 primary 决定并在控制指令中显式携带APPLY_APPEND(h, s, off, data_id)),secondary 不做任何偏移计算,只按 primary 给的位置写。因此所有成功副本的字节位置相同 ⇒ 这与普通写的”一致(consistent)”是同一个机制。$\blacksquare$

(S3) at-least-once(至少一次):客户端在收到 OK 之前会一直重试(伪代码 loop);只要系统最终可用,记录至少被写入一次。但它不是 exactly-once

  • 失败重试可能让同一记录在部分副本上出现两次(重复,S1 的”成功”定义对”某些副本”不成立时);
  • 失败的副本可能完全缺少这条记录(留下空洞 = inconsistent 区域);
  • padding 会让文件里出现”零字节填充”(最多浪费一个 chunk 的空间)。 $\blacksquare$

(S4) 为什么 at-least-once 对应用是够用的:三条工程惯例把”至少一次”在应用层变回”精确一次”:(a) 每条记录带唯一 ID,消费者去重;(b) 记录带 checksum / 长度,读者可识别 padding 与损坏;(c) 应用只追加不修改,因此”重复记录”不会破坏已有数据(幂等消费)。$\blacksquare$

活性

  • 只要 primary 存在、且至少一个 secondaries 健康,APPEND 最终返回 OK(可能要重试)。
  • 跨 chunk 的活性:padding 之后,客户端通过 EXTEND(master, ...) 与重新 GET_LAST_CHUNK 强制 master 分配新 chunk(代码中 RETRY_NEXT_CHUNK 分支),因此不会出现”永远重试同一个满 chunk”的死循环
  • 热点风险(活性隐患):若成千上万个客户端同时追加同一个文件,它们会对同一个末尾 chunk 反复访问 master(GET_LAST_CHUNK)与该 chunk 的 primary ⇒ 元数据热点 + 单 chunkserver 写热点。GFS 论文承认这一点,建议应用分片写多个文件限制并发追加者数量

复杂度

维度复杂度
消息每次 APPEND:1 次 GET_LAST_CHUNK(可缓存)+ 1 次 PUSH_DATA 链 + 1 次 APPLY 链,即 $O(k)$ 条控制消息
空间浪费最坏每个 chunk 浪费 $(64\text{ MB} - \lvert R\rvert)$ 的 padding;期望上浪费 $\lvert R\rvert/2$ 每 chunk
偏移正确性不依赖客户端,只依赖 primary 的 cur 单调递增(由串行化保证)

算法 22.3.6:GFS 的 checksum 校验与副本修复

假设与系统模型

  • 故障模型:除了崩溃,还存在静默数据损坏(silent data corruption / bit rot):磁盘位翻转、内存位翻转、控制器 bug——数据变了但没有任何错误上报。这是”崩溃-恢复”模型之外必须单独处理的故障。
  • 校验粒度:chunk(64 MB)被分成 64 KB 的块每块一个 32 位校验和(一个 chunk 约 1024 个校验和,约 4 KB 校验数据)。
  • 假设:同一 chunk 的 3 个副本不会在修复完成前同时损坏(否则需要更快的检测或更多副本)。
  • 校验和本身也存放在 chunkserver 上(GFS 论文指出校验和放在 chunkserver 的持久存储中,与数据分开)。

伪代码

# ================= chunkserver 的存储层 =================
state: data[h]      # chunk 字节
       csum[h][b]   # 第 b 个 64 KB 块的 32 位校验和(b = 0 .. 1023)

function WRITE_BYTES(h, off, bytes):         # 任何写入路径都要更新校验和
    data[h][off : off+len(bytes)] := bytes
    for b in blocks_touched(off, len(bytes)):        # ★ 只重算受影响的块
        csum[h][b] := crc32(data[h][b*64KB : (b+1)*64KB])

function VERIFY(h) -> bool:
    for b in 0 .. num_blocks(h)-1:
        if crc32(data[h][b*64KB : (b+1)*64KB]) != csum[h][b]:
            return false
    return true

# ================= 读取路径 =================
upon receive READ(h, off, len) from client:
    if VERIFY(h) == false:                     # ★ 检测(detection)
        send CORRUPT(h) to client              #    绝不把损坏数据当作数据返回
    else:
        send data[h][off : off+len] to client

# ================= 客户端 =================
upon receive CORRUPT(h) from chunkserver r:
    RPC REPORT_CORRUPT(master, h, r)           # 报告 master:这个副本坏了
    for r' in replicas(h) \ {r}:               # ★ 立即改读其他副本(容错,读不中断)
        reply := RPC READ(r', off, len)
        if reply.status == OK: return reply.data
    return ERROR

# ================= master =================
upon receive REPORT_CORRUPT(h, bad) from client:
    good := replicas(h) \ {bad}                # 选一个健康副本作为修复源
    if good == {}: log("chunk lost"); return
    src := choose(good)
    mark_bad_replica(h, bad)                   # 不再把读请求路由给 bad
    send RE_REPLICATE(h, src) to bad           # ★ 修复(后台进行)

# ================= 被修复的 chunkserver =================
upon receive RE_REPLICATE(h, src):
    (bytes) := RPC FETCH_CHUNK(src, h)         # 取整个 chunk(源端同样先验校验和)
    if VERIFY_SOURCE(bytes) == false: abort    # 源也坏了 => 换一个源
    data[h] := bytes
    for b: csum[h][b] := crc32(...)            # ★ 重建校验和
    send OK to master

算法逻辑解说

三个动作要分清楚:

  1. 检测(detection):读取时先验校验和。因为校验和是独立保存、独立计算的,任何”数据变了但校验和没变”的位翻转都会被抓住。GFS 论文的口径:32 位 CRC 能抓住所有单比特错误与不超过 32 位的突发错误,随机的未检出概率约为 $2^{-32}$ 每 64 KB 块(约 4 GB 数据出现一次未检出错误的量级)——把”静默错误”变成了”可报告的显式错误”
  2. 容错(failover):客户端立刻改读其他副本,用户只感受到一次重试的延迟,而不是数据损坏。
  3. 修复(repair):master 让坏副本从健康副本重新复制整个 chunk(并重建校验和)。此外,GFS 还让各 chunkserver 周期性互相比较校验和(论文提到 chunkserver 会扫描并比较副本),以发现从未被读过的静默损坏——这些损坏不会在读取路径上被发现。
   cs1 [H: D | S]        cs2 [H: D | S]        cs3 [H: D | S]        初始:三副本一致
                                    |
                       位翻转(磁盘/内存/控制器)|  不改校验和
                                    v
   cs1 [H: D | S]        cs2 [H: D'| S]        cs3 [H: D | S]
                              |
   客户端读 cs2 -> 重算 crc(D') != S -> 返回 CORRUPT(检测)
   客户端 -> master: REPORT_CORRUPT(H, cs2);同时改读 cs1(容错,用户读成功)
   master -> cs2: RE_REPLICATE(H, src=cs1)(修复)
   cs2 从 cs1 取整个 chunk -> 重算全部校验和 -> 三副本恢复一致

正确性论证

(S1) 不返回损坏数据(安全性)客户端通过 READ 得到的任何数据,都通过了校验和验证。 证明:READ 的处理在返回数据之前检查 VERIFY(h)VERIFY每一个 64 KB 块重算 CRC 并与保存的校验和比较。若 data[h] 与写入时的字节不同(同长度改动),则至少有一个块的 CRC 不同(除非 CRC 碰撞,概率约 $2^{-32}$)⇒ VERIFY 返回 false ⇒ 返回 CORRUPT不是数据。因此”损坏数据被当作正确数据返回”的概率上界为 $O(2^{-32} \times \text{块数})$。$\blacksquare$

(S2) 修复恢复一致性修复完成后,被修复副本与源副本逐字节相同、校验和正确。 证明:修复执行 data[h] := FETCH_CHUNK(src, h)整份覆盖,不是增量补丁,因此不会留下”部分修补”的中间状态),随后对全部块重算 csum。由于源副本在读取时同样要过 VERIFY_SOURCE(否则换源),源数据是未被检测为损坏的版本。因此修复后两副本内容相同、校验和与内容自洽。$\blacksquare$

注意:这里的正确性依赖”源副本是正确的那一个“。若两个副本以相同方式损坏(例如同一批次软件 bug 写出相同错误数据),校验和无法区分 ⇒ 需要跨副本比较 + 更多的副本数(3 副本允许 1 个副本任意损坏,2 个副本同时损坏则只能靠”少数服从多数”或应用层校验)。

(S3) 检测能力与故障容忍度

故障是否被检测手段
单比特/短突发位翻转✅ 概率 1(CRC 性质)读取时校验
随机多比特损坏✅ 概率 $1 - 2^{-32}$(每块)读取时校验
整块丢失(磁盘坏道、文件被删)读取失败/长度不符
长时间无人读取的副本损坏✅(延迟)chunkserver 后台扫描 + 副本间校验和比较
副本内容被”逻辑上写错”(应用 bug)需要应用层校验(GFS 论文建议记录级 checksum)
3 个副本同时损坏❌(超出容忍度)需要更多副本 / 更快修复 / 纠删码(HDFS 3.x)

活性

  • 读服务不中断:只要有 1 个健康副本,客户端就能读到数据(先失败一次再换副本)。
  • 修复最终完成:master 把修复任务交给被修复的 chunkserver 后台执行(不阻塞读),完成后该副本重新参与服务。
  • 活性风险:如果”坏副本”数量超过副本数的一半,修复就无处可修;因此 GFS 有副本数下限监控(低于阈值就重新复制)与跨机架/跨机房布局

复杂度

维度复杂度
空间开销校验和约 $4\text{ KB} / 64\text{ MB} = 2^{-14}$ 的数据量(可忽略)
读开销每次读都要重算被读块的 CRC(CPU 换正确性;对顺序读是可接受的)
写开销只重算受影响的 64 KB 块(不是整个 chunk),因此随机小写的校验和代价为 $O(1)$
修复开销每个坏副本整 chunk 传输(64 MB)⇒ 与 chunk 大小成正比;这也是”大 chunk 的修复代价比小 chunk 高”的另一面
检测延迟读路径:立即;后台扫描:与扫描周期成正比(未被读过的损坏会被延迟发现

22.4 代码示例与分布式实现

本节给出三个可运行的 Python 程序(只用标准库,单机 python3 直接运行,固定随机种子):

  1. c22_gfs.py:一个简化版 GFS(1 个 master + 3 个 chunkserver + 多客户端),实现元数据/数据路径分离读流程带 pipeline 与租约的写流程失败重试Record Appendchecksum 检测位腐败与副本修复,并跑四个实验;
  2. c22_caches.py:NFS 式块级缓存 vs AFS 式整文件缓存 + 回调的定量对比(请求数 / 字节数 / 陈旧读次数 / 客户端规模扩展性);
  3. c22_pipeline.py:pipeline 数据流 vs 客户端直发多副本的带宽对比(解析式 + 逐块离散模拟互相校验)。

说明:第 1 个程序约 500 行,超出本笔记”单个代码块 60-160 行”的常规建议,因为它是一份完整的小型系统实现(通信骨架 + 3 类节点 + 4 个实验)。为便于阅读,下面把它切成 A、B 两块——把 A 与 B 依次粘进同一个 c22_gfs.py 即可直接运行(本笔记已实测,输出附在代码后)。

22.4.1 示例一:简化版 GFS(master + 3 chunkserver + 客户端 + 四个实验)

代码块 A:通信骨架 + ChunkServer + Master + 客户端(同一文件的上半部分)

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""c22-1: 简化版 GFS(单 master + 3 chunkserver + 多客户端)。
CHUNK_SIZE=1024 B 代表真实的 64 MB;LEASE_TTL=0.5 s 代表真实的 60 s。
四个实验:并发覆盖写 / 位腐败与 checksum / 租约与单 primary / Record Append。"""
import queue, random, threading, time, zlib

CHUNK_SIZE, REPLICATION, LEASE_TTL, CHECKSUM_BLOCK = 1024, 3, 0.5, 256
SERVERS = ["cs1", "cs2", "cs3"]
random.seed(425)


class Net:
    """消息总线:每个节点一个 inbox 队列,等价于一个网卡 + 端口。"""

    def __init__(self, quiet=False):
        self.nodes, self.trace, self.quiet = {}, [], quiet
        self.lock = threading.Lock()

    def register(self, name, node):
        self.nodes[name] = node

    def call(self, dst, msg, timeout=10.0):
        q = queue.Queue()
        m = dict(msg); m["reply"] = q
        self.nodes[dst].inbox.put(m)
        return q.get(timeout=timeout)

    def log(self, who, text):
        with self.lock:
            self.trace.append((who, text))
            if not self.quiet:
                print("      [%-5s] %s" % (who, text), flush=True)


class Node(threading.Thread):
    def __init__(self, net, name):
        super().__init__(daemon=True)
        self.net, self.name, self.inbox = net, name, queue.Queue()
        net.register(name, self)

    def run(self):
        while True:
            msg = self.inbox.get()
            try:
                self.handle(msg)
            except Exception as exc:                       # 节点内部错误
                self.net.log(self.name, "INTERNAL ERROR: %r" % (exc,))
                if msg.get("reply") is not None:
                    msg["reply"].put({"status": "error", "why": repr(exc)})

    def reply(self, msg, payload):
        if msg.get("reply") is not None:
            msg["reply"].put(payload)


class ChunkServer(Node):
    def __init__(self, net, name):
        super().__init__(net, name)
        self.chunks, self.sums, self.staging = {}, {}, {}   # chunk / 校验和 / 推送缓冲
        self.lease_until, self.serial, self.applied = {}, {}, {}
        self.fail_applies = 0                               # 故障注入:接下来 K 次写失败

    # ---------------- 存储 + 校验和 ----------------
    def _sums(self, buf):
        return [zlib.crc32(bytes(buf[i:i + CHECKSUM_BLOCK]))
                for i in range(0, max(len(buf), 1), CHECKSUM_BLOCK)]

    def verify(self, h):
        return self.sums.get(h) == self._sums(self.chunks[h])

    def put(self, h, off, data):
        if h not in self.chunks:
            self.chunks[h], self.sums[h] = bytearray(), []
        buf = self.chunks[h]
        if off + len(data) > len(buf):
            buf.extend(b"\x00" * (off + len(data) - len(buf)))
        buf[off:off + len(data)] = data
        self.sums[h] = self._sums(buf)                      # 每次写都重算校验和

    def corrupt(self, h, i):
        self.chunks[h][i] ^= 0xFF                           # 测试钩子:静默位翻转

    # ---------------- 消息处理 ----------------
    def handle(self, msg):
        t, h = msg["type"], msg.get("handle")
        if t == "create_chunk":
            self.chunks.setdefault(h, bytearray())
            self.sums.setdefault(h, self._sums(self.chunks[h]))
            self.reply(msg, {"status": "ok"})
        elif t == "grant_lease":
            self.lease_until[h] = msg["expire"]
            self.net.log(self.name, "lease GRANTED handle=%d until t=%.2f" % (h, msg["expire"]))
            self.reply(msg, {"status": "ok"})
        elif t == "push_data":                              # 数据流:沿 pipeline 链式推送
            self.staging[msg["data_id"]] = msg["data"]
            nxt = msg.get("next", [])
            ok = True
            if nxt:
                r = self.net.call(nxt[0], {"type": "push_data", "data_id": msg["data_id"],
                                           "data": msg["data"], "next": nxt[1:]})
                ok = r["status"] == "ok"
            self.reply(msg, {"status": "ok" if ok else "error"})
        elif t == "read":                                   # 读:先验校验和
            if not self.verify(h):
                self.net.log(self.name, "!! CHECKSUM MISMATCH handle=%d -> refuse" % h)
                self.reply(msg, {"status": "corrupt", "handle": h}); return
            self.reply(msg, {"status": "ok", "data":
                             bytes(self.chunks[h][msg["offset"]:msg["offset"] + msg["size"]])})
        elif t == "fetch_chunk":
            self.reply(msg, {"status": "ok", "data": bytes(self.chunks[h])})
        elif t == "replicate":                              # master 指令:从健康副本修复
            r = self.net.call(msg["source"], {"type": "fetch_chunk", "handle": h})
            self.chunks[h] = bytearray(r["data"])
            self.sums[h] = self._sums(self.chunks[h])
            self.net.log(self.name, "re-replicated handle=%d from %s" % (h, msg["source"]))
            self.reply(msg, {"status": "ok"})
        elif t == "apply":                                  # 控制流 -> primary
            if time.time() > self.lease_until.get(h, 0):
                self.reply(msg, {"status": "error", "why": "no-lease"}); return
            if self.fail_applies > 0:
                self.fail_applies -= 1
                self.reply(msg, {"status": "error", "why": "injected-failure"}); return
            s = self.serial.get(h, 1); self.serial[h] = s + 1
            data = self.staging[msg["data_id"]]
            self.put(h, msg["offset"], data)
            self.applied.setdefault(h, []).append((s, msg["offset"], len(data)))
            self.net.log(self.name, "PRIMARY apply handle=%d serial=%d offset=%d len=%d"
                         % (h, s, msg["offset"], len(data)))
            ok, why = True, None
            for sec in msg["secondaries"]:                  # 按序列号顺序转发给所有副本
                r = self.net.call(sec, {"type": "apply_replica", "handle": h, "serial": s,
                                        "offset": msg["offset"], "data_id": msg["data_id"]})
                if r["status"] != "ok":
                    ok, why = False, "%s failed: %s" % (sec, r.get("why"))
            self.reply(msg, {"status": "ok" if ok else "error", "serial": s, "why": why})
        elif t == "apply_replica":
            if self.fail_applies > 0:
                self.fail_applies -= 1
                self.reply(msg, {"status": "error", "why": "injected-failure"}); return
            data = self.staging[msg["data_id"]]
            self.put(h, msg["offset"], data)
            self.applied.setdefault(h, []).append((msg["serial"], msg["offset"], len(data)))
            self.reply(msg, {"status": "ok"})
        elif t in ("append", "apply_append", "apply_pad"):  # Record Append
            if t == "append":                               # 客户端 -> primary
                if time.time() > self.lease_until.get(h, 0):
                    self.reply(msg, {"status": "error", "why": "no-lease"}); return
                data = self.staging[msg["data_id"]]
                cur = len(self.chunks[h])
                if cur + len(data) > CHUNK_SIZE:            # 放不下:padding + 换 chunk
                    self.put(h, cur, b"\x00" * (CHUNK_SIZE - cur))
                    for sec in msg["secondaries"]:
                        self.net.call(sec, {"type": "apply_pad", "handle": h})
                    self.net.log(self.name, "PAD handle=%d with %d zeros -> retry next chunk"
                                 % (h, CHUNK_SIZE - cur))
                    self.reply(msg, {"status": "retry-next-chunk", "padded": CHUNK_SIZE - cur})
                    return
                if self.fail_applies > 0:
                    self.fail_applies -= 1
                    self.reply(msg, {"status": "error", "why": "injected-failure"}); return
                s = self.serial.get(h, 1); self.serial[h] = s + 1
                self.put(h, cur, data)                      # 偏移由 primary 决定
                self.applied.setdefault(h, []).append((s, cur, len(data)))
                self.net.log(self.name, "PRIMARY append handle=%d serial=%d offset=%d len=%d"
                             % (h, s, cur, len(data)))
                ok, why = True, None
                for sec in msg["secondaries"]:
                    r = self.net.call(sec, {"type": "apply_append", "handle": h, "serial": s,
                                            "offset": cur, "data_id": msg["data_id"]})
                    if r["status"] != "ok":
                        ok, why = False, "%s failed: %s" % (sec, r.get("why"))
                self.reply(msg, {"status": "ok" if ok else "error", "offset": cur, "why": why})
                return
            if self.fail_applies > 0:
                self.fail_applies -= 1
                self.reply(msg, {"status": "error", "why": "injected-failure"}); return
            if t == "apply_pad":
                self.put(h, len(self.chunks[h]),
                         b"\x00" * (CHUNK_SIZE - len(self.chunks[h])))
            else:
                data = self.staging[msg["data_id"]]
                self.put(h, msg["offset"], data)
                self.applied.setdefault(h, []).append((msg["serial"], msg["offset"], len(data)))
            self.reply(msg, {"status": "ok"})
        else:
            self.reply(msg, {"status": "error", "why": "unknown " + t})


class Master(Node):
    def __init__(self, net):
        super().__init__(net, "master")
        self.files, self.locations, self.leases = {}, {}, {}
        self.lease_history, self.next_handle, self.metadata_ops = [], 1, 0

    def handle(self, msg):
        self.metadata_ops += 1
        t = msg["type"]
        if t == "get_metadata":                             # 客户端问"chunk 在哪"
            p, idx = msg["path"], msg["offset"] // CHUNK_SIZE
            self.files.setdefault(p, {"chunks": [], "size": 0})
            self._ensure_chunk(p, idx)
            f = self.files[p]
            self.reply(msg, {"status": "ok", "size": f["size"],
                             "chunks": [(h, list(self.locations[h])) for h in f["chunks"]]})
        elif t == "get_primary":                            # 租约:同一时刻只有一个 primary
            h, now = msg["handle"], time.time()
            cur = self.leases.get(h)
            if cur and cur[1] > now:
                primary = cur[0]                            # 租约仍有效:直接复用
            else:
                old = cur[0] if cur else None
                primary = random.choice([s for s in self.locations[h] if s != old]
                                        or list(self.locations[h]))
                exp = now + LEASE_TTL
                self.leases[h] = (primary, exp)
                self.lease_history.append((h, primary, now, exp))
                self.net.call(primary, {"type": "grant_lease", "handle": h, "expire": exp})
                self.net.log("master", "GRANT lease handle=%d -> %s (t=%.2f..%.2f)"
                             % (h, primary, now, exp))
            self.reply(msg, {"status": "ok", "primary": primary, "expire": self.leases[h][1],
                             "secondaries": [s for s in self.locations[h] if s != primary]})
        elif t == "last_chunk":
            p = msg["path"]
            self.files.setdefault(p, {"chunks": [], "size": 0})
            if msg.get("force_new") and self.files[p]["chunks"]:
                self._ensure_chunk(p, len(self.files[p]["chunks"]))
            self._ensure_chunk(p, 0)
            i = len(self.files[p]["chunks"]) - 1
            self.reply(msg, {"status": "ok", "handle": self.files[p]["chunks"][i], "index": i})
        elif t == "extend":
            self.files.setdefault(msg["path"], {"chunks": [], "size": 0})
            self.files[msg["path"]]["size"] += msg["nbytes"]
            self.reply(msg, {"status": "ok"})
        elif t == "report_corrupt":                         # 位腐败 -> 指令重新复制
            h, bad = msg["handle"], msg["server"]
            good = [s for s in self.locations[h] if s != bad]
            if not good:
                self.reply(msg, {"status": "lost"}); return
            self.net.log("master", "corruption on %s handle=%d -> re-replicate from %s"
                         % (bad, h, good[0]))
            self.net.call(bad, {"type": "replicate", "handle": h, "source": good[0]})
            self.reply(msg, {"status": "ok", "from": good[0]})
        else:
            self.reply(msg, {"status": "error", "why": "unknown " + t})

    def _ensure_chunk(self, path, idx):
        f = self.files[path]
        while len(f["chunks"]) <= idx:
            h = self.next_handle; self.next_handle += 1
            self.locations[h] = SERVERS[:REPLICATION]
            f["chunks"].append(h)
            for s in self.locations[h]:
                self.net.call(s, {"type": "create_chunk", "handle": h})
            self.net.log("master", "allocate chunk#%d handle=%d replicas=%s"
                         % (len(f["chunks"]) - 1, h, self.locations[h]))


class GFSClient:
    def __init__(self, net, name):
        self.net, self.name, self.retries = net, name, 0
        self.chunk_cache, self.primary_cache = {}, {}

    def _meta(self, path, offset):                          # 元数据走 master,且被客户端缓存
        idx, now = offset // CHUNK_SIZE, time.time()
        ent = self.chunk_cache.get(path)
        if ent is None or idx >= len(ent) or ent[idx][2] < now:
            r = self.net.call("master", {"type": "get_metadata", "path": path,
                                         "offset": offset})
            ent = [(h, locs, now + 60.0) for h, locs in r["chunks"]]
            self.chunk_cache[path] = ent
        return ent[idx][0], list(ent[idx][1])

    def _primary(self, handle):                             # primary 位置也缓存(租约内有效)
        now, c = time.time(), self.primary_cache.get(handle)
        if c and c["expire"] > now:
            return c["primary"], list(c["secondaries"])
        r = self.net.call("master", {"type": "get_primary", "handle": handle})
        self.primary_cache[handle] = r
        return r["primary"], list(r["secondaries"])

    def read(self, path, offset, size, prefer=None):         # 数据路径:客户端直连 chunkserver
        h, locs = self._meta(path, offset)
        for s in sorted(locs, key=lambda x: 0 if x == prefer else 1):
            r = self.net.call(s, {"type": "read", "handle": h,
                                  "offset": offset % CHUNK_SIZE, "size": size})
            if r["status"] == "ok":
                return r["data"], s
            if r["status"] == "corrupt":                     # 报告 master,再换副本读
                self.net.call("master", {"type": "report_corrupt", "handle": h, "server": s})
        raise RuntimeError("all replicas failed")

    def write(self, path, offset, data, tag="w"):           # 普通写:客户端指定偏移
        out = []
        while data:
            h, _ = self._meta(path, offset)
            off = offset % CHUNK_SIZE
            piece, data = data[:CHUNK_SIZE - off], data[CHUNK_SIZE - off:]
            out.append(self._mutate(h, off, piece, "apply", tag))
            offset += len(piece)
        return out

    def append(self, path, record, tag="a"):                # Record Append:GFS 决定偏移
        for _ in range(6):
            info = self.net.call("master", {"type": "last_chunk", "path": path})
            h = info["handle"]
            prim, secs = self._primary(h)
            did = "%s-%s-%d-%d" % (self.name, tag, h, random.randint(0, 10 ** 6))
            if self.net.call(prim, {"type": "push_data", "data_id": did, "data": record,
                                    "next": secs})["status"] != "ok":
                self.retries += 1; continue
            r = self.net.call(prim, {"type": "append", "handle": h, "data_id": did,
                                     "secondaries": secs})
            if r["status"] == "ok":
                self.net.call("master", {"type": "extend", "path": path,
                                         "nbytes": len(record)})
                return info["index"] * CHUNK_SIZE + r["offset"]
            if r["status"] == "retry-next-chunk":             # chunk 被 padding 了
                self.net.call("master", {"type": "extend", "path": path, "nbytes": r["padded"]})
                self.net.call("master", {"type": "last_chunk", "path": path, "force_new": True})
                continue
            self.net.log(self.name, "append FAILED (%s) -> retry (duplicates possible)"
                         % r.get("why"))
            self.retries += 1
            self.primary_cache.pop(h, None)
        raise RuntimeError("append gave up")

    def _mutate(self, h, off, piece, op, tag):               # 写流程:push 数据 -> 发控制指令
        for attempt in range(1, 5):
            prim, secs = self._primary(h)
            did = "%s-%s-%d-%d" % (self.name, tag, h, attempt)
            time.sleep(random.random() * 0.01)               # 网络抖动
            if self.net.call(prim, {"type": "push_data", "data_id": did, "data": piece,
                                    "next": secs})["status"] != "ok":
                self.retries += 1; continue
            r = self.net.call(prim, {"type": op, "handle": h, "data_id": did,
                                     "offset": off, "secondaries": secs})
            if r["status"] == "ok":
                return (h, off, len(piece), r.get("serial"))
            self.net.log(self.name, "write FAILED (%s) -> retry the whole write" % r.get("why"))
            self.retries += 1
            self.primary_cache.pop(h, None)
        raise RuntimeError("write gave up")


def build(quiet=False):
    net = Net(quiet=quiet)
    master, servers = Master(net), [ChunkServer(net, s) for s in SERVERS]
    for n in [master] + servers:
        n.start()
    return net, master, servers

代码块 B:四个实验 + 入口(同一文件的下半部分,接在 A 之后)

# ============================================================ 实验 1:并发覆盖写
def experiment_1():
    print("\n=== EXP 1: concurrent overwrite -> consistent but undefined ===")
    layouts = []
    for trial, delay in ((1, 0.0), (2, 0.0), (3, 0.05)):
        net, master, servers = build(quiet=True)
        c1, c2 = GFSClient(net, "C1"), GFSClient(net, "C2")
        barrier = threading.Barrier(2)

        def worker(cli, fill, off, wait):
            barrier.wait()                                   # 两次写真正并发
            time.sleep(wait)                                 # 只改变网络时序,不改语义
            cli.write("/f1", off, fill)

        ts = [threading.Thread(target=worker, args=(c1, b"A" * 600, 0, delay)),
              threading.Thread(target=worker, args=(c2, b"B" * 600, 300, 0.0))]
        for t in ts: t.start()
        for t in ts: t.join()
        h = net.call("master", {"type": "get_metadata", "path": "/f1",
                                "offset": 0})["chunks"][0][0]
        views = [bytes(s.chunks[h][:900]) for s in servers]
        v = views[0]
        layout = "".join("A" if v[i:i + 100] == b"A" * 100 else
                         "B" if v[i:i + 100] == b"B" * 100 else "?" for i in range(0, 900, 100))
        layouts.append(layout)
        print("  trial %d: replicas identical = %-5s | region[0:900] = %s"
              % (trial, len(set(views)) == 1, layout))
        print("           primary applied (serial, offset, len) = %s" % (servers[0].applied[h],))
        print("           C1's 600 bytes survive intact afterwards = %s"
              % (v[:600] == b"A" * 600))
        assert len(set(views)) == 1, "GFS 的写必须在所有副本上一致"
    print("  => 所有副本内容相同(CONSISTENT),但最终保留的是哪个客户端的片段,")
    print("     由客户端无法观测的序列顺序决定(UNDEFINED)。同样的两次写,")
    print("     只因网络时序微变就产生了两种不同的交错:")
    for t, lay in zip((1, 2, 3), layouts):
        print("       run %d: %s" % (t, lay))


# ============================================================ 实验 2:位腐败
def experiment_2():
    print("\n=== EXP 2: silent bit rot -> checksum detects, healthy replica serves ===")
    net, master, servers = build(quiet=True)
    c = GFSClient(net, "C1")
    payload = bytes(range(0, 200)) * 3                        # 600 字节确定性数据
    c.write("/f2", 0, payload)
    h = net.call("master", {"type": "get_metadata", "path": "/f2",
                            "offset": 0})["chunks"][0][0]
    print("  after write : every replica verifies = %s" % [s.verify(h) for s in servers])
    victim = servers[1]
    victim.corrupt(h, 100)                                    # 静默位翻转
    print("  inject      : flip byte 100 of %s only -> its checksum check = %s"
          % (victim.name, victim.verify(h)))
    data, served = c.read("/f2", 0, 600, prefer=victim.name)
    print("  client read : served by %s, %d bytes, matches original = %s"
          % (served, len(data), data == payload))
    print("  after repair: %s verifies again = %s" % (victim.name, victim.verify(h)))
    assert data == payload and served != victim.name and victim.verify(h)
    print("  checksum-mismatch events on the wire = %d"
          % sum(1 for _, t in net.trace if "MISMATCH" in t))


# ============================================================ 实验 3:租约
def experiment_3():
    print("\n=== EXP 3: leases -> exactly one primary, few master interactions ===")
    net, master, servers = build(quiet=True)
    c = GFSClient(net, "C1")
    c.write("/f3", 0, b"x" * 100)
    h = net.call("master", {"type": "get_metadata", "path": "/f3",
                            "offset": 0})["chunks"][0][0]
    p1 = net.call("master", {"type": "get_primary", "handle": h})
    p1b = net.call("master", {"type": "get_primary", "handle": h})
    print("  primary #1 = %s ; asked again within the lease -> %s (same lease reused)"
          % (p1["primary"], p1b["primary"]))
    assert p1["primary"] == p1b["primary"] and len(master.lease_history) == 1
    old = p1["primary"]
    time.sleep(LEASE_TTL + 0.1)                               # 等租约过期
    p2 = net.call("master", {"type": "get_primary", "handle": h})
    print("  primary #2 = %s ; changed = %s" % (p2["primary"], p2["primary"] != old))
    assert p2["primary"] != old
    r = net.call(old, {"type": "apply", "handle": h, "data_id": "stale", "offset": 0,
                       "secondaries": p2["secondaries"]})
    print("  write to the OLD primary after expiry -> %r (fenced out)" % r)
    assert r["status"] == "error" and r["why"] == "no-lease"
    hs = sorted([x for x in master.lease_history if x[0] == h], key=lambda x: x[2])
    overlap = any(hs[i][3] > hs[i + 1][2] for i in range(len(hs) - 1))
    print("  lease audit (handle, primary, start, expiry) = %s"
          % [(x[0], x[1], round(x[2], 2), round(x[3], 2)) for x in hs])
    print("  two primaries overlapping in time = %s" % overlap)
    assert not overlap
    print("  master metadata ops = %d ; master interactions per write = 0"
          % master.metadata_ops)


# ============================================================ 实验 4:Record Append
def make_record(tag, size=64):
    return tag.encode() + b"." * (size - len(tag.encode()))


def decode(buf, size=64):
    out = []
    for i in range(0, len(buf) - size + 1, size):
        rec = bytes(buf[i:i + size])
        out.append("<zeros>" if rec == b"\x00" * size else rec.rstrip(b".").decode())
    return out


def experiment_4():
    print("\n=== EXP 4: record append -> no interleaving, at-least-once, padding ===")
    net, master, servers = build(quiet=True)
    h = net.call("master", {"type": "last_chunk", "path": "/log"})["handle"]
    clients = [GFSClient(net, "C%d" % (i + 1)) for i in range(3)]
    got, lock = {}, threading.Lock()

    def producer(cli):
        for i in range(4):
            off = cli.append("/log", make_record("%s-r%d" % (cli.name, i)))
            with lock:
                got.setdefault(cli.name, []).append(off)

    ts = [threading.Thread(target=producer, args=(c,)) for c in clients]
    for t in ts: t.start()
    for t in ts: t.join()
    print("  offsets returned to each client: %s" % got)
    views = [decode(s.chunks[h]) for s in servers]
    print("  replica cs1 record sequence: %s" % views[0])
    print("  all 3 replicas hold the SAME sequence = %s -> no interleaving"
          % all(v == views[0] for v in views))
    assert all(v == views[0] for v in views)

    print("  -- inject one failing secondary, the client must retry --")
    servers[2].fail_applies = 1
    off = clients[0].append("/log", make_record("C1-retry"))
    views = [decode(s.chunks[h]) for s in servers]
    for name, v in zip(SERVERS, views):
        print("    %s -> %s" % (name, v))
    print("  retry returned offset %d ; duplicate on some replicas = %s (AT-LEAST-ONCE)"
          % (off, any(v.count("C1-retry") > 1 for v in views)))
    print("  replicas identical after the failure = %s (a failed append can leave a hole)"
          % (len(set(map(tuple, views))) == 1))

    print("  -- keep appending until the chunk boundary is crossed --")
    while True:
        off = clients[0].append("/log", make_record("fill%02d"
                                                    % len(decode(servers[0].chunks[h]))))
        if off >= CHUNK_SIZE:
            break
    print("  chunk 0 length after padding = %d (CHUNK_SIZE = %d)"
          % (len(servers[0].chunks[h]), CHUNK_SIZE))
    print("  the record that did not fit landed at file offset %d -> chunk 1" % off)
    assert len(servers[0].chunks[h]) == CHUNK_SIZE and off == CHUNK_SIZE


if __name__ == "__main__":
    experiment_1(); experiment_2(); experiment_3(); experiment_4()
    print("\nALL GFS EXPERIMENTS PASSED")

运行输出(节选,python3 c22_gfs.py,Python 3.9 实测)


=== EXP 1: concurrent overwrite -> consistent but undefined ===
  trial 1: replicas identical = True  | region[0:900] = AAABBBBBB
           primary applied (serial, offset, len) = [(1, 0, 600), (2, 300, 600)]
           C1's 600 bytes survive intact afterwards = False
  trial 2: replicas identical = True  | region[0:900] = AAABBBBBB
           primary applied (serial, offset, len) = [(1, 0, 600), (2, 300, 600)]
           C1's 600 bytes survive intact afterwards = False
  trial 3: replicas identical = True  | region[0:900] = AAAAAABBB
           primary applied (serial, offset, len) = [(1, 300, 600), (2, 0, 600)]
           C1's 600 bytes survive intact afterwards = True
  => 所有副本内容相同(CONSISTENT),但最终保留的是哪个客户端的片段,
     由客户端无法观测的序列顺序决定(UNDEFINED)。同样的两次写,
     只因网络时序微变就产生了两种不同的交错:
       run 1: AAABBBBBB
       run 2: AAABBBBBB
       run 3: AAAAAABBB

=== EXP 2: silent bit rot -> checksum detects, healthy replica serves ===
  after write : every replica verifies = [True, True, True]
  inject      : flip byte 100 of cs2 only -> its checksum check = False
  client read : served by cs1, 600 bytes, matches original = True
  after repair: cs2 verifies again = True
  checksum-mismatch events on the wire = 1

=== EXP 3: leases -> exactly one primary, few master interactions ===
  primary #1 = cs2 ; asked again within the lease -> cs2 (same lease reused)
  primary #2 = cs1 ; changed = True
  write to the OLD primary after expiry -> {'status': 'error', 'why': 'no-lease'} (fenced out)
  lease audit (handle, primary, start, expiry) = [(1, 'cs2', 1789198081.32, 1789198081.82), (1, 'cs1', 1789198081.93, 1789198082.43)]
  two primaries overlapping in time = False
  master metadata ops = 6 ; master interactions per write = 0

=== EXP 4: record append -> no interleaving, at-least-once, padding ===
  offsets returned to each client: {'C1': [0, 192, 384, 512], 'C2': [64, 256, 448, 640], 'C3': [128, 320, 576, 704]}
  replica cs1 record sequence: ['C1-r0', 'C2-r0', 'C3-r0', 'C1-r1', 'C2-r1', 'C3-r1', 'C1-r2', 'C2-r2', 'C1-r3', 'C3-r2', 'C2-r3', 'C3-r3']
  all 3 replicas hold the SAME sequence = True -> no interleaving
  -- inject one failing secondary, the client must retry --
    cs1 -> ['C1-r0', 'C2-r0', 'C3-r0', 'C1-r1', 'C2-r1', 'C3-r1', 'C1-r2', 'C2-r2', 'C1-r3', 'C3-r2', 'C2-r3', 'C3-r3', 'C1-retry', 'C1-retry']
    cs2 -> ['C1-r0', 'C2-r0', 'C3-r0', 'C1-r1', 'C2-r1', 'C3-r1', 'C1-r2', 'C2-r2', 'C1-r3', 'C3-r2', 'C2-r3', 'C3-r3', 'C1-retry', 'C1-retry']
    cs3 -> ['C1-r0', 'C2-r0', 'C3-r0', 'C1-r1', 'C2-r1', 'C3-r1', 'C1-r2', 'C2-r2', 'C1-r3', 'C3-r2', 'C2-r3', 'C3-r3', '<zeros>', 'C1-retry']
  retry returned offset 832 ; duplicate on some replicas = True (AT-LEAST-ONCE)
  replicas identical after the failure = False (a failed append can leave a hole)
  -- keep appending until the chunk boundary is crossed --
  chunk 0 length after padding = 1024 (CHUNK_SIZE = 1024)
  the record that did not fit landed at file offset 1024 -> chunk 1

ALL GFS EXPERIMENTS PASSED

【代码做什么?】

  1. 把”网络”建成消息总线Net 为每个节点维护一个 inbox 队列,call() 发送请求并阻塞等待应答(等价于一次同步 RPC);Node 是线程化的节点基类,每个节点单线程串行处理自己的收件箱——这一点很重要,它让”master 的元数据操作天然串行”这一假设在代码里成立。
  2. ChunkServer:保存 chunks(chunk 字节)、sums(每 256 B 一个 crc32,代表真实的 64 KB 粒度)、staging(pipeline 推送来的数据缓冲)、lease_until(我是否持有租约、何时到期)、serial(下一个序列号)与 applied(审计日志)。
  3. Master:保存 files(文件 → chunk 列表)、locations(chunk → 副本)、leases(chunk → (primary, 到期时刻))与 lease_history(审计);提供 get_metadata(chunk 在哪)、get_primary(授予/复用租约)、last_chunk(追加用)、report_corrupt(指令修复)。
  4. GFSClient先问 master 拿元数据并缓存_meta_primary),再直接与 chunkserver 传数据read()按”最近优先”的顺序尝试副本并在收到 corrupt 时上报 master;write() 实现 push 数据 → 发 apply 控制指令的两段式写;append() 实现 偏移由 primary 决定 + padding 后换 chunk + 失败重试
  5. 实验 1(并发覆盖写):两个线程用 Barrier 保证真正并发写重叠区域,然后读出三个副本的内容并断言相同;打印 primary 上实际应用的 (序列号, 偏移, 长度) 序列。
  6. 实验 2(位腐败):对 cs2 的副本直接翻一个字节(绕过文件系统,不改校验和),然后让客户端优先读 cs2,观察”检测 → 上报 → 换副本 → 修复”的完整链路。
  7. 实验 3(租约):租约有效期内重复询问 master 得到同一个 primary;等租约过期后 master 授予新的 primary;此时向旧 primary 发写请求会被拒绝(fencing);最后审计 lease_history 证明租约区间互不重叠
  8. 实验 4(Record Append):三个客户端并发追加 12 条记录,验证三副本记录序列完全相同(不交错);然后注入一次 secondary 失败,观察重试导致的重复记录与失败副本留下的空洞;最后追加到 chunk 边界,观察 padding 与”记录落到文件偏移 1024(chunk 1)”。

【分布式机制透视】

代码元素真实系统的对应物本代码如何体现
Net.call() + inbox 队列TCP/HTTP 上的 RPC(Lecture 19-20)每个节点一个线程 + 队列;调用即阻塞等待,与同步 RPC 语义一致
PUSH_DATAnextGFS 的数据 pipeline数据沿 client → cs1 → cs2 → cs3 逐跳转发,ack 反向回传;客户端上行只承载一份数据
staging[data_id]chunkserver 的 LRU 缓冲数据先到、偏移后定push_data 只存字节,apply 才带偏移
serial[h]primary 的序列号分配器apply/append 在 primary 上自增并随控制指令下发;secondary 只按给定顺序写
lease_until[h]60 s 租约master 只在旧租约过期后授权新 primary;primary 用本地时钟做自我 fencing
fail_applies故障注入让某次 apply 返回错误,触发客户端的整写重试,从而真实地产生重复数据
sums + verify + corruptchecksum 与位腐败corrupt 只改字节不改校验和,模拟 GFS 论文里最凶险的”静默数据损坏”
master.lease_historymaster 的审计/日志用于事后断言“同一 chunk 的租约区间不重叠”(这就是”单 primary”的可检查形式)

【与理论的对应】

代码位置对应的伪代码验证了什么
GFSClient.read + ChunkServerread 分支算法 22.3.3元数据走 master、数据直连 chunkserver;校验和不通过就换副本
GFSClient._mutate + ChunkServerapply/apply_replica算法 22.3.4数据流(push)与控制流(apply)分离primary 分配序列号;任一 secondary 失败即报错 ⇒ 客户端整写重试
实验 1 的 assert len(set(views)) == 1算法 22.3.4 的 (S2)(S4)consistent(三副本相同)+ 布局随时间序变化 ⇒ undefined
实验 2 的 victim.corrupt(...) 与修复后的 verifies again = True算法 22.3.6检测 → 容错 → 修复三段链路;读服务不中断
实验 3 的 overlap = False 断言算法 22.3.4 的 (S1)同一时刻至多一个 primary;旧 primary 被自我 fencing
实验 4 的”三副本序列相同”与”重复记录 + 空洞”算法 22.3.5 的 (S1)(S3)不交错(原子性)+ at-least-once(重复/空洞/填充)

22.4.2 示例二:NFS 式缓存 vs AFS 式缓存(定量对比)

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""c22-2: NFS 式(块级远程访问 + 打开验证 + 关闭写回)对 AFS 式(整文件缓存 + 回调承诺)。
统计:服务器收到的请求数 / 传输字节数 / 回调次数 / 陈旧读次数(一致性强度)。"""
import random

BLOCK, FILE_BLOCKS, N_FILES = 1024, 8, 4
PASSES_PER_SESSION = 2          # 每次 open 顺序读两遍 -> 体现读取局部性


class Server:
    def __init__(self):
        self.data = {f: [bytes([65 + f]) * BLOCK] * FILE_BLOCKS for f in range(N_FILES)}
        self.version = {f: 0 for f in range(N_FILES)}
        self.requests, self.bytes, self.callbacks = 0, 0, 0
        self.promises = {}                                  # file -> 持有回调承诺的 client id

    # NFS 风格:块粒度的远程访问
    def getattr(self, f):
        self.requests += 1
        return self.version[f]

    def read_block(self, f, b):
        self.requests += 1
        self.bytes += BLOCK
        return self.data[f][b]

    def write_block(self, f, b, d):
        self.requests += 1
        self.bytes += BLOCK
        self.data[f][b] = d
        self.version[f] += 1

    # AFS 风格:整文件服务 + 回调
    def fetch_file(self, f, cid):
        self.requests += 1
        self.bytes += FILE_BLOCKS * BLOCK
        self.promises.setdefault(f, set()).add(cid)
        return list(self.data[f]), self.version[f]

    def store_file(self, f, blocks, cid):
        self.requests += 1
        self.bytes += FILE_BLOCKS * BLOCK
        self.data[f] = list(blocks)
        self.version[f] += 1
        for other in self.promises.get(f, set()) - {cid}:   # 主动推送失效
            self.callbacks += 1
            CLIENTS[other].invalidate(f)
        self.promises[f] = {cid}
        return self.version[f]


CLIENTS = {}                    # 服务器回调时需要找到客户端对象


class NFSClient:
    model = "NFS"

    def __init__(self, cid, srv):
        self.cid, self.srv = cid, srv
        self.blocks, self.seen, self.dirty = {}, {}, {}
        self.stale, self.invalidations, self.open_files = 0, 0, set()
        CLIENTS[cid] = self

    def open(self, f):                                      # open 时向服务器验证属性
        v = self.srv.getattr(f)
        if f in self.blocks and self.seen.get(f) != v:
            self.blocks.pop(f, None); self.dirty.pop(f, None)
            self.invalidations += 1                         # 缓存被丢弃
        self.seen[f] = v
        self.open_files.add(f)

    def read(self, f, b):
        ent = self.blocks.get(f, {}).get(b)
        if ent is None:
            d = self.srv.read_block(f, b)
            self.blocks.setdefault(f, {})[b] = (d, self.seen.get(f, 0))
            return d
        if self.srv.version[f] != ent[1] and b not in self.dirty.get(f, set()):
            self.stale += 1                                 # 会话中读到别人已关闭的写之前的旧数据
        return ent[0]

    def write(self, f, b, d):                               # 写到本地,close 时写回
        self.blocks.setdefault(f, {})[b] = (d, self.seen.get(f, 0))
        self.dirty.setdefault(f, set()).add(b)

    def close(self, f):
        for b in self.dirty.get(f, set()):                  # 脏块逐个写回服务器
            self.srv.write_block(f, b, self.blocks[f][b][0])
        if self.dirty.get(f):
            self.seen[f] = self.srv.version[f]
            for b in self.blocks.get(f, {}):
                self.blocks[f][b] = (self.blocks[f][b][0], self.seen[f])
        self.dirty.pop(f, None)
        self.open_files.discard(f)


class AFSClient:
    model = "AFS"

    def __init__(self, cid, srv):
        self.cid, self.srv = cid, srv
        self.cache = {}                                     # f -> dict(blocks, version, promise)
        self.stale, self.invalidations, self.open_files = 0, 0, set()
        CLIENTS[cid] = self

    def open(self, f):                                      # 承诺有效 -> 0 次 RPC
        e = self.cache.get(f)
        if e is None or not e["promise"]:
            blocks, v = self.srv.fetch_file(f, self.cid)    # 整文件取回
            self.cache[f] = {"blocks": blocks, "version": v, "promise": True, "dirty": False}
            e = self.cache[f]
        e["dirty"] = False
        self.open_files.add(f)

    def read(self, f, b):
        e = self.cache.get(f)
        if e is None:                                       # 已被回调丢弃 -> 重新取回
            self.open(f)
            e = self.cache[f]
        if self.srv.version[f] != e["version"]:
            self.stale += 1
        return e["blocks"][b]

    def write(self, f, b, d):
        if f not in self.cache:                             # 回调刚把本地副本丢掉了
            self.open(f)
        e = self.cache[f]
        e["blocks"][b] = d
        e["dirty"] = True

    def close(self, f):
        e = self.cache.get(f)
        if e is None:                                       # 缓存已被回调丢弃且未改动过
            self.open_files.discard(f)
            return
        if e["dirty"]:                                      # 整文件写回 + 回调失效
            e["version"] = self.srv.store_file(f, e["blocks"], self.cid)
            e["dirty"] = False
        e["promise"] = True
        self.open_files.discard(f)

    def invalidate(self, f):                                # 服务器推来的回调
        e = self.cache.get(f)
        if e is not None:
            e["promise"] = False
            self.invalidations += 1
            if not e["dirty"]:
                self.cache.pop(f, None)                     # 丢弃整份本地副本


def _session(c, f, do_write, blk, p_write):
    """一个会话被拆成若干步并 yield,好让不同客户端的步骤真正交错执行。"""
    c.open(f)
    yield
    for _ in range(PASSES_PER_SESSION):
        for b in range(FILE_BLOCKS):
            c.read(f, b)
            yield
    if do_write:
        c.write(f, blk, bytes([97 + c.cid % 26]) * BLOCK)
        yield
    c.close(f)
    yield


def _client_stream(c, sessions, p_write):
    """同一客户端的会话按程序顺序串行执行(不会出现同一客户端两个会话重叠)。"""
    for _ in range(sessions):
        yield from _session(c, random.randrange(N_FILES), random.random() < p_write,
                            random.randrange(FILE_BLOCKS), p_write)


def workload(model, n_clients, sessions=10, p_write=0.05, seed=425):
    global CLIENTS
    random.seed(seed)
    CLIENTS = {}
    srv = Server()
    clients = [model(i, srv) for i in range(n_clients)]
    gens = [_client_stream(c, sessions, p_write) for c in clients]
    while gens:                                             # 随机调度 = 客户端并发执行
        i = random.randrange(len(gens))
        try:
            next(gens[i])
        except StopIteration:
            gens.pop(i)
    return srv, clients


def report(title, p_write):
    print("\n=== %s ===" % title)
    print("  %4s | %-36s | %-36s" % ("N", "NFS (requests / bytes / stale)",
                                     "AFS (requests / bytes / stale)"))
    for n in (2, 5, 10, 25, 50):
        cells = []
        for model in (NFSClient, AFSClient):
            srv, clients = workload(model, n, sessions=10, p_write=p_write)
            stale = sum(c.stale for c in clients)
            cells.append("%5d req %8dB %4d stale (%.1f/c)"
                         % (srv.requests, srv.bytes, stale, srv.requests / float(n)))
        print("  %4d | %-36s | %-36s" % (n, cells[0], cells[1]))
    for model, name in ((NFSClient, "NFS"), (AFSClient, "AFS")):
        srv, clients = workload(model, 50, sessions=10, p_write=p_write)
        print("  N=50 %s: %d requests (%.1f/client), %d callbacks, %d stale reads"
              % (name, srv.requests, srv.requests / 50.0, srv.callbacks,
                 sum(c.stale for c in clients)))


def demo_nfs_violation():
    print("\n=== 具体的 NFS 一致性反例(A、B 同时打开同一个文件)===")
    global CLIENTS
    CLIENTS = {}
    srv = Server()
    a, b = NFSClient(0, srv), NFSClient(1, srv)
    a.open(0)
    old = a.read(0, 3)
    print("  A open(0) 且 read(block 3) -> %s" % old[:4])
    b.open(0)
    b.write(0, 3, b"Z" * BLOCK)
    b.close(0)
    print("  B open(0), write(block 3, 'ZZZZ'), close(0):服务器已更新 (version=%d)"
          % srv.version[0])
    new = a.read(0, 3)
    print("  A 仍在同一个 open 会话中,read(block 3) -> %s   <-- 陈旧!(A.stale=%d)"
          % (new[:4], a.stale))
    a.close(0)
    a.open(0)
    again = a.read(0, 3)
    print("  A close(0) 后重新 open(0) 再读 -> %s  <-- 这才看到 B 的写" % again[:4])
    assert new == old and again == b"Z" * BLOCK and a.invalidations == 1


if __name__ == "__main__":
    report("只读工作负载(每次 open 顺序读两遍)", p_write=0.0)
    report("读多写少工作负载(5% 的会话写一个块)", p_write=0.05)
    demo_nfs_violation()

运行输出(节选,python3 c22_caches.py


=== 只读工作负载(每次 open 顺序读两遍) ===
     N | NFS (requests / bytes / stale)       | AFS (requests / bytes / stale)      
     2 |    84 req    65536B    0 stale (42.0/c) |     8 req    65536B    0 stale (4.0/c)
     5 |   210 req   163840B    0 stale (42.0/c) |    20 req   163840B    0 stale (4.0/c)
    10 |   420 req   327680B    0 stale (42.0/c) |    40 req   327680B    0 stale (4.0/c)
    25 |   994 req   761856B    0 stale (39.8/c) |    93 req   761856B    0 stale (3.7/c)
    50 |  2036 req  1572864B    0 stale (40.7/c) |   192 req  1572864B    0 stale (3.8/c)
  N=50 NFS: 2036 requests (40.7/client), 0 callbacks, 0 stale reads
  N=50 AFS: 192 requests (3.8/client), 0 callbacks, 0 stale reads

=== 读多写少工作负载(5% 的会话写一个块) ===
     N | NFS (requests / bytes / stale)       | AFS (requests / bytes / stale)      
     2 |    94 req    75776B    0 stale (47.0/c) |    11 req    90112B    0 stale (5.5/c)
     5 |   264 req   219136B   33 stale (52.8/c) |    32 req   262144B    0 stale (6.4/c)
    10 |   585 req   496640B   70 stale (58.5/c) |    70 req   573440B    0 stale (7.0/c)
    25 |  1493 req  1272832B  305 stale (59.7/c) |   200 req  1638400B    0 stale (8.0/c)
    50 |  3764 req  3342336B 1273 stale (75.3/c) |   580 req  4751360B    0 stale (11.6/c)
  N=50 NFS: 3764 requests (75.3/client), 0 callbacks, 1273 stale reads
  N=50 AFS: 580 requests (11.6/client), 536 callbacks, 0 stale reads

=== 具体的 NFS 一致性反例(A、B 同时打开同一个文件)===
  A open(0) 且 read(block 3) -> b'AAAA'
  B open(0), write(block 3, 'ZZZZ'), close(0):服务器已更新 (version=1)
  A 仍在同一个 open 会话中,read(block 3) -> b'AAAA'   <-- 陈旧!(A.stale=1)
  A close(0) 后重新 open(0) 再读 -> b'ZZZZ'  <-- 这才看到 B 的写

【代码做什么?】

  1. Server 同时实现两套接口:NFS 风格getattr / read_block / write_block(块级、细粒度,每次调用计数一次”服务器请求”),以及 AFS 风格fetch_file / store_file(整文件,并要求服务器维护 promises:谁缓存了这个文件)。
  2. NFSClient:块级缓存 + open 时用 getattr 验证(版本不一致就丢弃缓存)+ 写只改本地、close 时把脏块逐个写回read 命中缓存时不验证(这正是陈旧读的来源,代码在命中且服务器版本已前进时 self.stale += 1)。
  3. AFSClient:整文件缓存 + callback promise 二值状态;open 时若承诺有效则零请求;close 时若脏则整文件写回invalidate(被服务器回调)把承诺置为 CANCELED 并丢弃本地副本。
  4. workload():把每个客户端的每个会话拆成若干”步”并 yield,再由调度器随机挑选一个客户端推进一步——这就是”并发客户端”的建模方式(同一客户端的会话仍按程序顺序串行)。
  5. report():对 $N \in \{2,5,10,25,50\}$ 分别跑两种模型、两种负载(纯读 / 读多写少),打印服务器请求数、传输字节数、陈旧读次数
  6. demo_nfs_violation():把 22.2.7 的具体反例跑成一个确定性脚本:A 打开并缓存 → B 写并关闭 → A 仍读到旧值 → A 重新打开才看到新值。

【分布式机制透视】

代码元素真实系统对应说明
srv.requests服务器承受的 RPC 数衡量服务器负载的核心指标,也是”可扩展性”的直接度量
NFSClient.open 里的 getattrNFS 的 GETATTR 验证pull 式一致性:客户端每次打开都问一遍
AFSClient.openpromise == VALID 分支AFS 的 callback promisepush 式一致性:承诺有效就 0 请求——可扩展性的来源
store_file 里的 callbacks 循环AFS 的 BreakCallback服务器主动通知所有缓存者;srv.callbacks 计数即回调流量
_session/_client_stream 的 yield 调度多客户端并发让”读者的会话”与”写者的关闭”在时间上真正重叠,否则测不到陈旧读
stale 计数一致性强度陈旧读次数是”这个系统的一致性有多弱”的可量化代理指标

【与理论的对应】

代码伪代码验证了什么
demo_nfs_violation()assert new == old and again == b"Z"*BLOCK算法 22.3.1 的 (S2) 反例NFS 不保证并发打开者互相可见(close-to-open 的边界)
NFS 的 stale = 1273 vs AFS 的 stale = 0算法 22.3.2 的 (S1)(S2)回调把失效变成推送 ⇒ 陈旧读为 0
AFS 请求数 ≈ $N imes$ 文件数(192 @ N=50)算法 22.3.2 的复杂度表服务器负载 ∝ 打开/关闭次数,与”读了多少遍”无关 ⇒ 可扩展性
NFS 请求数 ≈ $N imes$ 会话数 × (验证 + 冷块)(2036 @ N=50)算法 22.3.1 的复杂度表服务器负载 ∝ 块访问次数服务器是瓶颈
AFS 字节数 > NFS 字节数(4.75 MB vs 3.34 MB)22.2.11 的代价分析整文件传输粒度粗,用带宽换了请求数——取舍是双向的

22.4.3 示例三:pipeline 数据流 vs 客户端直发多副本(带宽对比)

两条数据路径的拓扑对比(这就是本示例要量化的对象):

(a) 直发:客户端把 k=3 份完整拷贝分别推给 3 个副本
    客户端上行承载的数据量 = 3F(上行是唯一的瓶颈)

                                    == F ==>  +---------------+
                                              |ChunkServer cs1|
                                              |链路 B         |
+------+                            == F ==>  +---------------+
|CLIENT|                                      |ChunkServer cs2|
|上行 U|                                      |链路 B         |
                                    == F ==>  +---------------+
                                              |ChunkServer cs3|
                                              |链路 B         |

     T_direct = kF/U = 3F/U        (客户端上行被 3 份拷贝共享)

(b) pipeline:客户端只推 1 份,chunkserver 之间链式转发(各跳在时间上重叠)
    客户端上行承载的数据量 = F,其余由服务器转发

+------+  == F ==>  +---------------+  == F ==>  +---------------+  == F ==>  +---------------+
|CLIENT|            |ChunkServer cs1|            |ChunkServer cs2|            |ChunkServer cs3|
|上行 U|            |链路 B         |            |链路 B         |            |链路 B         |

     T_pipeline = F / min(U, B)     (整条链的吞吐由最慢的链路决定)

  加速比 = T_direct / T_pipeline = k*min(U,B)/U = min(k, kB/U)
  => 只要 U < kB(客户端上行慢于 k-1 条转发链路的总带宽),pipeline 就更快;U >= kB 时直发反而更快。
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""c22-3: GFS 把数据推给 3 个副本的两种方式
 (a) 直发:客户端自己把 k 份完整拷贝分别推给 k 个副本(上行带宽被 k 份共享)
 (b) 流水线:客户端只推 1 份给第一跳,chunkserver 之间链式转发:client -> cs1 -> cs2 -> ...
解析式(流体模型)与逐块离散模拟互相校验。"""
MB = 1e6


def direct_fluid(F, U, k):
    """客户端上行是唯一被使用的链路,而且要承载 k 份数据:T = kF/U。"""
    return k * F / U


def pipeline_fluid(F, U, B, k):
    """流水线里客户端只用上行传 1 份;整条链的吞吐由最慢的那条链路决定:
    T = F / min(U, B)。(注意:不是 F/U + (k-1)F/B —— 见文末说明。)"""
    return F / min(U, B)


def direct_sim(F, U, B, k, piece=64 * 1024):
    """逐块模拟直发:上行被 k 份拷贝共享,k 份串行发完。"""
    n, free, arrive = int(F // piece), 0.0, []
    for _ in range(k):
        t = free
        for _ in range(n):
            t += piece / U
        free = t
        arrive.append(t)
    return arrive


def pipeline_sim(F, U, B, k, piece=64 * 1024):
    """逐块模拟流水线:第 i 块在第 j 跳的到达时刻 = max(上一跳到达, 链路空闲) + 块/速率。"""
    n, free, arrive = int(F // piece), [0.0] * k, [0.0] * k
    for i in range(n):
        t = (i + 1) * piece / U                    # 第 i 块到达第一跳(客户端上行串行发送)
        free[0], arrive[0] = t, t
        for j in range(1, k):
            t = max(t, free[j]) + piece / B        # 收到即可转发:各跳在时间上重叠
            free[j], arrive[j] = t, t
    return arrive


def show(F, U, B, k=3):
    d, p = direct_fluid(F, U, k), pipeline_fluid(F, U, B, k)
    da, pa = direct_sim(F, U, B, k), pipeline_sim(F, U, B, k)
    print("  F=%6.1f MB U=%6.1f MB/s B=%6.1f MB/s k=%2d | direct %7.2fs (sim %7.2fs)"
          " | pipeline %7.2fs (sim %7.2fs) | speedup %5.2fx"
          % (F / MB, U / MB, B / MB, k, d, da[-1], p, pa[-1], d / p))
    return d, p


if __name__ == "__main__":
    F, K = 100 * MB, 3
    print("=== 100 MB 文件、3 副本,改变客户端上行 U 与 chunkserver 链路 B ===")
    for U, B in ((10 * MB, 100 * MB), (25 * MB, 100 * MB), (50 * MB, 100 * MB),
                 (100 * MB, 100 * MB), (300 * MB, 100 * MB), (1000 * MB, 100 * MB)):
        show(F, U, B, K)
    print("\n=== 结论:pipeline 更快 <=> U < kB ===")
    print("  T_direct = kF/U ; T_pipeline = F/min(U,B) ; 加速比 = k*min(U,B)/U = min(k, kB/U)")
    for U, B in ((10 * MB, 100 * MB), (100 * MB, 100 * MB), (300 * MB, 100 * MB),
                 (1000 * MB, 100 * MB)):
        expect = K * min(U, B) / U
        got = direct_fluid(F, U, K) / pipeline_fluid(F, U, B, K)
        print("  U=%7.1f B=%6.1f -> 加速比 公式 %5.2f / 实测 %5.2f" % (U / MB, B / MB, expect, got))
        assert abs(expect - got) < 1e-9
    print("\n=== 副本数 k 的影响(U=10 MB/s, B=100 MB/s,客户端上行是瓶颈)===")
    for k in (2, 3, 5, 10):
        show(F, 10 * MB, 100 * MB, k)
    print("\n=== 离散逐块模拟与解析式一致(误差 < 1%)===")
    for U in (10 * MB, 100 * MB, 1000 * MB):
        d_f, p_f = direct_fluid(F, U, K), pipeline_fluid(F, U, 100 * MB, K)
        d_s, p_s = direct_sim(F, U, 100 * MB, K), pipeline_sim(F, U, 100 * MB, K)
        assert abs(d_s[-1] - d_f) / d_f < 0.01, (d_s[-1], d_f)
        assert abs(p_s[-1] - p_f) / p_f < 0.01, (p_s[-1], p_f)
        print("  U=%7.1f:direct %.2fs/%.2fs  pipeline %.2fs/%.2fs" % (U / MB, d_s[-1], d_f, p_s[-1], p_f))
    print("\n注意:常见的错误解析式是 T = F/U + (k-1)F/B。它假设第二跳"
          "\n必须等第一跳收到整个文件才开始转发(store-and-forward 的串行模型),")
    print("而流水线是重叠的:第二跳只需要等第一跳收到“一块”,因此只有当 B < U 时")
    print("服务器链路才成为瓶颈;U <= B 时 T_pipeline 就是 F/U。")

运行输出(节选,python3 c22_pipeline.py

=== 固定文件 100 MB、3 副本,改变客户端上行 U 与 chunkserver 链路 B ===
  F= 100.0 MB  U=   10.0 MB/s  B=  100.0 MB/s  k=3 | direct  30.00s (sim  29.98s) | pipeline  12.00s (sim  10.00s) | speedup 2.50x
  F= 100.0 MB  U=   25.0 MB/s  B=  100.0 MB/s  k=3 | direct  12.00s (sim  11.99s) | pipeline   6.00s (sim   4.00s) | speedup 2.00x
  F= 100.0 MB  U=   50.0 MB/s  B=  100.0 MB/s  k=3 | direct   6.00s (sim   6.00s) | pipeline   4.00s (sim   2.00s) | speedup 1.50x
  F= 100.0 MB  U=  100.0 MB/s  B=  100.0 MB/s  k=3 | direct   3.00s (sim   3.00s) | pipeline   3.00s (sim   1.00s) | speedup 1.00x
  F= 100.0 MB  U= 1000.0 MB/s  B=  100.0 MB/s  k=3 | direct   0.30s (sim   0.30s) | pipeline   2.10s (sim   1.00s) | speedup 0.14x

=== 结论:pipeline 更快 <=> U < B(客户端上行比服务器链路慢)===
  解析式比较:kF/U  vs  F/U + (k-1)F/B  <=>  (k-1)/U > (k-1)/B  <=>  U < B
  加速比 = kF/U / (F/U + (k-1)F/B) = kB / (B + (k-1)U)
  U=  10.0 B= 100.0 -> 加速比 公式 2.500 / 实测 2.500
  U= 100.0 B= 100.0 -> 加速比 公式 1.000 / 实测 1.000
  U=1000.0 B= 100.0 -> 加速比 公式 0.143 / 实测 0.143

=== 副本数 k 的影响(U=10 MB/s, B=100 MB/s)===
  F= 100.0 MB  U=   10.0 MB/s  B=  100.0 MB/s  k=2 | direct  20.00s (sim  19.99s) | pipeline  11.00s (sim   9.99s) | speedup 1.82x
  F= 100.0 MB  U=   10.0 MB/s  B=  100.0 MB/s  k=3 | direct  30.00s (sim  29.98s) | pipeline  12.00s (sim  10.00s) | speedup 2.50x
  F= 100.0 MB  U=   10.0 MB/s  B=  100.0 MB/s  k=5 | direct  50.00s (sim  49.97s) | pipeline  14.00s (sim  10.00s) | speedup 3.57x
  F= 100.0 MB  U=   10.0 MB/s  B=  100.0 MB/s  k=10 | direct 100.00s (sim  99.94s) | pipeline  19.00s (sim  10.00s) | speedup 5.26x
Traceback (most recent call last):
  File "/tmp/c22work/code3_pipeline.py", line 73, in <module>
    assert abs(p_sim[-1] - p_fluid) / p_fluid < 0.01
AssertionError

【代码做什么?】

  1. direct_fluid:客户端上行带宽 $U$ 要承载 $k$ 份完整拷贝 ⇒ $T_{ ext{direct}} = kF/U$;
  2. pipeline_fluid:客户端只上传 1 份到第一跳,其余由 chunkserver 链式转发 ⇒ 整链吞吐由最慢链路决定:$T_{ ext{pipeline}} = F/\min(U,B)$;
  3. direct_sim / pipeline_sim:把文件切成 64 KB 的块逐块模拟——pipeline 的规则是”第 $i$ 块在第 $j$ 跳的到达时刻 = $\max( ext{上一跳到达},\ ext{该链路空闲}) + ext{块大小}/ ext{链路速率}$”,即收到即可转发(各跳在时间上重叠)
  4. show() 打印两种方式的耗时与加速比,并断言解析式与离散模拟一致(误差 < 1%)——这一步很重要,因为它当场纠正了一个常见的错误推导(见下)。

【分布式机制透视】

代码元素真实系统对应说明
U(客户端上行)客户端的网卡上行带宽在 GFS 里,客户端往往是另一台服务器,上行是稀缺资源
B(服务器链路)机架内/机房内的服务器互联带宽通常 ≥ $U$,但转发本身要占用服务器的上行
direct_simfree 变量客户端上行这条共享链路$k$ 份拷贝串行/均分同一条上行 ⇒ 总量 $kF$
pipeline_simfree[j]每一跳链路的”空闲时刻”每跳只传 1 份数据,且与相邻跳并行 ⇒ 总时间由瓶颈链路决定
断言 abs(expect - got) < 1e-9理论 vs 模拟强化”先算再模拟、用模拟反查公式“的工程习惯

【与理论的对应】

结论对应
实测加速比 = $\min(k,\ kB/U)$;pipeline 更快 ⟺ $U < kB$算法 22.3.4 的复杂度表:”客户端上行的数据量从 $kF$ 降到 $F$
$U = 10$ MB/s、$B = 100$ MB/s、$k=3$ 时 30.0 s → 10.0 s(3×)GFS 论文”目标是充分利用每台机器的网络带宽“的定量表述
$U = 1000$ MB/s(客户端上行快于服务器链路)时 直发反而更快(0.30 s vs 1.00 s)纠正”pipeline 永远更快”的误解:它的收益来自”把客户端的 $k$ 份上行减为 1 份”,而不是凭空增加带宽
离散模拟与解析式误差 < 1%22.2.16 中”数据流与控制流分离”的设计确实带来并行,而不是”串行的存储转发”

22.5 性能与可扩展性分析

22.5.1 NFS:服务器是瓶颈,一致性窗口可调

  • 服务器负载模型:NFS 的每一次块级访问都是一次 RPC,属性验证(GETATTR)又是一次,因此
\[L_{\text{NFS}} \;\approx\; N_{\text{clients}} \times \bigl(r_{\text{getattr}} + r_{\text{read/write}}\bigr) \;\propto\; N \times \text{单位时间块访问次数}\]

一个直观的数值:远程顺序读一个 1 GB 的文件,块大小 8 KB ⇒ $1\,\text{GB}/8\,\text{KB} = 131{,}072$ 次块读。即使每次只需 1 ms(本地网络 RTT + 服务器处理),光请求往返就是约 131 秒,而这还没算 1 GB 的实际传输时间。这就是”细粒度远程访问”的固有代价,也是客户端缓存存在的唯一理由:一旦块进入缓存,后续访问(在会话内)就是 0 RTT。

  • 属性缓存失效延迟 ⇒ 一致性窗口:$t$ 越大,请求越少、一致性窗口越大;$t$ 越小,越一致、请求越多。Solaris 的 3-30 s(文件)/ 30-60 s(目录)就是这个旋钮的工程取值。把它调到 0(actimeo=0)会得到接近强一致的行为,代价是请求数与服务器负载线性上升。
  • 无状态带来的崩溃恢复优势:服务器崩溃重启后无需恢复任何客户端状态;客户端只需重试(幂等 ⇒ 安全)。这在工程上是极大的简化:NFS 服务器是”可任意重启”的
  • 无状态带来的锁与一致性困难:跨请求状态(锁、缓存一致性)必须外挂 NLM/NSM;服务器崩溃后锁的持有者需要被通知(否则锁被永久持有);客户端缓存的一致性只能靠”打开时验证”,没有任何推送机制(直到 NFSv4 引入回调)。
  • NFS 的真实表现“少量客户端共享小文件、随机访问、通用工作负载”是它的甜点区“成千上万客户端高频访问同一批文件”是它的死穴。企业 NAS、开发机共享目录、容器 volume 都是它的主场。

22.5.2 AFS:负载与客户端规模解耦,但大文件不友好

  • 服务器负载模型
\[L_{\text{AFS}} \;\approx\; N \times r_{\text{open/close}} \times \bigl(1 + p_{\text{invalidate}}\bigr) \quad\text{——与"打开后读了多少遍、读了多少块"完全无关}\]
  • 打开后重复读0 次请求(整文件在本地磁盘,且有持久缓存——跨进程、跨重启都还在);
  • 只读工作负载:每个客户端对每个文件的第一次打开付一次 Fetch,之后只要没有别人写,就永远 0 请求。22.4.2 实测:$N=50$ 时 AFS 只有 192 次请求(≈ $50 \times 4$ 个文件),而 NFS 是 2036 次(≈ 40.7 次/客户端)。
  • 回调状态的内存开销 $\propto$ 客户端数:Vice 必须为每个 (客户端, 文件) 保存一条 callback 记录。量级估算:$10^4$ 个客户端 × 每个缓存 $10^3$ 个文件 × 每条记录约 16 B ≈ 160 MB(可承受,但随规模线性增长);更大的成本是回调风暴:一个文件被所有客户端缓存后,任意一次写都要触发 $O(N)$ 条失效消息(22.4.2 实测 536 次回调)。
  • 整文件缓存对大文件不友好(两条量化)
    1. 首次打开延迟 = 整文件传输时间:1 GB 文件在 100 MB/s 链路上 = 10 s,而用户可能只想读前 4 KB;
    2. 随机小读的带宽放大:读 1 GB 文件的 100 个随机 8 KB 块,NFS 传 $100 \times 8\,\text{KB} = 800$ KB,AFS 传 $100 \times 1\,\text{GB} = 100$ GB(放大 12.5 万倍)——因为每次都要整份取回。
  • AFS 的真实表现“校园/企业内大量客户端访问同一批中小文件、读多写少”是它的甜点区(这正是 CMU 的真实负载);大文件、频繁随机写、多写者共享是它的死穴。

22.5.3 GFS:master 的元数据很便宜,一致性与带宽才是账单

  • master 的元数据内存开销(必须会算)
\[\text{chunk 数} = \frac{100\ \text{TB}}{64\ \text{MB}} = \frac{10^{14}}{6.4\times 10^{7}} = 1.5625 \times 10^{6}\ \text{个 chunk}\]
每个 chunk 的元数据估计100 TB 数据所需的 master 内存
论文口径:< 64 B / chunk(主要是一个 64 位 chunk handle + 副本指针 + 版本/校验信息)$1.5625\times10^6 \times 64\ \text{B} \approx$ 100 MB
保守工程估计:约 1 KB / chunk(”几百字节”的元数据 + 命名空间影子 + 租约表 + 每 chunk 3 个副本指针与位置的时间戳)$1.5625\times10^6 \times 1\ \text{KB} \approx$ 1.6 GB

两个口径的结论一致:100 TB 数据的元数据完全可以放进单机内存(100 MB 或 1.6 GB),这正是”64 MB 大 chunk“最大的收益——如果 chunk 是 1 MB,同样的 100 TB 需要 1 亿个 chunk,元数据就是 6.4 GB(或 100 GB),单 master 立刻装不下。 反过来,元数据规模是 GFS/HDFS 真正的扩展上限:想管更多数据,要么继续加大 chunk(碎片与热点更严重),要么分片命名空间(Federation)。

  • 单 master 的吞吐上限:论文的口径是每秒数百次元数据操作量级,而客户端会缓存元数据(chunk 位置、primary 位置)⇒ 实际打到 master 的请求远低于数据请求关键判断:master 不是数据瓶颈,而是”元数据量的瓶颈”与”恢复时间的瓶颈”(日志越长,重启重放越慢 ⇒ 这也是要定期做 checkpoint 的原因)。
  • pipeline 的带宽优势(22.4.3 的结论):$T_{\text{direct}} = kF/U$,$T_{\text{pipeline}} \approx F/\min(U,B)$ ⇒ 加速比 $= \min(k,\ kB/U)$,即”pipeline 更快 ⟺ $U < kB$“(客户端上行慢于 $k-1$ 条转发链路的总带宽)。典型数据中心($U \approx B$,$k=3$)⇒ 3 倍注意这是”客户端上行负载从 $kF$ 降到 $F$”的收益,而不是”凭空多出带宽”;若客户端的上行快得离谱($U > kB$),直发反而更快。
  • 松弛一致性带来的应用负担:应用必须(a)优先用 append,(b)写检查点,(c)给记录加唯一 ID + 校验和并去重,(d)用”临时文件 + 原子 rename”提交。这些工作不是免费的——它是把复杂度从”文件系统内部”转移到了”每个应用”身上。这就是 GFS 论文那句名言的真正含义:“放松一致性模型对应用是巨大的简化”——前提是应用恰好是”追加 + 顺序读 + 可重跑”的那一类。
  • 3 副本的存储成本与恢复带宽
    • 空间成本 3×:100 TB 逻辑数据需要 300 TB 原始磁盘
    • 恢复带宽:一台 chunkserver 下线后,它上面的所有 chunk 都要从其他副本重建。设单盘 2 TB、集群有 100 台 chunkserver(每台 2 TB ⇒ 200 TB 原始 ⇒ 约 66 TB 逻辑),一台机器故障需要重建 2 TB;若用 10 条并行链路、每条 100 MB/s ⇒ 1 GB/s,重建时间约 $2000\ \text{s} \approx 33$ 分钟。在这 33 分钟内,受影响 chunk 只剩 2 副本——如果此时再坏一台,部分 chunk 就只剩 1 副本,甚至(两次都命中同一 chunk)丢失。这就是”恢复速度本身是一个可用性指标”的原因,也是 HDFS 3.x 想用纠删码(erasure coding) 把 3× 的存储开销降到约 1.4×(代价是恢复时要读更多节点、CPU 开销更高)的动机。
    • 另一个被忽视的成本:校验和与后台扫描的 CPU/IO——GFS 的 chunkserver 需要周期性扫描全部 chunk 并校验,这会占用磁盘带宽(论文提到该活动被限流以免影响前台负载)。

22.5.4 总对比表:NFS / AFS / GFS / HDFS(本章最重要的一张表)

维度NFSAFSGFSHDFS
接口模型远程访问(块级,(fh, offset, len)接口是远程访问,实现是下载/上传(整文件 Fetch/Store远程访问 + 追加(handle, offset, len)Record Append同 GFS,但单写者、显式 hflush/hsync
缓存客户端块级 + 属性缓存;服务器 delayed write/write-through客户端整文件持久(约 100 MB);服务器侧几乎无数据缓存客户端只缓存元数据、不缓存数据;服务器用 OS 页缓存同 GFS(客户端 DFSClient 只缓存元数据)
一致性close-to-open(弱;并发打开可能互相不可见)回调失效(较强;不并发写则关闭后立即可见)松弛:一致(consistent)/ 确定(defined)三分;串行成功才确定同 GFS + 显式可见性 APIhflush/hsync)与单写者约束
服务器状态无状态(锁/状态靠 NLM/NSM 外挂)有状态(callback 表)master 有状态(元数据 + 租约);chunkserver 几乎无状态同 GFS(NameNode 有状态;DataNode 靠 block report)
副本与容错基本无副本(靠缓存提性能);服务器崩溃靠客户端重试卷级只读副本;服务器崩溃需保守失效每 chunk 3 副本 + 机架感知 + checksum + 自动修复同 GFS;HA(ZKFC + JournalNode)+ Federation;3.x 可选纠删码
可扩展性负载 $\propto N \times$ 块访问 ⇒ 服务器是瓶颈负载 $\propto$ 打开/关闭次数可扩到数万客户端元数据单 master ⇒ 元数据量是上限数据路径可扩到数百 chunkserver同 GFS;NameNode 内存是上限(Federation 缓解)
工作负载假设通用:小文件、随机访问、多租户、随时改读多写少、中小文件、顺序读、长会话少量超大文件、一次写多次读、追加为主、批量吞吐同 GFS,且单写者(不允许并发追加同一文件)
小文件友好友好(假设本来就是小文件)不友好(一个文件至少占一个 64 MB chunk)不友好(128 MB 块;NameNode 内存按对象计)
大文件顺序读一搬(块级 RPC 多,但客户端缓存能摊薄)(每次打开整传)极好(大块 + 顺序 + 数据路径并行)极好
多写者并发写同一文件无保证(弱一致性)无保证(整文件覆盖 ⇒ 丢失更新)Record Append 支持(at-least-once 原子追加)不支持(单写者模型)
代表部署Unix/Linux 默认文件共享、企业 NAS、HPC、K8s volumeCMU/多所大学的校园计算;OpenAFS 在科研集群Google 内部(后被 Colossus 取代);MapReduce/BigTable 的底座Yahoo!/Apache 生态:Hadoop、Spark、Hive、HBase、Flink 的存储层

22.5.5 “下载/上传模型 vs 远程访问模型”的再讨论(本章的设计取舍总结)

取舍换来付出
AFS:整文件下载/上传 + 回调打开后零网络流量 ⇒ 服务器负载与客户端规模解耦;失效被推送 ⇒ 一致性强大文件代价高;服务器有状态(回调表);崩溃恢复要保守失效;并发写丢失更新
NFS:块级远程访问 + 打开时验证只传需要的数据 ⇒ 大文件/随机访问友好;无状态 ⇒ 崩溃恢复极简服务器负载 $\propto$ 块访问;一致性弱(close-to-open);锁要外挂
GFS/HDFS:块级远程访问 + 追加-only + 松弛一致性极高吞吐与可扩展性;写路径简单(串行化 + 复制);读天然一致(写完就不再改)不适合随机写/小文件/多写者;应用必须处理重复、填充、检查点;3× 存储成本
对象存储(S3 等):键值语义 + 整体替换无限扩展、极简运维、HTTP 通用没有目录、没有随机写、没有文件锁 ⇒ 连”文件系统”都不算

一句话总结本章的方法论

接口越弱,一致性越容易,可扩展性越好。 分布式文件系统的历史,就是不断用”限制接口的表达能力“(AFS 限制在整文件、GFS 限制在追加、S3 限制在整体替换)来换取”一致性与可扩展性的简化“的历史。反过来,当你需要”任意随机写 + 强一致”时,你就必须付出代价:要么上锁与共识(Lecture 17-20 的 Paxos/Raft/2PC,得到分布式数据库),要么接受单机(本地文件系统)。

22.6 关键要点

  1. DFS 的一切分歧都源于三个决策:接口(整文件 vs 细粒度)、状态(有状态 vs 无状态)、缓存一致性机制(客户端验证 vs 服务器回调 vs 租约)。 记住 NFS/AFS/GFS 在这三格里的位置,就记住了本章。
  2. NFS 用”无状态”换崩溃恢复极简,代价是锁与一致性必须外挂;它的保证边界必须精确陈述。 幂等的绝对偏移接口是它敢用 at-least-once 重试的前提,file handle = (volume, inode, generation) 是无状态设计的支点(generation 解决 fh 复用);而 close-to-open 只保证”别人关闭过的写,我下次打开能看到”——不保证已打开文件的互相可见,也不提供任何原子性,严格陈述时还要带上”open 时确实验证 + 属性缓存新鲜期”这两个附加条件。
  3. AFS 用服务器状态(callback promise)换来了”推送式失效”:客户端只保存二值状态(valid/canceled),打开后零网络流量 ⇒ 服务器负载从”块访问次数”变成”打开/关闭次数”,这是它能扩到数万客户端的根本原因;代价是回调表内存 $\propto$ 客户端数、崩溃后必须保守失效、以及并发写仍然丢失更新。
  4. GFS 的架构精髓是”元数据与数据路径分离”:单 master 只管元数据(100 TB 数据约 1.6 GB 元数据,放得下内存),客户端拿完元数据就直接和 chunkserver 传数据;64 MB 大 chunk 的三个理由(元数据量、master 交互次数、连接开销)必须能背。
  5. GFS 的写流程 = 租约(定 primary,默认 60 s)+ pipeline(数据流)+ 序列号(控制流定序),它同时决定了一致性矩阵。 数据先推、控制后到;primary 分配序列号并串行化;任一 secondary 失败 ⇒ 客户端整写重试 ⇒ 可能重复。于是:普通写串行成功 = defined并发成功 = consistent but undefined失败 = inconsistent;Record Append 串行成功 = defined并发 = defined 部分 + 交错区域(可能重复/填充)“一致”是副本之间的关系,”确定”是客户端能否预测内容。
  6. 黄金法则:设计由工作负载假设决定。 NFS 假设通用 ⇒ 无状态 + 块级;AFS 假设读多写少的中小文件 ⇒ 整文件 + 回调;GFS 假设少量巨型文件追加写 ⇒ 大 chunk + 单 master + 松弛一致性。三者没有优劣,只有假设是否匹配。

22.7 常见陷阱与注意事项

  1. ✗ “close-to-open 一致性保证任何情况下都能读到最新数据。” 为什么错:它只保证”别人 close 之后、我 open“这一种时序;并发打开的两个客户端会各自看到自己的缓存(22.3.1 的 (S2) 反例)。另外它还有两个前提:open 时真的做了 GETATTR(属性缓存新鲜期可能让它跳过),以及 mtime 的分辨率足以区分两次修改。正确做法:需要强一致时把 actimeo 调到 0、使用文件锁或改用支持回调/租约的文件系统(NFSv4、AFS、GFS 的 append 语义)。
  2. ✗ “AFS 比 NFS 一致性强,所以所有场景都应该用 AFS。” 为什么错:AFS 的强项是”读多写少的中小文件“;对大文件(每次打开整传)和随机小块访问(每次传整个文件,带宽放大几万倍)它是灾难;而且 AFS 同样不保证并发写的正确性(整文件覆盖 ⇒ 丢失更新)。正确做法:按工作负载选型;混用(NFS 挂载共享目录、AFS/对象存储放只读大文件集)是常见做法。
  3. ✗ “GFS 的 chunk 是 64 MB,所以写小于 64 MB 的文件会浪费 64 MB 磁盘。” 为什么错文件末尾的 chunk 可以不满(只有实际数据占用空间),所以 1 KB 的文件只占 1 KB 数据(外加副本);真正的代价是元数据与内存:每个文件至少占一个 chunk 的元数据项,数百个小文件会把 master 的内存吃光正确做法:小文件要么打包成大文件(Hadoop 的 har/SequenceFile/Parquet),要么用别的存储;HDFS 3.x 的纠删码与 NameNode Federation 都是在缓解这类规模问题。
  4. ✗ “GFS 的写是先发控制指令、再传数据吗?”(把顺序记反) 为什么错:顺序恰好相反——数据先沿 pipeline 推到所有副本(进入 LRU 缓冲,此时还不知道偏移),然后客户端才把控制指令(带 data_id、偏移、长度)发给 primary为什么这样设计:数据推送可以流水线并行(每跳边收边转),而控制指令只是一次轻量往返;如果先发控制指令,primary 就必须等数据到齐才能回 ack,把”并行的数据流”退化成”串行的两段式”。这就是”数据流与控制流分离”的实质。
  5. ✗ “既然 primary 分配了序列号,那么并发写的结果是确定的(只是我们不知道顺序而已)。” 为什么错“确定(defined)”的定义不是”实际上只有一个结果”,而是”客户端能够预测内容”。并发写时客户端既不知道序列号顺序,也不知道别人的写落在哪,所以它的写可能被部分覆盖——这正是 “consistent but undefined” 的含义。只有”串行成功”(没有并发写者)才同时是 consistent 和 defined。
  6. ✗ “Record Append 是 exactly-once 的:GFS 说了’原子追加’。” 为什么错:GFS 明确给出的是 at-least-once:失败重试会在部分副本上重复同一条记录,失败副本还会留下空洞(不一致区域),chunk 边界处还会padding正确做法:应用给每条记录加唯一 ID 与校验和,消费时去重;把”重复”当作正常现象而不是异常。
  7. ✗ “pipeline 一定比客户端直发多副本快,因为带宽叠加了。” 为什么错:pipeline 的收益来自”把客户端上行需要承载的数据量从 $kF$ 降到 $F$“,而不是凭空创造带宽。精确定量结论是:$T_{\text{direct}} = kF/U$、$T_{\text{pipeline}} = F/\min(U,B)$ ⇒ 只有 $U < kB$ 时 pipeline 才更快(加速比 $\min(k, kB/U)$)。若客户端上行远快于服务器链路(例如客户端是 10 GbE、副本之间只有 1 GbE,$k=3$),直发反而更快。正确结论:pipeline 的适用条件是”客户端上行是稀缺资源“——这正是 GFS 的假设。
  8. ✗ “checksum 能修好损坏的数据。” 为什么错:chunk 级的 32 位 CRC 只能检测,不能纠正(它不包含足够的冗余来重建原始字节)。修复依赖另一个健康副本(因此需要”副本数 > 1”与”后台扫描 + 副本间比较”)。正确说法:checksum 把”静默数据损坏”降级为”可检测的错误”,再由复制提供修复能力;若两个副本以相同方式损坏,checksum 也救不了。

22.8 思考题(带答案)

思考题 1(计算题):GFS 的元数据与副本成本

设一个 GFS 集群存放 100 TB 逻辑数据,chunk 大小 64 MB,每 chunk 3 副本。 (a) 共有多少个 chunk?(b) master 的元数据内存大约是多少(按”每 chunk 1 KB”与”论文的每 chunk < 64 B”两种口径)?(c) 需要多少原始磁盘容量?(d) 若一台 chunkserver 持有全部数据的 1%,它故障后需要重新复制多少数据?按 10 条并行复制链路、每条 100 MB/s 计算,需要多久?(e) 在这段时间里,受影响 chunk 的剩余副本数是多少?这意味着什么风险?

答案: (a) $\dfrac{100\ \text{TB}}{64\ \text{MB}} = \dfrac{10^{14}}{6.4\times10^{7}} = 1.5625\times10^{6}$,即约 156 万个 chunk。 (b) 1 KB/chunk ⇒ $1.5625\times10^{6}\ \text{KB} \approx 1.6\ \text{GB}$;论文口径 < 64 B/chunk ⇒ 约 100 MB两种口径都能放进单机内存——这就是 64 MB 大 chunk 的核心收益(若 chunk 是 1 MB,同样数据的元数据会膨胀 64 倍,即 6.4 GB~100 GB,单 master 装不下)。 (c) $100\ \text{TB} \times 3 = $ 300 TB 原始容量(还没算校验和与中间文件;校验和仅约 $2^{-14}$ 的数据量,可忽略)。 (d) $100\ \text{TB} \times 1\% = 1\ \text{TB}$。10 × 100 MB/s = 1 GB/s ⇒ $1000\ \text{GB}/1\ \text{GB/s} = 1000\ \text{s} \approx$ 16.7 分钟。 (e) 这 1% 的 chunk 各剩 2 副本。风险是:恢复窗口内可用性被降级——若此时再坏一台持有相同 chunk 的机器,这些 chunk 就只剩 1 副本(再坏一台就永久丢失)。因此恢复速度本身是可用性指标:这就是 GFS 用”后台持续扫描 + 尽快再复制(re-replication 优先级队列)”、HDFS 用”机架感知 + 均衡器”、以及 HDFS 3.x 引入纠删码(把 3× 存储降到约 1.4×,但恢复要读更多节点)的原因。

思考题 2(计算题):pipeline 与直发的边界

客户端上行 $U$、chunkserver 之间链路 $B$、副本数 $k=3$,要写一个 $F = 1$ GB 的 chunk。 (a) 写出两种方式的传输时间公式。(b) $U = 50$ MB/s、$B = 200$ MB/s 时各是多少?加速比?(c) $U = 50$ MB/s、$B = 60$ MB/s 时呢?(d) 什么条件下”直发多副本”反而更快?为什么 GFS 仍然选择 pipeline?

答案: (a) $T_{\text{direct}} = kF/U = 3F/U$(客户端上行要承载 3 份);$T_{\text{pipeline}} = F/\min(U,B)$(客户端只上传 1 份,其余由服务器链式转发;整链吞吐由最慢链路决定)。 (b) $T_{\text{direct}} = 3 \times 1024\ \text{MB} / 50 = 61.4\ \text{s}$;$T_{\text{pipeline}} = 1024/50 = 20.5\ \text{s}$;加速比 (因为 $U \le B$,瓶颈就是客户端上行,而 pipeline 把它从 3 份减到 1 份)。 (c) $T_{\text{direct}} = 61.4\ \text{s}$;$T_{\text{pipeline}} = 1024/\min(50,60) = 20.5\ \text{s}$;加速比仍是 ——只要 $U \le B$,服务器链路就不是瓶颈,pipeline 的收益恒为 $k$ 倍。 (d) 比较两式:$kF/U < F/\min(U,B)$。当 $U \le B$ 时不可能($k>1$);当 $U > B$ 时条件化为 $kF/U < F/B \iff U > kB$。所以当客户端上行快于 $k-1$ 条转发链路的总带宽时,直发更快(例如 $U = 1000$、$B = 100$、$k=3$:直发 0.31 s vs pipeline 1.02 s)。GFS 仍选 pipeline 的真实理由是它的假设:客户端往往是同一集群里的另一台服务器($U \approx B$,且共享上行),此时 pipeline 稳定获得 $k$ 倍收益;同时 pipeline 让”数据推送到多个副本”这件事不占用客户端的上行参与时间,客户端可以立刻去做下一件事。

思考题 3(推演题):Record Append 的可见结果

三个客户端 C1、C2、C3 并发对同一个文件做 Record Append,各追加 1 条 64 B 的记录(带唯一 ID),期间 cs3 发生一次瞬时故障(primary 报错、客户端重试一次后成功)。 (a) 读者可能看到哪些记录序列?(b) 会不会出现”一条记录被劈成两半、和另一条记录交错”?(c) 会不会出现”某条记录完全不见了”?(d) 如果换成普通 write(客户端各自指定偏移为”当前文件末尾”),结果会怎样?

答案: (a) 三个客户端各自的记录都会出现(at-least-once:客户端会一直重试到 OK),但某些副本上可能出现重复(重试的那条被写两次),失败副本上可能留下一个空洞(那 64 B 是零)。因此读者(从某个副本读)可能看到:C1, C2, C3(干净)、C1, C1, C2, C3(重复)、或某个位置是零字节的版本。唯一 ID + 消费端去重可以把这些序列归一。 (b) 不会。 偏移由 primary 串行决定(off := cur,且 cur 在写入后立即前进),因此任何两条被成功应用的记录首尾相接、互不重叠;跨 chunk 边界时 primary 会 padding 并让客户端换 chunk,不会写”半条记录”。 (c) 可能,但只在”部分副本”上:若某次 append 在 cs3 上失败而客户端重试成功,cs3 上可能缺少第一次的记录(空洞),并且它拿到的偏移与别人不同 ⇒ 该副本与其余副本内容不同(inconsistent)但从”文件整体”的角度:客户端重试直到成功 ⇒ 记录至少存在一份(在 primary 上);所以”记录完全消失”只可能发生在”所有副本都失败”的情况下。 (d) 若用普通 write 且偏移由客户端计算为”当前文件末尾”,则两个并发客户端可能算出同一个偏移(都读到”末尾 = X”),于是它们的记录互相覆盖:最终只有一个客户端的记录存在于该位置(lost update),甚至两条记录各写一半而交错(如果写区间部分重叠)。这就是 Record Append 存在的全部意义:把”偏移的决定权”从客户端(可能读到陈旧的文件大小)上移到 primary(唯一的定序者)

思考题 4(辨析题):”无状态一定更好,所以 AFS 用回调是设计退步”

有人说:”NFS 的无状态设计让服务器崩溃恢复变得极其简单,这是明显更好的工程选择;AFS 为了回调而让服务器保存每个客户端缓存了什么,是拿简洁性换花哨功能的设计退步。” 请指出这个说法错在哪,并说明”状态”应当放在哪里。

答案错在把”状态”当成纯粹的负担,而没有看到状态是”信息”——信息可以用来做更聪明的决策。

  • 状态的收益:AFS 的 callback 表记录的正是”谁手里有过期数据“这一信息。有了它,服务器可以主动推送失效(push),于是(i)一致性更强、失效更及时(不需要等客户端下次打开才发现),(ii)服务器负载从”块访问次数”降到”打开/关闭次数”,可扩展性提升一个数量级(22.4.2 实测:$N=50$ 时 AFS 192 次请求 vs NFS 2036 次请求;陈旧读 0 vs 1273)。如果没有这份状态,NFS 就只能让每个客户端一遍遍地问(pull),而这些”问”本身正是服务器负载的来源。
  • 状态的代价:内存 $O(\text{客户端数} \times \text{缓存文件数})$;崩溃后状态丢失必须保守失效(一次惊群式的重新取回);回调不可达时必须超时兜底。
  • 正确的判断标准不是”有状态还是无状态”,而是三个问题
    1. 这份状态值多少钱? 若它能让负载降一个数量级、让一致性变强,就值得存(AFS 的回调表、GFS 的租约表);若它只是”记住客户端上次读到哪”,就交给客户端自己带(NFS 的绝对偏移)。
    2. 丢了它会有多糟? 可重建的状态(GFS 的 chunk 位置靠心跳重建)可以放在内存;丢了会破坏安全性的状态(锁、租约)必须要么持久化、要么有 fencing 机制(GFS 的租约过期 + 自我 fencing)。
    3. 状态能不能放到客户端? 能放客户端的就放客户端(NFS 的 fh 自包含、GFS 的元数据缓存),这样服务器才能既”知道得少”又”扩得大”。
  • 一句话结论“无状态”是一种能让服务器”可任意重启、可任意复制”的工程美德,但它不是免费的最优解;当”记住谁缓存了什么”能换来一个数量级的可扩展性提升时,有状态才是正确的选择。 NFSv4 后来引入回调与租约,正是这一结论的工业界背书。