分片

分布式数据库通常通过两种方式在节点间分布数据:

  1. 在多个节点上保存相同数据的副本,这就是 复制(replication)
  2. 如果不想让每个节点都存储全部数据,可以将大规模数据集拆成更小的 分片(shard)或 分区(partition),再把不同分片存放在不同节点上

通常情况下,每条数据(每条记录、每行或每个文档)属于且仅属于一个分片。

分片通常与复制结合使用,使得每个分片的副本存储在多个节点上。这意味着,即使每条记录只属于一个分片,它仍然可以存储在多个不同的节点上以获得容错能力。

图 7-1 复制与分片结合使用:每个节点对某些分片充当领导者,对另一些分片充当追随者

分片,在不同软件中有许多不同名称:Kafka 称其为 分区(partition),CockroachDB 称为 范围(range),HBase 和 TiDB 称为 区域(region),Bigtable 和 YugabyteDB 称为 表分片(tablet),Cassandra、ScyllaDB 和 Riak 称为 虚节点(vnode),Couchbase 则称为 虚桶(vBucket)

一些数据库把分区和分片视为两个不同概念,在 PostgreSQL 中,分区是把一张大表拆成存储在同一台机器上的多个文件(这样做有若干好处,比如可以极快地删除整个分区);分片则是把数据集拆分到多台机器上。

有一种说法,shard 原本是 System for Highly Available Replicated Data(高可用复制数据系统)的首字母缩写

分片的利与弊

对数据库进行分片,主要是为了获得 可伸缩性(scalability) :当数据量或写入吞吐量大到单个节点无法承受时,分片可以把数据和写入分散到多个节点上。(如果瓶颈是读取吞吐量,则未必需要分片,可以采用 读扩展)。

分片是实现 水平扩展(horizontal scaling,也称 横向扩展,scale-out 架构)的主要手段之一。

复制可以提供容错和离线运行能力,因而无论规模大小都有用;分片却是一种重量级方案,主要适用于大规模场景。如果数据量和写入吞吐量仍可由单台机器处理(如今单机的能力可不容小觑!),通常最好避免分片,坚持使用单分片数据库。

分片往往会增加复杂性。通常需要选择一个 分区键(partition key) ,据此决定每条记录应放入哪个分片;分区键相同的记录都会进入同一分片。这个选择十分重要:如果知道记录在哪个分片,访问就很快;如果不知道,就只能低效地搜索所有分片,而且日后很难更改分片方案。

分片通常很适合键值数据,因为可以直接按键分片;关系数据则比较棘手,因为你可能需要通过二级索引搜索,或连接散落在不同分片中的记录。

关于分片,一次写入可能需要更新多个不同分片中的相关记录,单节点事务相当普遍,但要保证多个分片之间 的一致性,就需要 分布式事务(distributed transaction)

有些数据库支持分布式事务,但这类事务通常比单节点事务慢得多,可能称为整个系统的瓶颈,还有些系统根本 不支持分布式事务。

有些系统会在单台机器上使用分片,通常每个CPU核心上运行一个单线程进行,以利用CPU的并行能力。例如,Redis、VoltDB 和 FoundationDB 都采用每个核心一个进程的方式,并依靠分片把负载分摊到同一台机器的各个 CPU 核心上。

面向多租户的分片

软件即服务(SaaS)产品和云服务通常采用 多租户(multitenant)模式,每个租户对应一个客户。例如在电子邮件营销服务中,每家注册企业通常都是一个独立租户,因为各家企业的简报订阅信息、投递数据等彼此无关。

用分片实现多租户有以下优点:

资源隔离

如果某个租户执行计算开销很大的操作,只要它与其他租户位于不同分片,其他租户的性能就不太容易受到影响。

权限隔离

如果访问控制逻辑存在漏洞,只要各租户的数据集在物理上彼此隔离,意外让一个租户访问另一租户数据的可能性就会降低。

单元化架构

分片不仅可以用在数据存储层,也可以用来划分运行应用代码的服务。在 单元化架构(cell-based architecture)中,为一组特定租户服务的应用与存储会组成一个自包含的 单元(cell),不同单元大体可以彼此独立地运行。这种方法能够实现 故障隔离(fault isolation):一个单元里的故障只影响该单元,不会殃及其他单元中的租户。比如游戏服务器的大区、分区、服。

