重读 DDIA : Chapter-7 Sharding

反思了一下为啥现在我会对 副本分片 这两概念那么的情有独钟,大概是因为工作久了,对于线上的系统,只要不需要频繁人工介入止损、性能不退化以及不半夜触发电话报警就烧高香了,至于其它的都是次要的。而副本和分片这两者就可以预防上述问题的出现,算是我优质睡眠和高效工作的大恩人。首先说明一下,这里提到的都是有状态的数据系统,也即实例需要维护有状态的数据。

为什么需要分片

上一节中提到的副本,可以降低数据丢失的概率,本节中的分片,则可以将数据丢失的概率进一步降低的同时,提升服务的性能和可用性。

首先我们需要明确 “分片” 是为了解决什么问题。在没有分片的时候,例如 MongoDB 的一个 Replica Set ,又或者是 Kafka 的单 Partition ,虽然可以通过多副本的形式保证数据不丢,但是随着数据量的增长,每一个实例上存储的数据量越来越多,即使单实例的性能足够强,处理数据读写请求都 ok ,但是假如实例所在的机器故障或者重启,那么新实例同步数据的时间会随着数据量的增加而变长,随着时间的拉长,同步失败以及多台机器同时挂的概率都会增加,这依然会对服务的可用性造成影响。另外数据同步期间会占用额外的带宽以及服务性能,这对服务本身的性能也是有影响的。

一种解决方式就是对数据进行拆分,将其分散到一批 Replica Set 中,这种方式就称之为分片。对数据进行拆分后,每一个分片上需要维护的数据量就会有所减少,这样在进行实例迁移时,新实例完成数据同步的时间也会相对缩短。如果数据的拆分策略适配数据的业务场景,那么也会有效地减少单分片需要处理的请求负载,从而提升系统整体的数据处理性能。需要注意的是,由于一个分片只会持有一部分数据,所以它不保证服务整体的可用性,也即如果一个分片挂了,那么整个系统就会缺了一部分数据,也即整个系统也挂了。

如果业务的数据量预期整体较少,并且整体读写的频率也不高,那么就完全没必要使用分布式的分片,直接单分片 + 多副本即可,或者在单台机器上启多个进程模拟分片。

这点我还是挺有体会的,25 年给新的离线架构系统做监控数据采集,采集的模式为从各个模块的所有实例上拉 metrics 数据,进行一些加工(比如打上一些 tag 以及聚合)后返回给 prometheus。当时评估单个模块整体的实例量级大概在 10w - 20w 左右,1分钟采集一次,在写 prometheus exporter 的时候,如果把 goroutine 拉满(网络IO处理的 goroutine 可以设置高一些,数据处理的 goroutine 数量比 CPU 核心数少一些,整体通过 chanel 流式处理),单实例应该就可以满足一个模块所有实例的监控采集需求。后面验证也是如此,单 exporter 部署在一台物理机上,采集延时大概在 10s 以内。

由于为了稳定性,使用的是集团云提供的 Prometheus,它对于单次采集有数据条数和大小的限制,而由于 exporter 单次采集的监控项是模块粒度的,前面提到,大概在 10w - 20w 左右 ,超过了单次采集大小的限制。此时 分片 的思想就有用处了,我可以在一个 Prometheus 采集配置中配置多个 exporter 地址,每个 exporter 采集一个模块一部分实例的监控数据,并且这个 partition 的数量是可以动态调整的,也就是说如果实例数再增加,我只用调高 partition 的配置就 ok 了,后面确实遇到了这种问题,调整了多次。

但此时就有一个问题,如何协调多个 exporter ?我当时大概评估了一下,首先单台物理机足够支持 exporter 对全量实例进行监控采集,所以没有必要搞成分布式的模式,将 exporter 分散在多台机器上,其次,也没必要启多个进程,因为还需要考虑进程通信的问题以及进程状态的维护,最后想了一个最简单的办法,exporter 主程序启动 partition 数量对应的 http server ,每个 http server 实际上就是一个 partition exporter ,负责一部分实例的监控数据采集。至于采集实例的分配,也很简单,由于实例数没有特别多,也即 10w 左右,所以我干脆把所有实例的信息都放内存里基于 id 排个序(物理机内存有 100 GB ,足够了,实在不行也可以参考 sort 用额外存储做归并排序),然后根据 partition 数量进行切分,每个 exporter 就取属于它的那一份实例列表进行采集即可。

数据如何分片

既然以及将数据拆分到了多个分片上,那么就需要重点关注数据是如何写入和查询的了,其中的核心点就是数据到分片的映射策略,也即如何决策一条数据应该被写到哪个分片上,以及如何知道一条数据在哪个分片上。

