像是副本、分片之类的概念,我是在 22 年开始正式工作,接手 MongoDB 和 Kafka 后才逐步开始了解的,之前自己本地部署的服务都是单机的,并且在实习期间也没有关注相关系统这方面的特性。随着工作经历的积累,渐渐地对相关的概念有了一定的理解。
概述
本章涉及的概念 Replication ,也即副本,存在的主要目的是为了一个集群中单实例的下线不影响整体的数据处理,比如说有一个存储有数据的实例所在的机器挂了,或者需要升级,那么由于它是多副本的,所以其它的副本实例可以顶上。对于 Kafka 和 MongoDB ,一般配置为三副本,最多可以容忍有两个实例下线。书中提到对于多副本的场景,一般是 leader 承接读写流量,follower 承接读流量,不过目前组内自运维的 Kafka 和 MongoDB 都是只有 leader 承接读写流量,follower 仅作同步数据。其中对于 MongoDB ,如果业务有大批量导出数据的需求,为了不对 leader 造成额外的压力影响正常读写,会要求业务从 follower 导。
书中在开头额外介绍了 Backups 和 Replication 的区别,备份是数据的一个快照,产生后不变;副本是数据的一个复制,会随数据的变化而变化。
数据写入
对于副本的数据同步方式,有 sync(同步)和 async(异步)两种,目前 Kafka 和 MongoDB 均为 async 。印象中 Kafka 有一个参数可以控制至少需要有几个副本同步完成才可以被认为写入成功,印象我看都配置的是 1 ,也即 leader 写入就可以了,一开始觉得这有一定的风险,不过后面发现由于部分 Kafka 集群部署在混布的环境中,机器压力比较大,经常 ISR 中只有一个同步副本(也即只有 leader)的情况,如果不配未 1 ,那么就会影响写入吞吐。如果要求数据不可丢,且对写入性能影响最小,书中建议可以设置至少一个 follower 是 sync 的,其它是 async 。
副本数据同步
如果 leader 挂了,那么需要从 follower 中选出一个新的 leader 出来,然后从新的 leader 处同步数据。首先需要定义 “leader 挂了” ,一般是使用超时时间,接下来是进行选举,要么是 follower 投票产生,要么是由一个更高层的 controller 来指定。对于前者来说,就是经典的分布式共识算法里的 leader 选举,而对于后者,一个例子是 Kafka 分片的 leader ,它就是由 controller 指定的。
在完成 leader 切换后,异步复制的场景下,新的 leader 可能没有相关的数据。后续处理这部分数据直接丢弃,这会导致被认为是写入的数据实际丢了。如果丢弃的数据会和外部系统联动,那么就会有不一致的问题。另外判断 leader 挂了的超时时间也需要根据系统的实际状态来设定。
组里自运维的老版本的 Kafka 集群有时会出现新的 leader 无法分配的问题,现象是原 leader 所在的 broker 挂了后,leader 无法切换导致分片写入失败,数据断流。这个是因为 controller 有啥未知的故障导致的,此时只用将 zk 上的 /controller 删了,同时切新的 controller 就可以了。有时因为集群中的 isr 变更的太频繁了,再加上旧的 controller 卡死,会在 zk 的一个路径下积压大量的 isr change notification 数据,新的 controller 在加载这部分数据时,ls 路径下的子目录时可能会爆 buffer ,导致其无法启动,此时只能人工将其中较老的数据人工清理了,印象中当时这个问题查了挺长时间的,业务反馈说他在 Kafka Manger 上创建了 topic ,但是发现其分片一直没给分配 broker,我很快就意识到是 controller 有问题,然后还发现新的 controller 日志里一直打印 zk 的 response 把 buffer 打爆的,当时额外写了一个 Golang 程序遍历了 zk 上所有的数据,发现都不大,最后怎么发现是 ls 的时候因为子目录过多而导致的问题我有点记不太清了,只记得当时从下午查到了快晚上九点,晚饭都没吃。后面为了防止这种问题再出现,额外部署了一个定时任务来清理历史过久的 notification 数据。
对于 follower 选举,一个线上比较常见的故障, MongoDB 单分片丢主故障的问题,大概是这样的:一个分片常态下是三副本,位于三个实例上,当有实例需要迁移时,PAAS 平台会先扩一个新的实例,等新的实例副本数据同步完毕后,再删老的实例。在同步数据期间,是四副本,也即有四个实例。假如需要迁移的老实例因为某种原因直接挂了,然后另一个实例异常也挂了(常常是主节点),那么此时四副本中只有两个实例了,由于 leader 挂了所以触发了选举,但是此时无法选主(2/4),整个分片就挂了。解决方案就是连到一个正常的 secondary 实例上,将其中的一个异常实例从 rs 的配置里移除,此时是三副本,一个离线实例,两个 secondary 实例,然后强制更新 rs 就可以了。
针对 MongoDB 扩容场景下同步数据导致 leader 挂掉的问题,我专门研究和分析过,发现都是 leader 实例由于压力过大,将内存数据 dump 到了本地的一个 wt 文件中,导致磁盘打满进程退出。我怀疑是在同步数据的时候,由于是全量数据同步,给 leader 带来了额外的负担,又了解到 MongoDB 支持从 Secondary 节点同步数据,所以就搞了一个定时任务,每一分钟扫一下所有集群中的所有副本集,如果发现有 Startup2 状态(第一次同步数据)的实例从 Primary 节点同步数据,然后就将其切为从 Secondary 节点来同步。这个定时任务上线后,基本就没再发生过这类问题了。
对于同步的数据格式,书中列了三种格式:1)原始的请求语句,2)WAL,3)Logical log。对于原始语句,一般需要进行额外的处理,比如将其中实时生成的数据固定下来,例如 now()。对于 WAL ,需要做好存储引擎版本的兼容,因为不同版本的 WAL 格式可能有差异。对于 Logical log ,比如 MySQL 的 binlog ,是用于描述 row 粒度的数据是如何变更的,独立于存储引擎,所以可以做到存储引擎版本的兼容,并且可以给外部的系统来使用,例如数据分析系统。
除了由数据系统本身提供的数据同步外,还有一些场景,例如在做 MongoDB 表数据热迁移时,会用到额外的工具来做。应该是 23 年的时候,当时需要将几个表从一个 MongoDB 集群迁移至另一个集群,但不能停写。当时用的是一个叫做 MongoShake 的工具,大概的功能是先做 base 数据同步,然后订阅各个分片的 log 做增量同步。源集群需要配置表所在 shard 的地址信息,而目标集群只用配置 mongos 地址即可。根据工具介绍,为了防止数据同步死循环的出现,这个工具的提供方,好像是阿里云,的 MongoDB 的 log 中加了额外的信息,MongoShake 可以识别这个信息,对于包含这个信息的源数据不做同步。
副本读写一致性
对于写少读多的场景,一般会让 follower 也承接读流量来提供系统的读吞吐,但是由于副本数据同步有 lag ,所以从 follower 处读到的会是老数据。书中列了几个 async 的场景:
reading your own writes : writer 写入数据后立刻读这个数据,如果是从 follower 读且 follower 同步数据不及时,就会导致 writer 读不到写入的数据。一种方案是通过某种方式识别这种场景,然后让 leader 承接这部分读流量。
monotonic reads : reader 多次读某条数据,由于同步有 lag ,会导致读到某个副本上老版本的数据,出现类似于时间回退的效果。一种方案是让相同的读请求落在同一个副本上(当然如果这个副本挂了就无效了)。
consistent prefix read : 读一批数据,多次读是乱序的,这种在分片的场景比较常见。比如 t1 和 t2 时刻写入的数据落在了分片 1 和 2 中,2 副本同步的比 1 快,就会导致 reader 先读到 t2 再读到 t1 。
由于我这边维护的系统一般情况下仅由 leader 承接读流量,所以基本没遇到这种问题。不过在使用别人的系统时倒是遇到过,比如读一个 kv 的数据,发现每次请求读到的数据版本不一样,这个大概率就是因为有问题 follower 副本导致的。
Multi-leader 和 Leaderless replication
对于规模比较小的业务,使用 Single-leader 就足够了,但是如果涉及到跨地域的场景,例如数据服务部署在多个地域(region),如何在一个地域的数据变更时,让其它地域也感知到呢?如果采用 Single-leader 的话,leader 只能在其中一个机房,而跨机房访问 leader 成本较高。如果业务数据是天然可以按照地域进行切分的,对跨地域的数据读写取延时不敏感,可以在每一个地域的机房都设置一个 leader ,每个地域数据的读写都由这个 leader 来进行,然后各个 leader 之间异步进行数据同步。不过此时就会有数据冲突的问题,一般最简单的方法就是给数据的变更打上时间戳,然后采用 LWW (Lastest Write Win)的策略,仅保留时间戳最新的变更,老的都丢弃。其实既然已经选型了这种模式,那么对应的业务场景下就需要可以容忍因为有冲突而导致的数据丢失,否则最好还是老老实实地用 Single-leader 。对于 Multi-leader ,书中介绍一个非常有意思的场景,就是多端同步的应用,比如日历 app 、Git 以及在线协作文档。每一个客户端应用都可以看做是一个 region。
对于搜索的架构来说,从离线架构的视角看,六地域(现在砍成了五地域)机房的正排/摘要/索引数据都是只读的,是由离线 yq (阳泉,robin 的老家)机房生产后,通过主动推送或被动拉取的形式同步到六地域机房的。之前值周的时候处理过几次小说业务评论数据的问题,这个业务的数据流基本也是这样的,六机房的在线服务处理用户的评论数据时,会将数据推到 yq 机房的 Kafka 里,然后走离线建库再生效到六机房。类似于是一个 Star 的拓扑,阳泉机房位于中心,然后连接其它的机房。
对于 Leaderless replication ,我看的更是云里雾里了 … 这直接让每个副本都可以承接读写流量,然后通过某些策略对数据进行 merge 。在数据读取和写入时,需要并行地访问多个副本,以保证可以写入成功以及读取到最新的数据。书中给了一个公式,假如是 n 个副本 ,成功写入的节点数为 w ,查询的节点数为 r ,那么要求 w + r > n 。一般来说,会设置 w = r = (n + 1) / 2 。对于公式 w + r > n 的理解,其实可以从反着去理解,也即假如 w + r = n 会发生什么?在计算情况下,假如 w = r = n/2 ,那么如果我写入的 1/2 个副本和读取的 1/2 个副本刚好不重合,那么就读不到最新写入的数据了,如果 w + r < n 就更读不到了。那么假如是普通 3 副本的场景,就要求同时读写 2 个节点,可以保证写入的数据不丢且读到最新的数据。在 Multi-leader 的场景下就会有冲突的问题,那么在 Leaderless replication 的场景下冲突就会跟频繁了,书中介绍了不少冲突发现和解决的内容,但是我感觉目前应该用不太到,等后面遇到的时候再仔细研究吧。
碎碎念
另外吐槽一点,虽然在系统正常时,如果是仅由 leader 承接读写流量,那么多副本看起来确实没啥用处,但是在 leader 故障后,却是可以保证系统整体是可以继续运行的。系统如此,工作也是如此,比较理想的情况是确保一项工作/系统有多个人了解,这样在负责人离职后,有人可以快速接手,但是现在的现状常常是仅有负责人了解,其他人完全不了解(也没人力去了解),这样一旦负责人离职,仅靠离职交接是完全交接不清楚的。不过这个也可以从侧面印证,这种单副本的工作/系统往往对整个公司来说是不重要的,至少在领导的眼里不重要 …