按租户备份和恢复

分别备份每个租户的分片,就能从备份中恢复某个租户的状态,而不影响其他租户。租户意外删除或覆盖重要数据时,这一能力很有用。

法规合规性

GDPR 等数据隐私法规赋予个人访问并删除关于自己的全部存储数据的权利。如果每个人的数据都存放在独立分片中,实现这一权利就只需对相应分片执行简单的数据导出和删除操作。

数据驻留

如果数据驻留法规要求某个租户的数据必须存放在特定司法管辖区,那么区域感知数据库可以把该租户的分片分配到指定区域。

逐步推出模式变更

模式迁移可以逐步推出,每次只迁移一个租户,这样能在问题波及所有租户之前将其发现,从而降低风险,不过很难以事务方式完成。

使用分片实现多租户的主要挑战是:

键值数据的分片

假设有大量数据并且想要分片,如何决定在哪些节点上存储哪些记录呢?

分片的目标是将数据和查询负载均匀分布在各个节点上。如果每个节点公平分担数据和负载,那么理论上,10 个节点应该能够处理单个节点 10 倍的数据量和 10 倍的读写吞吐量(暂时忽略复制)。

在添加或移除节点时,我们希望能够 再平衡(rebalance)负载,使它均匀分布在增加后的 11 个节点上,或移除节点后剩余的 9 个节点上。

如果分片不公平,某些分片承载的数据或查询比其他分片更多,我们就称其为 倾斜(skew) 。倾斜会大幅降低分片的效果。在极端情况下,全部负载都可能集中到一个分片上,10 个节点中有 9 个闲置,瓶颈却卡在唯一繁忙的节点上。负载高得不成比例的分片称为 热分片(hot shard)热点(hot spot) ;如果某个键的负载特别高(例如社交网络中的名人账号),则称为 热键(hot key)

需要一种算法,以记录的分区键为输入,指出这条记录属于哪个分片。在键值存储中,分区键通常就是键或键的第一部分;在关系模型中,它可以是表中的某一列,不一定非得是主键。为了缓解热点,这种算法还必须便于再平衡。

按键的范围分片

一种分片方法,是为每个分片指定一段连续的分区键范围

图 7-2 印刷版百科全书按键范围分片

各段键范围不一定等宽,因为数据本身很可能分布不均。如 A B 开头的单词就比 X 开头的多得多。 为了均匀分布数据,分片边界必须根据数据进行调整。

分片边界既可以由管理员手工选择,也可以由数据库自动确定。例如,Vitess(MySQL 的分片层)采用手动的键范围分片;Bigtable、其开源版本 HBase、MongoDB 的范围分片选项、CockroachDB、RethinkDB 和 FoundationDB 则采用自动方式 。YugabyteDB 同时支持手动和自动拆分表分片。

每个分片内部都按顺序存储键,使用B树或SSTable,很容易执行范围查找,把键当作联合索引,在一次查询中获取多条相关记录。

再平衡键范围分片数据

首次建立数据库时,还没有数据可供划定键范围。一些数据库(如 HBase 和 MongoDB)允许在空数据库上配置一组初始分片,这称为 预拆分(pre-splitting)。采用这种办法,必须事先大致了解键将如何分布,才能选出合适的键范围边界。

此后,随着数据量和写入吞吐量增长,采用键范围分片的系统会把现有分片拆成两个或更多较小分片,每个新分片都保存原键范围中的一段连续子范围;这些较小的分片随后可以分散到多个节点上。如果大量数据被删除,几个相邻且已经变小的分片也可能需要合并成一个较大的分片。这个过程类似于 B 树顶层发生的变化。

对于自动管理分片边界的数据库,分片的拆分通常由一下情况触发:

键范围分片的优点是,分片数量能够随数据量调整。数据很少时,只需少量分片,开销也很小;数据量巨大时,每个分片的大小仍会被限制在可配置的上限之内。