对于 Kafka 这种消息队列服务,一般情况下数据是顺序批量读写的,基本不存在需要 seek 指定数据的情况。在数据写入的时候,可以指定数据对应的 key ,然后对这个 key 进行 hash ,映射到对应的 partition 上。在数据读取的时候,消费者是从设定的 parittion 上顺序读取数据的,一般情况不用关心数据到 partition 的映射策略。在搜索离线建库的场景下,一般是使用 url 来作为 key ,因为我们希望同一条 url 的数据,可以被顺序处理,也即都被写到一个 partition 中,这样消费者一个线程消费一个 partition ,可以保证相同的 key 的新老数据被顺序处理。不过也有一些场景,业务预期同一个 url 在很短的时间内会触发很多次建库,同时对数据处理先后的顺序不敏感(数据处理无状态,且在出口处有基于时间戳的历史数据过滤),在这个场景下,在生成 key 的时候,会在 url 后面在拼一个时间戳 + 随机数,这样就可以让相同 url 的数据足够分散,否则往同一个 partition 中写入了大量的数据。

对于 MongoDB 这种数据库服务,需要考虑的地方就比较多了。目前是有两种分片策略:

  • Sharding by Key Range :根据 Key Range 来分片,就是将普通索引作为 Shard 的分片键,例如 sh.shardCollection("xxxx", {"idx":1}) 。普通索引使用的是 B+树,其中的数据是有序的,所以相邻数据会位于一个 shard 上,如果设置的不好,会有相邻 key 热点的问题,比如说将时间戳作为分片键,那么一段时间内写入的数据都会落到一个 shard 上。不过由于数据整体是有序的,所以在进行范围查询的时候,可以快速定位需要查询的 shard ,而不用通过广播的形式查询所有的 shard 了。

  • Sharding by Hash of Key:根据 Hash Key 来分片,就是将 Hash 所以作为 Shard 的分片键,例如 sh.shardCollection("xxx", {"idx": "hashed"}) 。对于 MongoDB 来说,Hash 索引底层使用的还是 B+树,而不是 Hash 表,只不过是在数据插入 B+树时,先对 key 进行了 Hash 操作。采用这种分片方式,可以让数据在各个 shard 上分布的较为均匀,不容易造成读写热点,但是代价是在进行范围查询的时候,mongos 需要请求所有的 shard 。

这两种方式没有优劣,我们需要做的就是基于实际的业务场景来选用合适的分片键。目前离线建库对业务提供的MongoDB ,默认采用的 Hash 索引。基于我从入职到现在协助业务排查问题时捋的和 MongoDB 交互相关的业务逻辑,原先对业务提供的 MongoDB 集群,主要的功能是将建库 url 对应的数据进行存储 ,所以基本执行的都是 kv 查询。在这种场景下,没有进行范围查询的需求,并且为了避免同一站点的 url 的数据建到一个 shard 上,所以使用 url hash 作为分片键是没问题的。但是因为这个自运维的 MongoDB 集群的成本比公司集团云提供的 MongoDB 要低的多的多(据一位前同事说差了八到十倍),所以不但原先业务场景中写入的数据有了新的查询场景,而且后来接入了很多非建库业务,这就引入了不少一级和二级索引的范围查询/更新以及二级索引的 kv 查询/更新,而这些场景必定是需要通过访问所有的 shard 来完成的。不过好在目前两套业务集群将近 30 个分片主要的瓶颈不是读写性能(读 8w qps 写 2w qps ioutil 50% 整体延时 1ms以内)而是存储(单实例 ssd 900 GB , 目前用量 70%),所以目前也没有明显的性能问题,倒是因为存储的问题扩了好几次 shard 。

另外还有一个组内自用的 MongoDB 集群,原先是 3 个 shard ,出了一次事故后扩过一次容,现在是 6 个 shard 。说起来这个事故的场景就是上一节中介绍的使用分片的原因。这个集群里有一个超大表,出故障时将近 3 TB,现在已经 3.3 TB 了。正是由于这个大表的存储,所以每次 shard 中有实例进行迁移时,经常会因为这个表的量级太大,同步时间过长且期间因为网络抖动导致数据同步失败。那次事故大概是这样的,有三个副本实例,新实例从其中一个 Secondary 实例同步数据,一直同步失败,此时另一个 Secondary 所在的机器挂了,而 Leader 实例也因为异常 wt 把磁盘打满也挂了,那么这时只有正在同步数据的新实例以及它的数据源 Secondary 实例,但问题是此时只有两个实例了,而副本集中有四个实例,导致选不了主,这个分片直接就挂了。后续 case study 时,一个 todo 就是新增集群 shard 数量,直接 double 到 6 个 shard,之后就没有类似的问题了。

在创建 MongoDB 表 unique 索引时,假如这个表开了分片,那么 unique 索引中必须包含分片键,否则索引无法创建。这个设计确实是合理的,因为如果不加这个限制,那么数据在写入的时候,为了保证 unique 的特性,需要先查询所有的 shard 来确认是否有重复的数据,这会极大影响性能。而如果 unique 索引的 key 中就包含分片键,那么只用访问分片对应的 shard 即可。

新版的 MongoDB 额外支持了混合分片键(compound shard key),采用的是 hash + range 的模式,这样子对于不同的 hash key,可以做到足够分散不至于有读写热点,而对于同一个 hash key ,由可以利用 range 的特性来进行高效的范围查询。

