)
数据分区(Data Sharding)面对一些海量数据集或非常高的查询压力使用数据复制技术还不够我们还需要将数据拆分成为分区也称为分片。这里我们所讨论的分区在不同系统有着不同的称呼例如它对应于MongoDB, Elasticsearch和SolrCloud中的shard, HBase的region, Bigtable中的tablet, Cassandra和Riak中的vnode以及Couchbase中的vBucket。总体而言分区是最普遍的术语。分区通常是这样定义的即每一条数据或者每条记录每行或每个文档只属于某个特定分区。实现的方法有多种稍后将逐一介绍。实际上每个分区都可以视为一个完整的小型数据库虽然数据库可能存在一些跨分区的操作。**采用数据分区的主要目的是提高可扩展性。**不同的分区可以放在一个无共享集群无共享架构也称为水平扩展。当采用这种架构时运行数据库软件的机器或者虚拟机称为节点。每个节点独立使用本地的CPU,内存和磁盘。节点之间的所有协调通信等任务全部运行在传统网络以太网之上且核心逻辑主要依靠软件来实现。的不同节点上。这样一个大数据集可以分散在更多的磁盘上查询负载也随之分布到更多的处理器上。对单个分区进行查询时每个节点对自己所在分区可以独立执行查询操作因此添加更多的节点可以提高查询吞吐量。超大而复杂的查询尽管比较困难但也可能做到跨节点的并行处理。分区数据库最初在20世纪80年代由Teradata和Tandem NonStop SQL等率先推出最近又被一些NoSQL数据库和基于Hadoop的数据仓库重视起来。这些系统有些是为事务型负载设计的有些是分析型负载设计的。二者的差异会显著影响系统的优化策略然而分区技术的基本原理则可以普遍适用。本章我们将首先介绍切分大型数据集的若干方案讨论数据索引如何影响分区。接下来讨论分区的再平衡这对动态添加或删除节点非常重要。最后我们将介绍数据库如何将请求路由到正确的分区并执行查询。数据分区与数据复制分区通常与复制结合使用即每个分区在多个节点都存有副本。这意味着某条记录属于特定的分区而同样的内容会保存在不同的节点上以提高系统的容错性。一个节点上可能存储了多个分区。下图展示了主从复制模型与分区组合使用时数据的分布情况。_____________________________________________________________________________ | | | [ 节点 1 ] [ 节点 2 ] | | ┌──────────────────────────┐ ┌─────────────────────┐| | │ ┌──────────────────────┐ │ │ ┌─────────────────┐ │| | │ │ 分区 1 │ │ │ │ 分区 4 │ │| | │ │ (主副本) │ │ │ │ (从副本) │ │| | │ └──────────┬───────────┘ │ │ └────────▲────────┘ │| | └────────────│─────────────┘ └──────────│──────────┘| | │ │ | | │ (P1 复制流) │ | | ├──────────────────────────────┐ │ | | │ │ (P4 复制流) | | │ ┌─────────────┼──────────────────┤ | | │ │ │ │ | | ┌────────────▼────────────────▼─┐ ┌────▼──────────────────┴────────┐ | | │ ┌───────────┐ ┌────────────┐ │ │ ┌───────────┐ ┌─────────────┐ │ | | │ │ 分区 1 │ │ 分区 4 │ │ │ │ 分区 1 │ │ 分区 4 │ │ | | │ │ (从副本) │ │ (从副本) │ │ │ │ (从副本) │ │ (主副本) │ │ | | │ └───────────┘ └────────────┘ │ │ └───────────┘ └──────▲──────┘ │ | | └───────────────────────────────┘ └───────────────────────│────────┘ | | [ 节点 3 ] [ 节点 4 ] │ | | │ | | 分区 4 的写请求 ───────┘ | | │ | | ○ | | ---- (用户) | | / \ | | | | ───────────► 每分区对应的复制数据流 | |_____________________________________________________________________________|上图组合使用复制和分区每个节点同时充当某些分区的主副本和其他分区的从副本每个分区都有自己的主副本例如被分配给某节点而从副本则分配在其他一些节点。一个节点可能即是某些分区的主副本同时又是其他分区的从副本。所有复制相关的原理同样适用于对分区数据的复制。考虑到分区方案的选择通常独立于复制因此本章将力求简洁而省略与复制相关的内容。键-值数据的分区假设面临海量数据现在需要切分它们那么该如何决定哪些记录放在哪些节点上呢分区的主要目标是将数据和查询负载均匀分布在所有节点上。如果节点平均分担负载那么理论上10个节点应该能够处理10倍的数据量和10倍于单个节点的读写吞吐量忽略复制。而如果分区不均匀则会出现某些分区节点比其他分区承担更多的数据量或查询负载称之为倾斜。倾斜会导致分区效率严重下降在极端情况下所有的负载可能会集中在一个分区节点上这就意味着10个节点9个空闲系统的瓶颈在最繁忙的那个节点上。这种负载严重不成比例的分区即成为系统热点。避免热点最简单的方法是将记录随机分配给所有节点上。这种方也可以比较均匀地分布数据但是有一个很大的缺点当试图读取特定的数据时没有方法知道数据保存在哪个节点上所以不得不并行查询所有节点。可以改进上述方法。现在我们假设数据是简单的键-值数据模型这意味着总是可以通过关键字来访问记录。例如像一个纸质百科全书可以通过标题来查找某一个条目而所有的条目按字母序排序因此可以做到快速查找条目。基于关键字区间分区一种分区方式是为每个分区分配一段连续的关键字或者关键字区间范围以最小值和最大值来指示如百科全书按照关键字区间进行分区。如果知道关键字区间的上下限就可以轻松确定哪个分区包含这些关键字。如果还知道哪个分区分配在哪个节点上就可以直接向该节点发出请求。关键字的区间段不一定非要均匀分布这主要是因为数据本身可能就不均匀。为了更均匀地分布数据分区边界理应适配数据本身的分布特征。分区边界可以由管理员手动确定或者由数据库自动选择我们将在本章后面的“分区再平衡”中更详细地讨论。采用这种分区策略的系统包括Bigtable, Bigtable的开源版本HBase, RethinkDB和2.4版本之前MongoDB。每个分区内可以按照关键字排序保存参阅“SSTables和LSM Trees”。这样可以轻松支持区间查询即将关键字作为一个拼接起来的索引项从而一次查询得到多个相关记录参阅“多列索引”。例如对于一个保存网络传感器数据的应用系统选择测量的时间戳年月日时分秒作为关键字此时区间查询会非常有用它可以快速获得某个月份内的所有数据。然而基于关键字的区间分区的缺点是某些访问模式会导致热点。如果关键字是时间戳则分区对应于一个时间范围例如每天一个分区。然而当测量数据从传感器写入数据库时所有的写入操作都集中在同一个分区即当天的分区这会导致该分区在写入时负载过高而其他分区始终处于空闲状态。为了避免上述问题需要使用时间戳以外的其他内容作为关键字的第一项。例如可以在时间戳前面加上传感器名称作为前缀这样首先由传感器名称然后按时间进行分区。假设同时有许多传感器处于活动状态则写入负载最终会比较均匀地分布在多个节点上。接下来当需要获取一个时间范围内、多个传感器的数据时可以根据传感器名称各自执行区间查询。基于关键字哈希值分区对于上述数据倾斜与热点问题许多分布式系统采用了基于关键字哈希函数的方式来分区。一个好的哈希函数可以处理数据倾斜并使其均匀分布。例如一个处理字符串的32位哈希函数当输入某个字符串它会返回个0和$ 2^32-1 $之间近似随机分布的数值。即使输入的字符串非常相似返回的哈希值也会在上述数字范围内均匀分布。用于数据分区目的的哈希函数不需要在加密方面很强例如Cassandra和MongoDB使用MDS, Voldemort使用Fowler-Noll-Vo函数。许多编程语言也有内置的简单哈希函数主要用于哈希表但是要注意这些内置的哈希函数可能井不适合分区例如Java的Object.hashCode和Ruby的Object#hash同一个键在不同的进程中可能返回不同的哈希值。一且找到合适的关键字哈希函数就可以为每个分区分配一个哈希范围而不是直接作用于关键字范围关键字根据其哈希值的范围划分到不同的分区中。如下图所示____________________________________________________________________________________________ | | | 2014-04-19 17:08:10 ───► 7,372 ──────┐ | | 2014-04-19 17:08:11 ───► 18,805 ──────┼───────────┐ | | 2014-04-19 17:08:12 ───► 50,537 ──────┼───────────┼─────────────────────┐ | | 2014-04-19 17:08:13 ───► 31,579 ──────┼───────────┼───────────┐ │ | | 2014-04-19 17:08:14 ───► 62,253 ──────┼───────────┼───────────┼─────────┼────┐ | | 2014-04-19 17:08:15 ───► 24,510 ──────┼────────┐ │ │ │ │ | | 这里取 MD5 │ │ │ │ │ │ | | 哈希值的头 │ │ │ │ │ │ | | 2 个字节 │ │ │ │ │ │ | | ▼ ▼ ▼ ▼ ▼ ▼ | | ┌─────┬─────┬─────┬─────┬─────┬─────┬─────┬─────┐ | | │ p0 │ p1 │ p2 │ p3 │ p4 │ p5 │ p6 │ p7 │ | | └─────┴─────┴─────┴─────┴─────┴─────┴─────┴─────┘ | | 0 16,383 32,767 49,151 65,535 | |_____________________________________________________________________________________________|这种方法可以很好地将关键字均匀地分配到多个分区中。分区边界可以是均匀间隔也可以是伪随机选择在这种情况下该技术有时被称为一致性哈希。一致性哈希由Karger等人首先提出是一种平均分配负载的方法最初用于内容分发网络 CDN 等互联网缓存系统。它采用随机选择的分区边界未规避中央控制或分布式共识。请注意此处的一致性与副本一致性或ACID的一致性没有任何关联它只描述了数据动态平衡的一种方法。正如后面“分区再平衡” 一节将介绍的这种特殊的分区方法对于数据库实际效果并不是很好所以目前很少使用虽然某些数据库的义档仍采用一效性哈希的术语但其实并不准确。然而通过关键字哈希进行分区我们丧失了良好的区间查询特性。即使关键字相邻但经过哈希之后会分散在不同的分区中区间查询就失去了原有的有序相邻的特性。在MongoDB中如果启用了基于哈希的分片模式则区间查询会发送到所有的分区上而Riak、CouchbaseJ和Voldemort干脆就不支持关键字上的区间查询。Cassandra则在两种分区策略之间做了一个折中。Cassandra中的表可以声明为由多个列组成的复合主键。复合主键只有第一部分可用于哈希分区而其他列则用作组合索引来对Cassandra SSTable中的数据进行排序。因此它不支持在第一列上进行区间查询但如果为第一列指定好了固定值可以对其他列执行高效的区间查询。组合索引为一对多的关系提供了一个优雅的数据模型。例如在社交网站上一个用户可能会发布很多消息更新。如果更新的关键字设置为user_id, update_timestamp的组合那么可以有效地检索由某用户在一段时间内所做的所有更新且按时间戳排序。不同的用户可以存储在不同的分区上但是对于某一用户消息按时间戳顺序存储在一个分区上。负载倾斜与热点如前所述基于哈希的分区方法可以减轻热点但无法做到完全避免。一个极端情况是所有的读/写操作都是针对同一个关键字则最终所有请求都将被路由到同一个分区。这种负载或许并不普遍但也并非不可能例如社交媒体网站上一些名人用户有数百万的粉丝当其发布一些热点事件时可能会引发一场访问风暴出现大量的对相同关键字的写操作其中关键字可能是名人的用户ID或者人们正在评论的事件ID。此时哈希起不到任何帮助作用因为两个相同ID的哈希值仍然相同。大多数的系统今天仍然无法自动消除这种高度倾斜的负载而只能通过应用层来减轻倾斜程度。例如如果某个关键字被确认为热点一个简单的技术就是在关键字的开头或结尾处添加一个随机数。只需一个两位数的十进制随机数就可以将关键字的写操作分布到100个不同的关键字上从而分配到不同的分区上。但是随之而来的问题是之后的任何读取都需要些额外的工作必须从所有100个关键字中读取数据然后进行合并。因此通常只对少量的热点关键字附加随机数才有意义而对于写入吞吐量低的绝大多数关键字这些都意味着不必要的开销。此外还需要额外的元数据来标记哪些关键字进行了特殊处理。也许将来某一天数据系统能够自动检测负载倾斜情况然后自动处理这些倾斜的负载。但截至目前仍然需要开发者自己结合应用来综合权衡。分区与二级索引我们之前所讨论的分区方案都依赖于键-值数据模型。键-值模型相对简单即都是通过关键字来访问记录自然可以根据关键字来确定分区并将读写请求路由到负责该关键字的分区上。但是如果涉及二级索引情况会变得复杂。二级索引通常不能唯一标识一条记录而是用来加速特定值的查询例如查找用户123的所有操作找到所有含有hogwash的文章查找所有颜色为红色的汽车等。二级索引是关系数据库的必备特性在文档数据库中应用也非常普遍。但考虑到其复杂性许多键-值存储如HBase和Voldemort并不支持二级索引但其他一些如Riak则开始增加对二级索引的支持。此外二级索引技术也是Solr 和Elasticsearch等全文索引服务器存在之根本。二级索引带来的主要挑战是它们不能规整的地映射到分区中。有两种主要的方法来支持对二级索引进行分区基于文档的分区和基于词条的分区。基于文档分区的二级索引假设有一个销售二手车的网站数据信息如下图所示__________________________________________________ __________________________________________________ | | | | | [ 分区 0 ] | | [ 分区 1 ] | | ┌────────────────────────────────────────────┐ | | ┌────────────────────────────────────────────┐ | | │ 关键字索引 │ | | │ 关键字索引 │ | | │ 191 - {color:red, make:Honda...} │ | | │ 515 - {color:silver, make:Ford...} │ | | │ 214 - {color:black, make:Dodge...} │ | | │ 768 - {color:red, make:Volvo...} │ | | │ 306 - {color:red, make:Ford...} │ | | │ 893 - {color:silver, make:Audi...} │ | | ├────────────────────────────────────────────┤ | | ├────────────────────────────────────────────┤ | | │ 二级索引 (基于文档的分区) │ | | │ 二级索引 (基于文档的分区) │ | | │ color:black - [214] │ | | │ color:black - [] │ | | │ color:red - [191, 306] ◄────┐ │ | | │ color:red - [768] ◄────┐ │ | | │ color:yellow - [] │ │ | | │ color:silver - [515, 893] │ │ | | │ make:Dodge - [214] │ │ | | │ make:Audi - [893] │ │ | | │ make:Ford - [306] │ │ | | │ make:Ford - [515] │ │ | | │ make:Honda - [191] │ │ | | │ make:Volvo - [768] │ │ | | └──────────────────────────────────│─────────┘ | | └──────────────────────────────────│───────────┘ | └─────────────────────────────────────│────────────┘ └─────────────────────────────────────│──────────────┘ │ │ └──────────────────────────┬───────────────────────────┘ │ (对所有分区执行查询 然后合并结果) │ ○ ---- 查找红色的汽车 / \ (用户)每个列表都有一个唯一的文档ID用此ID对数据库进行分区例如ID 0到499归分区0ID 500 到999划为分区l 。现在用户需要搜索汽车可以按汽车颜色和厂商进行过滤所以需要在颜色和制造商上设定二级索引在文档数据库中这些都是字段在关系数据库中则是列。声明这些索引之后数据库会自动创建索引。例如每当一辆红色汽车添加到数据库中数据库分区会自动将其添加到索引条目为“color: red”的文档ID列表中。在这种索引方法中每个分区完全独立各自维护自己的二级索引且只负责自己分区内的文档而不关心其他分区中数据。每当需要写数据库时包括添加删除或更新文档等只需要处理包含目标文档ID的那一个分区。因此文档分区索引也被称为本地索引而不是全局索引后者将在本章后面介绍。但读取时需要注意除非对文档ID做了特别的处理否则不太可能所有特定颜色或特定品牌的汽车都放在一个分区中例如上图中红色汽车就出现在分区0和分区l中。因此如果想要搜索红色汽车就需要将查询发送到所有的分区然后合并所有返回的结果。这种查询分区数据库的方法有时也称为分散/聚集显然这种二级索引的查询代价高昂。即使采用了并行查询也容易导致读延迟显著放大。尽管如此它还是广泛用于实践MongoDB、Riak、Cassandra、Elasticsearch、SolrCloud和VoltDB都支持基于文档分区二级索引。大多数数据库供应商都建议用户自己来构建合适的分区方案尽量由单个分区满足二级索引查询但现实往往难以如愿尤其是当查询中可能引用多个二级索引时例如同时指定颜色和制造商两个条件。基于词条的分区另一种方法我们可以对所有的数据构建全局索引而不是每个分区维护自己的本地索引。而且为避免成为瓶颈不能将全局索引存储在一个节点上否则就破坏了设计分区均衡的目标。所以全局索引也必须进行分区且可以与数据关键字采用不同的分区策略。__________________________________________________ __________________________________________________ | | | | | [ 分区 0 ] | | [ 分区 1 ] | | ┌────────────────────────────────────────────┐ | | ┌────────────────────────────────────────────┐ | | │ 关键字索引 │ | | │ 关键字索引 │ | | │ 191 - {color:red, make:Honda...} │ | | │ 515 - {color:silver, make:Ford...} │ | | │ 214 - {color:black, make:Dodge...} │ | | │ 768 - {color:red, make:Volvo...} ─┼──┐ | │ 306 - {color:red, make:Ford...} │ | | │ 893 - {color:silver, make:Audi...} │ | | ├────────────────────────────────────────────┤ | | ├────────────────────────────────────────────┤ | | │ 二级索引 (基于词条的分区) │ | | │ 二级索引 (基于词条的分区) │ | | │ color:black - [214] │ | | │ color:silver - [515, 893] │ | | │ color:red - [191, 306, 768] ───┐ │ | | │ color:yellow - [] │ | | │ make:Audi - [893] │ │ | | │ make:Honda - [191] │ | | │ make:Dodge - [214] │ │ | | │ make:Volvo - [768] │ | | │ make:Ford - [306, 515] │ | | | | | | | └─────────────────────────────────────│──────┘ | | └────────────────────────────────────────────┘ | └────────────────────────────────────────│─────────┘ └──────────────────────────────────────────────────┘ │ (查询请求只发给对应词条的分区) │ │ │ │ ○ └────── ---- 查找红色的汽车 / \ (用户)以上图为例所有数据分区中的颜色为红色的汽车被收录到在索引color:red中而索引本身也是分区的例如从a到r开始的颜色放在分区0中从s到z的颜色放在分区1中。类似的汽车制造商的索引也被分区两个分区的边界分别是字母f和字母h。我们将这种索引方案称为词条分区它以待查找的关键字本身作为索引。例如颜色color: red。名字词条源于全文索引一种特定类型的二级索引。和前面讨论的方法一样可以直接通过关键词来全局划分索引或者对其取哈希值。直接分区的好处是可以支持高效的区间查询例如查询汽车报价在某个值以上而采用哈希的方式则可以更均匀的划分分区。这种全局的词条分区相比于文档分区索引的主要优点是它的读取更为高效即它不需要采用scatter/gather对所有的分区都执行一遍查询相反客户端只需要向包含词条的那一个分区发出读请求。然而全局索引的不利之处在于写入速度较慢且非常复杂主要因为单个文档的更新时里面可能会涉及多个二级索引而二级索引的分区又可能完全不同甚至在不同的节点上由此势必引入显著的写放大。理想情况下索引应该时刻保持最新即写入的数据要立即反映在最新的索引上。但是对于词条分区来讲这需要一个跨多个相关分区的分布式事务支持写入速度会受到极大的影响所以现有的数据库都不支持同步更新二级索引参阅第7 章和第9章。实践中对全局二级索引的更新往往都是异步的也就意味着如果在写入之后马上去读索引那么刚刚发生的更新可能还没有反映在索引中。例如 Amazon DynamoDB的二级索引通常可以在1秒之内完成更新但当底层设施出现故障时也有可能需要等待很长的时间。其他使用全局索引的系统还包括Riak的搜索功能和Oracle数据仓库后者允许用户来选择是使用本地索引还是全局索引。分区再平衡随着时间的推移数据库可能总会出现某些变化(1) 查询压力增加因此需要更多的CPU来处理负载。(2) 数据规模增加因此需要更多的磁盘和内存来存储数据。(3) 节点可能出现故障因此需要其他机器来接管失效的节点。所有这些变化都要求数据和请求可以从一个节点转移到另一个节点。这样一个迁移负载的过程称为再平衡或者动态平衡。无论对于哪种分区方案分区再平衡通常至少要满足(1) 平衡之后负载、数据存储、读写请求等应该在集群范围更均匀地分布。(2) 再平衡执行过程中数据库应该可以继续正常提供读写服务。(3) 避免不必要的负载迁移以加快动态再平衡并尽量减少网络和磁盘I/O影响。动态再平衡的策略将分区对应到节点上存在多种不同的分配策略这里逐一介绍为什么不用取模对节点数取模方法的问题是如果节点数N发生了变化会导致很多关键字需要从现有的节点迁移到另一个节点。例如假设hash(key) 123456假定最初是10个节点那么这个关键字应该放在节点6123456 mod 10 6当节点数增加到11时它需要移动到节点3123456 mod 11 3当继续增长到12个节点时又需要移动到节点0123456 mod 12 0。这种频繁的迁移操作大大增加了再平衡的成本。因此我们需要一种减少迁移数据的方法。固定数量的分区幸运的是有一个相当简单的解决方案首先创建远超实际节点数的分区数然后为每个节点分配多个分区。例如对于一个10节点的集群数据库可以从一开始就逻辑划分为1000个分区这样大约每个节点承担100个分区。接下来如果集群中添加了一个新节点该新节点可以从每个现有的节点上匀走几个分区直到分区再次达到全局平衡。如果从集群中删除节点则采取相反的均衡措施。选中的整个分区会在节点之间迁移但分区的总数量仍维持不变也不会改变关键字到分区的映射关系。这里唯一要调整的是分区与节点的对应关系。考虑到节点间通过网络传输数据总是需要些时间这样调整可以逐步完成在此期间旧的分区仍然可以接收读写请求。原则上也可以将集群中的不同的硬件配置因素考虑进来即性能更强大的节点将分配更多的分区从而分担更多的负载。目前Riak、Elasticsearch、Couchbase和Voldemort都支持这种动态平衡方法。使用该策略时分区的数量往往在数据库创建时就确定好之后不会改变。原则上也可以拆分和合并分区稍后介绍但固定数量的分区使得相关操作非常简单因此许多采用固定分区策略的数据库决定不支持分区拆分功能。所以在初始化时已经充分考虑将来扩容增长的需求未来可能拥有的最大节点数设置一个足够大的分区数。而每个分区也有些额外的管理开销选择过高的数字可能会有副作用。如果数据集的总规模高度不确定或可变例如开始非常小但随着时间的推移可能会变得异常庞大此时如何选择合适的分区数就有些困难。每个分区包含的数据量的上限是固定的实际大小应该与集群中的数据总量成正比。如果分区里的数据量非常大则每次再平衡和节点故障恢复的代价就很大但是如果一个分区太小就会产生太多的开销。分区大小应该“恰到好处”不要太大 也不能过小如果分区数量固定了但总数据量却不确定就难以达到一个最佳取舍点。动态分区对于采用关键字区间分区的数据库如果边界设置有问题最终可能会出现所有数据都挤在一个分区而其他分区基本为空那么设定固定边界、固定数量的分区将非常不便而手动去重新配置分区边界又非常繁琐。因此一些数据库如HBase和RethinkDB等采用了动态创建分区。当分区的数据增长超过一个可配的参数阈值HBase上默认值是lOGB它就拆分为两个分区每个承担一半的数据量。相反如果大量数据被删除并且分区缩小到某个阈值以下则将其与相邻分区进行合井。该过程类似于B树的分裂操作。每个分区总是分配给一个节点而每个节点可以承载多个分区这点与固定数量的分区一样。当一个大的分区发生分裂之后可以将其中的一半转移到其他某节点以平衡负载。对于HBase 分区文件的传输需要借助HDFS底层分布式文件系统。动态分区的一个优点是分区数量可以自动适配数据总量。如果只有少量的数据少量的分区就足够了这样系统开销很小如果有大量的数据每个分区的大小则被限制在一个可配的最大值。但是需要注意的是对于一个空的数据库 因为没有任何先验知识可以帮助确定分区的边界所以会从一个分区开始。可能数据集很小但直到达到第一个分裂点之前所有的写入操作都必须由单个节点来处理 而其他节点则处于空闲状态。为了缓解这个问题HBase和MongoDB允许在一个空的数据库上配置一组初始分区这被称为预分裂。对于关键字区间分区预分裂要求已经知道一些关键字的分布情况。动态分区不仅适用于关键字区间分区 也适用于基于哈希的分区策略。MongoDB从版本2.4开始同时支持二者并且都可以动态分裂分区。按节点比例分区采用动态分区策略拆分和合并操作使每个分区的大小维持在设定的最小值和最大值之间因此分区的数量与数据集的大小成正比关系。另一方面对于固定数量的分区方式其每个分区的大小也与数据集的大小成正比。两种情况分区的数量都与节点数无关。Cassandra和Ketama则采用了第三种方式使分区数与集群节点数成正比关系。换句话说每个节点具有固定数量的分区。此时当节点数不变时每个分区的大小与数据集大小保持正比的增长关系当节点数增加时分区则会调整变得更小。较大的数据量通常需要大量的节点来存储因此这种方法也使每个分区大小保持稳定。当一个新节点加入集群时它随机选择固定数量的现有分区进行分裂然后拿走这些分区的一半数据量将另一半数据留在原节点。随机选择可能会带来不太公平的分区分裂但是当平均分区数量较大时Cassandra默认情况下每个节点有256个分区新节点最终会从现有节点中拿走相当数量的负载。Cassandra在3.0时推出了改进算法可以避免上述不公平的分裂。随机选择分区边界的前提要求采用基于哈希分区可以从哈希函数产生的数字范围里设置边界。这种方法也最符合本章开头所定义一致性哈希。一些新设计的哈希函数也可以以较低的元数据开销达到类似的效果。自动与手动再平衡操作动态平衡另一个重要问题我们还没有考虑它是自动执行还是手动方式执行全自动式的再平衡即由系统自动决定何时将分区从一个节点迁移到另一个节点不需要任何管理员的介入与纯手动方式即分区到节点的映射由管理员来显式配置之间可能还有一个过渡阶段。例如Couchbase, Riak和Voldemort会自动生成一个分区分配的建议方案但需要管理员的确认才能生效。全自动式再平衡会更加方便它在正常维护之外所增加的操作很少。但是也有可能出现结果难以预测的情况。再平衡总体讲是个比较昂贵的操作它需要重新路由请求并将大量数据从一个节点迁移到另一个节点。万一执行过程中间出现异常会使网络或节点的负载过重并影响其他请求的性能。将自动平衡与自动故障检测相结合也可能存在一些风险。例如假设某个节点负载过重对请求的响应暂时受到影响而其他节点可能会得到结论该节点已经失效接下来激活自动平衡来转移其负载。客观上这会加重该节点、其他节点以及网络的负荷可能会使总体情况变得更槽甚至导致级联式的失效扩散。出于这样的考虑让管理员介入到再平衡可能是个更好的选择。它的确比全自动过程响应慢一些但它可以有效防止意外发生。请求路由现在我们已经将数据集分布到多个节点上但是仍然有一个悬而未决的问题当客户端需要发送请求时如何知道应该连接哪个节点如果发生了分区再平衡分区与节点的对应关系随之还会变化。为了回答该问题我们需要一段处理逻辑来感知这些变化并负责处理客户端的连接例如想要读/写关键字“foo”需要连接哪个IP地址和哪个端口号。这其实属于一类典型的服务发现问题服务发现并不限于数据库任何通过网络访问的系统都有这样的需求尤其是当服务目标支持高可用时在多台机器上有冗余配置。许多公司已经开发了自己的内部服务发现工具其中很多已经开源。概括来讲这个问题有以下几种不同的处理策略如下所示的三种情况(1) 允许客户端链接任意的节点例如采用循环式的负载均衡器。如果某节点恰好拥有所请求的分区则直接处理该请求否则将请求转发到下一个合适的节点接收答复并将答复返回给客户端。(2) 将所有客户端的请求都发送到一个路由层由后者负责将请求转发到对应的分区节点上。路由层本身不处理任何请求它仅充一个分区感知的负载均衡器。(3) 客户端感知分区和节点分配关系。此时客户端可以直接连接到目标节点而不需要任何中介。不管哪种方法核心问题是作出路由决策的组件可能是某个节点路由层或客户端如何知道分区与节点的对应关系以及其变化情况这其实是一个很有挑战性的问题所有参与者都要达成共识这一点很重要。否则请求可能被发送到错误的节点而没有得到正确处理。分布式系统中有专门的共识协议算法但通常难以正确实现。许多分布式数据系统依靠独立的协调服务如ZooKeeper跟踪集群范围内的元数据。每个节点都向ZooKeeper中注册自己ZooKeeper维护了分区到节点的最终映射关系。其他参与者如路由层或分区感知的客户端可以向Zoo Keeper订阅此信息。一旦分区发生了改变或者添加、删除节点ZooKeeper就会主动通知路由层这样使路由信息保持最新状态。例如HBaseSolrCloud和Kafka使用ZooKeeper来跟踪分区分配情况。MongoDB有类似的设计但它依赖于自己的配置服务器和mangos守护进程来充当路由层。Cassandra和Riak则采用了不同的方法它们在节点之间使用gossip协议来同步群集状态的变化。请求可以发送到任何节点由该节点负责将其转发到目标分区节点。这种方式增加了数据库节点的复杂性但是避免了对ZooKeeper之类的外部协调服务的依赖。Couchbase并不支持自动再平衡功能这简化了设计。它通过配置一个名为moxi的路由选择层向集群节点学习最新的路由变化。当使用路由层或随机选择节点发送请求时客户端仍然需要知道目标节点的IP地址。IP地址的变化往往没有分区————节点变化那么频繁采用DNS通常就足够了。并行查询执行到目前为止我们只关注了读取或写入单个关键字这样简单的查询对于文档分区的二级索引里面要求分散/聚集查询。这基本上也是大多数NoSQL分布式数据存储所支持的访问类型。然而对于大规模并行处理massively parallel processing, MPP这一类主要用于数据分析的关系数据库在查询类型方面要复杂得多。典型的数据仓库查询包含多个联合、过滤、分组和聚合操作。MPP查询优化器会将复杂的查询分解成许多执行阶段和分区以便在集群的不同节点上并行执行。尤其是涉及全表扫描这样的查询操作可以通过并行执行获益颇多。数据仓库中快速并行执行查询可以作为单独的话题。小结本章我们探讨了将大规模数据集划分成更小子集的多种方法。数据量如果太大单台机器进行存储和处理就会成为瓶颈因此需要引入数据分区机制。分区的目地是通过多台机器均匀分布数据和查询负载避免出现热点。这需要选择合适的数据分区方案在节点添加或删除时重新动态平衡分区。我们讨论了两种主要的分区方法(1) 基于关键字区间的分区。先对关键字进行排序每个分区只负责一段包含最小到最大关键字范围的一段关键字。对关键字排序的优点是可以支持高效的区间查询但是如果应用程序经常访问与排序一致的某段关键字就会存在热点的风险。采用这种方法当分区太大时通常将其分裂为两个子区间从而动态地再平衡分区。(2) 哈希分区。将哈希函数作用于每个关键字每个分区负责一定范围的哈希值。这种方法打破了原关键字的顺序关系它的区间查询效率比较低但可以更均匀地分配负载。采用哈希分区时通常事先创建好足够多但固定数量的分区让每个节点承担多个分区当添加或删除节点时将某些分区从一个节点迁移到另一个节点也可以支持动态分区。混合上述两种基本方法也是可行的例如使用复合键键的一部分来标识分区而另一部分来记录排序后的顺序。我们还讨论了分区与二级索引二级索引也需要进行分区有两种方法(1) 基于文档来分区二级索引本地索引。二级索引存储在与关键字相同的分区中这意味着写入时我们只需要更新一个分区但缺点是读取二级索引时需要在所有分区上执行scatter/gather。(2) 基于词条来分区二级索引全局索引。它是基于索引的值而进行的独立分区。二级索引中的条目可能包含来自关键字的多个分区里的记录。在写入时不得不更新二级索引的多个分区但读取时则可以从单个分区直接快速提取数据。最后我们讨论了如何将查询请求路由到正确的分区包括简单的分区感知负载均衡器以及复杂的并行查询执行引擎。理论上每个分区基本保持独立运行这也是为什么我们试图将分区数据库分布、扩展到多台机器上。但是如果写入需要跨多个分区情况就会格外复杂例如如果其中一个分区写入成功但另一个发生失败。参考https://www.digitalocean.com/community/tutorials/understanding-database-sharding Understanding Database Sharding