这种方法的缺点是,拆分分片代价很高:必须把其中的全部数据重写到新文件里,类似于日志结构存储引擎的压实操作。需要拆分的分片往往本就处于高负载,拆分开销还会雪上加霜,甚至使它彻底过载。

按键的哈希分片

如果希望相邻的分区键进入同一个分片,键范围分片就很有用,如果不关心分区键是否相邻,常见做法是先计算分区键的哈希值,再将它映射到分片。

好的哈希函数可以把倾斜的数据均匀打散。用于分片的哈希函数不必具备密码学强度,例如 MongoDB 使用 MD5,Cassandra 和 ScyllaDB 则使用 Murmur3。

许多编程语言都内置了哈希表的简单哈希函数,但它们未必适合分片,例如 Java 的 Object.hashCOde() 和 Ruby 的 Object#hash 可能 让同一个键在不同进程中得到不同的哈希值,因此不能用于分片。

哈希取模节点数

算出键的哈希值之后,该如何选择存储它的分片?简单的方法是用节点数 取模(许多编程语言使用 % 运算符)hash(key) % n

模 N 方法的问题在于,只要节点数 N 发生变化,大多数键就必须从一个节点移到另一个节点。

图 7-3 通过对键进行哈希并取模节点数来将键分配给节点。更改节点数会导致许多键从一个节点移动到另一个节点

模 N 很容易计算,却会导致极其低效的再平衡,因为大量记录在节点之间进行了不必要的迁移。我们需要一种只移动必要数据的办法。

固定数量的分片

简单而常用的解决方案,是创建远多于节点数的分片,再给每个节点分配多个分片。例如,一个运行在 10 节点集群上的数据库可以从一开始就划分成 1,000 个分片,每个节点分得 100 个。键会存入编号为 hash(key) % 1,000 的分片,而系统另行记录每个分片存放在哪个节点上。

如果向集群加入一个节点,系统可以把现有节点上的一部分分片重新分配给新节点,直到分片再次均匀分布。图 7-4 展示了这一过程。移除节点时,则反向执行同样的操作。

图 7-4 向每个节点有多个分片的数据库集群添加新节点

在这种模型中,只有完整的分片在节点之间移动,成本低于拆分分片。分片的数量不会改变,键所指定的分片也不会改变;唯一改变的是分片所在的节点。这种变更并非即时——在网络上传输大量数据需要时间——所以传输期间发生的读写,仍按原有的分片到节点映射处理。

分片数量通常会选成一个因数很多的数字,使数据集能够均匀分配到多种不同规模的节点集群中,例如不必要求节点数是 2 的幂。甚至还可以照顾集群中的硬件差异:给性能更强的节点分配更多分片,让它们承担更大比例的负载。

Citus(PostgreSQL 的分片层)、Riak、Elasticsearch 和 Couchbase 等系统都采用这种分片方法。只要首次创建数据库时能较准确地估计所需分片数,它就很好用:此后可以轻松增删节点,不过节点数不能超过分片数。

如果发现最初配置的分片数不合适——例如系统规模已经大到所需节点数超过分片数——就必须执行代价高昂的重新分片。这个过程要拆开每个分片、写出新文件,并占用大量额外磁盘空间。有些系统不允许在数据库继续接受写入时重新分片,因此很难在不停机的情况下改变分片数量。

如果数据集总量变化很大(例如开始时很小,随后可能增长许多倍),选择合适的分片数就很困难。由于每个分片包含总数据量的固定比例,其大小会随集群中的数据总量同比增长。分片太大,再平衡和从节点失效中恢复都会十分昂贵;分片太小,又会带来过多管理开销。分片大小不大不小、“恰到好处”时性能最佳,但在分片数固定而数据集大小不断变化时,这一状态很难维持。

按哈希范围分片

如果无法事先预测需要多少分片,最好采用一种能让分片数量轻松适应工作负载的方案。前述键范围分片具备这一性质,但大量写入集中到相邻键时容易形成热点。一种解决办法是将键范围分片与哈希函数结合,使每个分片包含一段 哈希值 范围,而不是一段 键 范围。