在建库场景下,一个库种(业务)对应多个库层,这个库层也可以被看做是分片。一般来说,一个 url 会先算一个 sign ,可以视为 hash ,然后基于这个 sign 再计算需要落到哪个库层。具体的算法没有研究过,但是目前确实可以保证各个库层的数据分布地相对均匀。之前在做离线建库架构迁移的时候,在进行大批量数据测试时,一个核心的关注点就是新老架构产出的库层的 url list 是否有 diff ,因为 url 预处理、使用的 sign 以及 shard 的算法的差异会导致 url 最终被建到不同的库层上。如果 url 的 diff 不消除,就没法继续比其它数据的 diff 了。

分片信息存在哪儿

当数据被分散到各个分片上后,不论是新数据的写入,还是存量数据的查询和更新,都需要知道:1)分片所在的地址。2)分片上 key 的范围。一般来说都这个信息会维护在外部存储或者有服务本身进行维护。

对于 Kafka 来说,老版本是将元数据存到 ZooKeeper 上,作为 source of truth ,然后由负责 controller 的 broker 缓存这部分信息,并将变更分发给各个 broker ,每个 broker 内部都存有一份全量元信息的 cache 。在 client 进行元数据查询时,直接连任意一个 broker 查询即可。新版使用的是 KRaft 模式,也即有一个内部的 __cluster_meta topic 存储有全量的元信息以及变更,所有的 broker 都会从 KRaft leader 处同步这个信息,数据同步直接复用了 Topic 副本数据同步机制。

在为组里的建库业务调研 KRaft Kafka 时,遇到了一个比较难受的问题,__cluster_meta 这个 topic 非常非常重要,优先级是高于其它所有的 topic 的。在部署时,给 broker 挂载了一块独占的 HDD ,原先是将 __cluster_meta 和其它 topic 一起都放在这个 HDD 上的,但是在进行压测时,在将 HDD ioutil 打满的情况下,__cluster_meta 的数据同步会出现滞后的问题,此时这个 broker 的状态就会有些异常了(具体异常的场景有点忘记了)。最后正式部署的时候就将 __cluster_meta 放到了部署及运行环境的磁盘上,虽然这个磁盘是公用的,但是 ioutil 一般不高。后续发现这样做也是有问题的,因为虽然 PAAS 平台对服务实例的 CPU 和内存做了硬限,但是有些实例的磁盘使用没限制住,偶发的异常实例 coredump 直接把公用的磁盘打满了,导致 __cluster_meta 数据无法写入,broker 就直接挂了,这个比之前提到的异常场景还严重。联系了好几次 OP 但是都没解决,最后只能我人工登实例去删那些 coredump 文件。

对于 MongoDB 来说,和 Kafka 依赖外部存储不同的是,它会将分片的信息存储在一个独立的 Replica Set 中,称之为 Config Server 。当连接 mongos 访问 config 数据库时,实际查询的就是 Config Server 中的信息,其中记录着各个 shard 对应的 Replica Set 地址,以及每一个 shard 存储的 key 的范围。在我的印象中,线上使用的三套 MongoDB 集群的 Config Server 一直都没出过问题。

数据如何均衡

分片的数量可能随着业务的增长而需要增加,在新增分片后,就需要往新的分片上搬存量分片的数据,并且需要确保相关数据的读写请求能被路由到正确的分片上。对于 hash 的场景,可以使用一致性 hash 算法,使用这种方式可以保证在新增分片后,只用重新分配极少量的数据。

Kafka topic 在做 partition 扩容后,基本不会进行存量数据搬迁的操作,但是在新增 broker 后,需要将存量 broker 上的部分 partition 副本迁移到新的 broker 上。MongoDB 在新增一个 shard 后,存量 shard 上的数据会以 chunk 为单位来进行搬迁,在完成搬迁后,需要重建老 shard 的数据以最终释放磁盘空间。

多租户

书中还提到了分片的一个作用,就是多租户的场景来使用,也即将不同用户的数据放到不同的分片上,以实现数据隔离。目前离线建库用于采集与存储监控数据的 Prometheus ,我感觉就是类似于这种形式。组里有一批自运维的物理机,Prometheus 就部署在上面。原先仅部署了两到三个 Prometheus ,后续随着模块和实例的增加,经常出现 Prometheus oom、wal 来不及 compact、磁盘打满、fd 打满各种问题,后面就根据模块的优先级和规模对监控采集任务进行了拆分,在其它的物理机上又部署了新的 Prometheus 来承接原先 Prometheus 中拆分出来的监控采集任务。这个算是人工进行了分片。至于路由的话,我们使用的是 Grafana 来配置监控看板和报警规则的,所以就在 Grafana 的配置中注册上 Prometheus 的地址信息及 id 即可。我后面好奇查了一下,发现 Prometheus 本身居然没有提供多副本以及分片的机制 …