书中之前提过,设计系统时,如果一个故障场景没考虑到,那么当它发生时大概率没法解决 …
partial failure
对于单机服务来说,要么它是正常运行的,要么因为系统/代码等问题挂了(fault),主要就是这两种状态,所以整体来说是 determinism 的。而对于分布式服务来说,服务是部署在一批机器上的,比较常见的情况是有一部分服务挂了,也即 partial failure 。如果这种情况处理不好的话,就会导致一些 nondeterminism 的问题,这类问题比较难排查。假如说可以处理好 partial failure 的情况,那么就可以得到一些好处,比方说可以对系统进行不停机地升级,也就是一次只升级一部分实例,这些实例在重启期间无法对外通过服务,其实也可以看做是 partial failure。
网络
之所以被称为 “分布式” ,原因是因为这个服务是由多个实例组成的,实例之间如果需要进行通信,必然需要通过网络来进行,没法像单实例多线程那样通过共享内存/文件来通信。跨机器的网络通信是不稳定的,因为这种通信模式一般来说是异步的,以 TCP 为例,当发送者调用系统接口写 socket 后,系统会对数据进行拆包并放到队列中,然后一次性发送一批网络包,当接收端收到数据后,也会先放到队列中,然后在被处理后发 ack 告知发送端数据已经被成功发送,最后应用层的 socket 接口才返回。这其中会有多个地方导致数据丢失,比如说如果发送队列满了、数据包在传输时丢了、接收队列慢了、接受者来不及发 ack 就挂了等。数据的发送者没法精确感知到数据包的状态,所以一般都是通过设定一个超时时间来判断数据是否发送成功的。对于超时时间具体的值,设置的时候需要谨慎一些,因为如果超时时间设置的过短,会导致无用的重试次数过多导致系统负载升高最终压垮整个系统,但如果超时时间设置的过长,会导致无法及时感知到故障实例。
实际上也有同步传输的通信模式,但这要求为一次通信预留固定的带宽,即使通信期间实际没有数据传输,这部分预留的带宽也不能被其它的通信使用。这虽然保证了通信的稳定性,但浪费了通信的带宽资源。采用异步通信的模式的好处是可以最大化利用带宽资源,应对 burst traffic。为了防止同一台机器上部分实例把网络带宽打满,可以通过 TC 或 ebpf 来实现网络流控。
时间
除了网络通信外,还有一个问题就是时间。由于服务实例分散在多台机器上,这些机器上的时间不一定完全一致,有两类时间:1)Time-of-day clocks 是当前的日期和时刻 ,2)Monotonic clocks 是从某一时刻开始累加的值,适合来计算时间间隔(比如是否 timeout)。如果有些数据处理依赖多个实例来完成,并且需要基于各个实例执行的时间来做一些判断,那么如果使用上述两种时间都会产生误差,此时就可以引入 logical clocks ,它是基于计数器,而非真实时间来实现的。基于 logical clocks 可以判断各个实例上发生事件的先后顺序。一般来说,实例 i 维护向量 VC[i] ,当本地事件发生时,VC[i] += 1 ,请求其它实例时会带上这个向量,收到其它实例请求时,逐项取 max,再 VC[i] += 1 。在进行判断时,如果 VC1 中所有的值都小于 VC2 ,则代表 VC1 在 VC2 前发生,反之在后发生。如果有大有小,则代表两者是同时发生的,没法判断先后。Google 的数据中心机器提供了 TrueTime ,它会返回一个时间区间 [earliest, latest] ,表示当前时间是在这个区间中的。
共识
分布式系统中因为部分实例故障而导致的问题,可以通过 quorum 来解决,也即系统的状态不依赖于单实例,而是由集群中多数实例,来进行判断。比如说如果系统中多数实例无法访问某一个实例,虽然这个实例本身是正常的,但是它依旧可以被认定为故障实例。之所以不能依赖单实例,原因之一是单实例的故障除了彻底挂掉外,还可能是被 pause 了或者与外部断连,当它恢复正常运行时,它认为的集群的状态和正式的状态大概率已经不一致了,类似于是 “乃不知有汉,无论魏晋”。为了不让故障实例破坏系统整体状态,需要引入 fencing token 机制,也即系统会授予正常的实例一个 token ,数据的操作需要凭借这个 token 来完成,而问题实例使用的 token 是老的,所以在被认定为故障后,就无法继续改变系统的状态。一个典型的样例就是分布式选主算法中的 epoch ,每当一个新任 leader 被选出时,epoch + 1,老的 leader 发给成员请求中的 epoch 由于比新的小,所以成员都会拒绝这个请求。
书中举的例子是分布式锁 + 写文件,也即一个文件预期同时只能被一个 client 写,所以 client 在写之前会请求 lock service 取锁,当 client a 拿到锁之后短暂挂掉了,锁到期释放,接着被 client b 拿到了,然后 b 就开始写文件。a 恢复后,认为自己依旧持有锁,所以它也往文件里写,此时有两个 client 同时写文件,导致故障。即使 client a 在每次写文件前都检查自己是否持有锁,也是会有问题的,在它第一次持有锁时发的写请求,由于网络延时,到 client b 拿到锁后写文件期间才发到,此时从现象上看也是有两个 client 同时在写文件的。一种解决的方式就是在文件写入端校验这个写请求是否来来自于最新持有锁的 client。it is not safe to assume that only one node is holding a lease at any one time 。
书中提到了两个概念,分别是 “safely” 和 “liveness” 。前者表示 “nothing bad happens” ,而后者表示 “something good eventually happens” 。根据描述,有点像是 “强一致性” 和 “最终一致性” ,或者是 CAP 中的 CP 和 AP 。当 P(网络分区)发生时,对于 safely ,只要出现不一致性(C),就直接拒接请求,这就影响了可用性(A),而对于 liveness ,由于认为数据最终肯定会一致的(C),所以当出现不一致的时候不会拒绝外部请求(A)。
小结
所以由于分布式系统的特性,必然涉及到部署在多台机器上的实例进行协作来完成数据处理,而协作使用的网络和时间的问题,会引发出一系列的 nondeterminism 问题。在设计系统时,需要 “focus on what could go wrong,even though it may be unlikely” ,因为 “anything that can go wrong will go wrong” 。这让我想起了最近遇到的一件事情,去年为调研新版 Kafka 部署了一个十节点的集群,当时遇到了 raft 丢主的问题。丢主后,由于无法从外部干预 raft log ,所以只能通过清空包含 raft log 的 metadata 目录重建集群的方式来解决这个问题,这意味着集群中 topic 的信息都丢了,这对于生产环境是不可接受的,所以优先补了 raft voter 和 leader 的监控和电话报警。但是后面觉得没有一个干预的方式还是不保险,所以就改了 Kafka 的源码,加了一个人工干预 raft log 中 voter 信息配置项,想着万一哪天用得上呢(详见之前的一篇博文《kafka raft 模式无法选主问题介绍及止损方式研究》)。在写本文的上一周,隔壁组有个同事跟我说,他之前部署 Kafka 时用的是我改完源码的版本,线上有个很重要的 Kafka 集群中的实例因为 oom 全挂了,是用我新增的这个干预配置项来完成对集群的恢复的。所以幸好留了这个口子,虽然目前我这边部署的集群没用上,但是帮了隔壁组也是挺有价值的。