一个 16 位哈希函数,它会返回 0 到 65,535 = 2¹⁶ − 1 之间的数(实际使用的哈希通常至少有 32 位)。即使输入键十分相似(例如连续的时间戳),它们的哈希值也会均匀分布在这个范围内。

图 7-5 为每个分片分配连续的哈希值范围

与键范围分片一样,哈希范围分片也可以在分片过大或负载过重时将其拆分。这个操作依然昂贵,但可以按需执行,因此分片数量会随数据量调整,而不是预先固定不变。

它相对于键范围分片的缺点,是无法高效地对分区键执行范围查询,因为范围内的键如今散布在所有分片中。不过,如果键由两列或更多列组成,而分区键只是其中第一列,仍然可以对第二列及之后的列高效执行范围查询:只要范围查询中的所有记录拥有相同分区键,它们就会落在同一个分片中。和关系型数据库B树的组合索引一个道理 多列索引之前前面的相同才能用索引范围查询索引第二列。

YugabyteDB 和 DynamoDB 采用哈希范围分片,MongoDB 也把它作为一种可选方案。Cassandra 和 ScyllaDB 则采用这种方法的一个变体,如 图 7-6 所示:它们把哈希值空间划分成若干范围,范围数与节点数成正比(图 7-6 中每个节点有 3 个范围;实际默认值是 Cassandra 每个节点 8 个、ScyllaDB 每个节点 256 个),各范围之间的边界随机选定。这样有些范围会比其他范围大,但每个节点拥有多个范围之后,这些不均衡往往能相互抵消。

图 7-6 Cassandra 和 ScyllaDB 将可能的哈希值范围(这里是 0–1023)拆成边界随机的连续区间,并为每个节点分配多个区间

添加或移除节点时,系统会相应增删范围边界,并拆分或合并分片。在 图 7-6 的例子中,加入节点 3 之后,节点 1 把自己两个范围中的一部分交给节点 3,节点 2 也把一个范围中的一部分交给节点 3。这样,新节点便能分得大致公平的一份数据,同时避免在节点间传输不必要的数据。

一致性哈希

一致性哈希(consistent hashing)算法是一种哈希函数,它把键映射到指定数量的分片,并满足两个性质:

可以去看一遍哈希环,分布式一致性哈希。

倾斜的工作负载与缓解热点

一致性哈希可以保证键大致均匀地分布到各节点,却不能保证实际负载也同样均匀。如果工作负载高度倾斜——也就是某些分区键下的数据量远大于其他键,或者某些键的请求速率远高于其他键——仍然可能有些服务器不堪重负,另一些服务器却几乎闲置。

例如社交媒体网站,拥有大量粉丝的名人,可能产生针对同一个键的大量读写(可能是用户ID、评论ID、帖子ID等等)。

这种情况需要更加灵活的分片策略。如果系统按键范围(或哈希范围)定义分片,就可以把一个热键单独放进一个分片,甚至给它分配一台专用机器。

可以在应用层补偿倾斜。例如,如果已知某个键非常热,一种简单办法是在键的开头或末尾添加随机数。只需两位十进制随机数,就能把针对该键的写入均匀拆成 100 个不同的键,让它们分布到不同分片。不过,写入分散到不同键之后,读取就得付出额外代价:必须从全部 100 个键读取数据,再把结果合并起来。

热键分散后,每个分片承受的读取量并没有减少,降低的只有写入负载。这种技术还需要额外的记录工作:只有少数热键值得添加随机数;对于写入吞吐量很低的绝大多数键,这样做只会徒增开销。因此,还需要记录哪些键已被拆分,并设计一个流程,把普通键转换成需要特殊管理的热键。

负载还会随时间变化,使问题更加复杂。例如,某条突然爆火的社交媒体帖子可能连续几天承受很高负载,之后又很快归于平静。此外,有些键是写入热点,有些则是读取热点,二者需要采用不同的处理策略。

一些系统(尤其是面向大规模场景设计的云服务)能够自动处理热分片;例如,Amazon 把相关机制称为 热度管理(heat management) 或 自适应容量(adaptive capacity)。

运维:自动手动再平衡

再平衡自动还是手动进行?

有些系统无需人工介入,会自动决定何时拆分分片、何时把分片从一个节点迁移到另一个节点;另一些系统则要求管理员显式配置分片。两者之间也有折中方案:例如,Couchbase 和 Riak 会自动生成建议的分片分配,但必须由管理员确认提交后才会生效。

全自动再平衡很方便,因为日常维护所需的运维工作更少;这样的系统甚至可以自动伸缩,以适应工作负载的变化。DynamoDB 等云数据库宣称,能够在几分钟内自动增删分片,应对负载的大幅升降。

自动分片管理也可能难以预测。再平衡代价很高,因为它要重新路由请求,并在节点间迁移大量数据。如果处理不够谨慎,这一过程可能使网络或节点过载,拖累其他请求的性能。系统在再平衡期间还必须继续处理写入;如果已经接近最大写入吞吐量,分片拆分的速度甚至可能赶不上新写入到达的速度。

这种自动化机制如果再与自动失效检测结合,可能十分危险。假设某个节点过载,暂时无法及时响应请求;其他节点据此断定它已经失效,于是自动对集群进行再平衡,把负载从该节点移走。这会给其他节点和网络施加额外负载,让局面进一步恶化,甚至引发级联失效:其他节点也相继过载,并被错误地判定为已经宕机。

出于这个原因,让人参与再平衡过程是一件好事。这比全自动流程慢,但有助于防止运维意外。

请求路由

现在来看问题:如果想读写某个特定的键,怎样知道应该连接哪个节点,也就是哪个IP地址和端口?

这个问题称为 路由请求(request routing),与 服务发现十分类似。二者最大的区别在于:运行应用代码的服务实例通常是无状态的,负载均衡器可以把请求发给任意实例;而在分片数据库中,某个键的请求只能交给持有该键所在分片副本的节点处理。

请求路由必须了解键到分片、以及分片到节点的映射,有几种方法:

  1. 允许客户端连接任意节点(例如通过轮询负载均衡器)。如果该节点恰好持有请求涉及的分片,就直接处理请求;否则,它把请求转发给正确的节点,收到响应后再转交给客户端。
  2. 客户端的所有请求都先发送到一个路由层,由路由层判断哪个节点应当处理每个请求,再相应地转发。路由层本身并不处理请求,只充当一个能够感知分片的负载均衡器。
  3. 让客户端了解分片方式以及分片到节点的分配关系。这样,客户端无须经过任何中间层,就能直接连接到正确的节点。
图 7-7 将请求路由到正确节点的三种不同方式

在所有情况下,有一些关键问题:

许多分布式数据系统依靠 ZooKeeper、etcd 等独立协调服务来记录分片分配,服务使用共识算法实现容错并防止脑裂。每个节点都在 ZooKeeper 中注册,ZooKeeper 维护分片到节点的权威映射;路由层或能感知分片的客户端等其他参与者,可以订阅 ZooKeeper 中的信息。只要分片易主,或有节点加入、退出,ZooKeeper 就会通知路由层,使其路由信息保持最新。

图 7-8 使用 ZooKeeper 跟踪分片到节点的分配

例如,HBase 和 SolrCloud 使用 ZooKeeper 管理分片分配,Kubernetes 使用 etcd 记录每个服务实例的运行位置。MongoDB 的架构与之相似,不过它依靠自有的 配置服务器(config server)实现,并以 mongos 守护进程作为路由层。Kafka、YugabyteDB 和 TiDB 则使用内置的 Raft 共识协议实现这项协调功能。

Cassandra、ScyllaDB 和 Riak 采用另一种办法:节点之间通过 流言协议(gossip protocol) 传播集群状态的变化。它提供的一致性比共识协议弱得多,因而可能出现脑裂,使集群的不同部分对同一个分片持有不同的节点分配。无主数据库可以容忍这种情况,因为它们本就只提供较弱的一致性保证。

无论使用路由层还是把请求发送给随机节点,客户端仍然要先找到可供连接的 IP 地址。IP 地址的变化没有分片到节点的分配那么频繁,因此通常用 DNS 就足够了。

以上请求路由主要关注如何为单个键找到对应分片,这最适用于分片的 OLTP 数据库。分析型数据库通常也会分片,但其查询执行方式截然不同:查询一般不是在单个分片中执行,而是要并行聚合并连接来自许多分片的数据。

分片与二级索引

上面分片方案,要求客户端知道待访问记录的分区键。这在键值数据模型中最容易做到:分区键是主键的第一部分(或整个主键),因此可以据此确定分片,并把读写请求路由到负责该键的节点。

涉及二级索引时,情况会复杂得多,二级索引通常不能唯一标识一条记录,而是用来搜索某个特定值出现在哪里:例如,查找用户 123 的所有操作、所有包含单词 hogwash 的文章,或所有颜色为 red 的汽车。

键值存储通常没有二级索引,但它是关系数据库的基础能力,在文档数据库中也十分常见,更是 Solr、Elasticsearch 等全文检索引擎的 立身之本。二级索引的问题在于,它无法干净利落地映射到分片。对带有二级索引的数据库进行分片,主要有两种办法:本地索引和全局索引

本地二级索引

假设运营一个二手车交易网站,每条车辆信息都有唯一ID,并以该ID作为分区键进行分片,(例如,ID 0 到 499 归分片 0,ID 500 到 999 归分片 1,依此类推)。

如果要让用户搜索车辆,并按颜色与品牌筛选,就需要在 color 和 make 上建立二级索引(在文档数据库中它们是字段,在关系数据库中则是列)。声明索引后,数据库会自动维护它。例如,每增加一辆红色汽车,所在分片就会自动把它的 ID 加入索引条目 color:red 对应的 ID 列表。,这种 ID 列表也称为 倒排列表(postings list)

图 7-9 本地二级索引:每个分片只索引其自己分片内的记录

每个分片独立,各自维护自己的二级索引,只覆盖本分片中的记录,不关心其他分片存了什么数据,每次写入数据库 添加、删除、更新记录 只需处理包含该记录的分片。这种二级索引称为 本地索引(local index) ;在信息检索领域,它也称为 按文档分区的索引(document-partitioned index)

读取本地二级索引时,如果已经知道目标记录的分区键,就只需在对应分片上搜索。如果只想获得 部分 结果而不要求全部,也可以把请求发给任意分片。

如果需要全部结果,又事先不知道这些记录的分区键,就必须把查询发送到所有分片,再合并返回结果,因为匹配的记录可能散布在每个分片中。

这种查询分片数据库的方式,会让二级索引上的读取查询变得相当昂贵。即使并行查询所有分片,也很容易出现尾部延迟放大。还会限制应用的可伸缩性,增加分片能容纳更多数据,但如果每次查询仍要由所有分片处理,查询吞吐量并不会随之提高。

这种本地二级索引,MongoDB、Riak、Cassandra、Elasticsearch、SolrCloud 和 VoltDB 都有在用。

全局二级索引

可以构建一个覆盖所有分片数据的 全局索引(global index) 。不能只把索引存放在单个节点上,否则很可能成为瓶颈,使分片失去意义。因此全局索引本身也必须分片,但可以采用与主键索引不同的分片方式。

图 7-10 全局二级索引反映来自所有分片的数据,并且本身按索引值进行分片

这种索引也称为 按词项分区(term-partitioned)

先用二级索引查找相应的ID,再根据ID找记录。如果列表很长,通过网络传输它们再计算交集,速度可能就很慢。

全局二级索引的另一个难题,是写入比本地索引复杂:写入一条记录可能影响索引的多个分片(文档中的每个词项都可能位于不同分片)。因此,二级索引很难与底层数据保持同步。一种办法是使用分布式事务,以原子方式更新保存主记录的分片及其二级索引分片。

CockroachDB、TiDB 和 YugabyteDB 都使用全局二级索引;DynamoDB 则同时支持本地和全局二级索引。在 DynamoDB 中,写入会异步反映到全局索引,因此从全局索引读到的结果可能是陈旧的(类似于 “复制延迟的问题”)。尽管如此,如果读取吞吐量高于写入吞吐量,而且倒排列表不太长,全局索引仍然很有用。