一个中小型分布式KV数据库的设计

今天和大家分享的是一个中小型分布式KV数据库的设计,数据容量定位于千万至亿级,因为这个级别可以满足大多数中小型互联网企业的存储需求,设计和开发者可以腾出手来,把高可用和高可运维做为开发的核心。

为什么十亿或百亿级的分布式系统不好做?

分布式存储的思想就是切分,把一个大型的系统按照一定的组织关系切分成若干小块,每一小块存储一部分数据。对存储分块来说,分10块和分1000块理论上对于技术实现基本没有区别,主要是运营期的维护,机器多了、数据量大了各种问题都会放大,因而很多大型互联网巨头在基础支撑开发和运维方面投入都非常大,高可用和高可运维非常重要。

还有,如何组织这些分块存储,需要有一个映射关系,专业的说就是索引。对于一个亿级以内的系统,索引通常可以保存在一个节点或者说单台主机,并且做到一条记录一个索引,不用考虑其他的分片规则。目前主要有两种分片规则:哈希和范围,无论是哪一种,后期数据规模大了需要迁移或者拆分的时候,一定会施加很多约束,这个暂且不讲吧。我们先粗略计算一下,如果采用灵活的一对一索引,1亿条记录需要的存储容量。以每条索引占100字节计,1亿条索引就是100 * 100M = 10G,这无论是做硬盘索引还是内存索引都是可以接受的。当然考虑到性能,通常会把索引加载到内存中。

现在让我们再考虑下百亿级系统,索引无论如何是单台主机无法放得下的,这时有两种方向可以考虑:一是压缩索引,即前面讲的哈希和范围,二是索引也分块,按一定的规则把索引分布存储在不同的索引服务器中,很显然这个问题又复杂化了。

业界内评价高可用系统有一个著名的CAP理论,即一致性、可用性、分区容灾。前面我们一直在讲可用性,一致性应该也涵盖了一部分,但分区容灾需要以多副本为前提,也就是说无论是数据和索引都需要多个备份。个人认为,开发一个稳定的系统,光数据分布和备份已是十分不易,若索引也分布和备份,开发量和复杂度可想而知,故本文暂不考虑十亿级系统。

设计指标:

在设计系统之前,我们先大致定义下系统的设计指标:

存储容量:5000W级
并发访问:5000次/S
系统扩容:弹性扩充,自动均衡
故障维护:低度人工介入
机器要求:10台以内

关于5000W级记录所需要的存储容量,这个跟所存储的数据类型有关,如果是普通的用户记录数据,按每条10K计,需要的硬盘为500G,而如果是存储图片或小视频,以每个10M计,需要的硬盘为500T。两者对于机器的要求区别还是非常大的,相差了整整1000倍。

关于并发访问的性能,同样的这里我只能暂且以用户记录来估计,若存储服务器分布于3台机器,每台服务器可平均负载2000次/S,并且支持存储服务器弹性扩充性能。但如果是图片或小视频的访问,访问一次需要的时间和流量以及系统负载都非常大,如果不使用CDN的话,广域网纯粹依靠数据库系统本身,应该性能非常低。

关于系统扩容,系统需要考虑单机性能瓶劲,当存储或CPU出现负载不足时,支持扩充新的机器加入以扩展性能。一个低可维护的系统,需要停止业务或者是部分停机,由运维人员手动迁移数据到新的机器上,然后再恢复业务。我们这里的设计要求是,运维人员只需要简单的对机器进行配置,剩下的数据迁移工作由系统自动完成,整个过程不停机,用户访问不受影响。

关于故障维护,出现单机故障,在人工不干预的情况下,仍能提供服务。单机故障期间,其业务访问由其他主机接管,故障恢复后,数据同步自动进行,且基本不影响正常访问。

关于机器要求评估,考虑到分区容灾,以存储普通的记录计,暂定10台以内,具体布署要看整体的系统设计情况。

系统设计:

总体架构:

图示如下:


image.png
注册服务器:

注册服务器主要提供整个分布式系统的服务器感知功能,它的功能包括:
a) 管理服务器注册
b) 维护活跃的服务器列表
c) 提供服务器列表的查询

技术选型:
a) 选择开源软件,如zookeeper、consul、etcd,不同的软件对开发的技术栈和标准不一样,需要架构师根据自己团队的实际情况进行选择。
b) 自研,其实就是造轮子,这对于研发实力和资金雄厚的团队可以考虑,好处就是完全Hold得住,缺点也比较明显,短时间内功能和稳定性肯定不如上面的几款知名开源软件。

性能和可靠性:
a) 对于中小型系统来说,通常逻辑服务器的节点数在百以内,性能一般不会成为瓶劲。
b) 可靠性比较重要,通常可以容忍短时间内的故障,这时整个系统不能注册新服务,也不能做服务器存活状态的广播,除注册服务器外的个别的服务器故障也应该在系统的容错范围内,但长时间不恢复的话一亘故障服务器的数量达到一定级别,系统将不再可用。

索引服务器:

索引服务器在一些系统中又称为主服务器,需考虑以下设计:
a) 维护关键字到数据服务器的映射。
b) 维护数据服务器的负载状态。
c) 对外提供关键字到数据服务器的查询。

技术选型:
通常没有现成的索引服务器可用,必须自研。索引服务器的组成包括两大部分:核心存储和外围逻辑。核心存储的开发比较复杂,可选择成熟的技术方案,如mysql、leveldb等等。外围逻辑包括数据服务器负载的存储,但通常保存在内存中即可。其它就是接口的设计和支持。

性能和可靠性:
a) 索引服务器作为实际数据寻址的关键纽带,通常情况下每一次读写操作,首先都需要向索引服务器发起关键字查询定位,若索引服务器的性能出现瓶劲,则整个系统的性能将受到限制。
b) 可靠性非常重要,基本不容忍故障,出现故障必须能够快速隔离或恢复,这对于严重依赖人为干预的系统来说,基本上达不到要求。所以必须在设计上充分考虑高可用性。
c) 索引数据必须有多个副本,索引服务器也必须以集群形式存在,不必多个节点同时提供服务,但必须做到主点出现故障时,快速自动切换备点。

数据服务器:

数据服务器主要管理核心数据的存储,需考虑以下设计:
a) 管理核心数据的存储。
b) 当数据发生更新时,必须及时同步到索引服务器。
c) 定期向索引服务器报告cpu、disk等负载状态。
d) 支持数据迁移,支持自动弹性扩充。

技术选型:
同索引服务器一样,没有现成的开源可用,必须自研。同样的组成包括核心存储和外围逻辑两部分,核心存储也选择成熟的技术方案,如mysql、leveldb等等。

性能和可靠性:
a) 数据服务器通常是集群整体提供服务,单点的性能通常不会成为整体性能的瓶劲。
b) 数据服务器通常是对等布署,节点之间弱交互,即除了数据迁移时有交互外,通常不直接交互。
c) 数据也必须是多副本,且多副本同时提供服务,短时间内单个节点的故障不会对系统造成影响,但长时间的故障必须人工干预,移除故障节点补入新的节点。

访问服务器:

访问服务器主要对外提供读写访问操作,需考虑以下设计:
a) 对外支持核心数据的增删改查操作。
b) 访问发生时,首先查询索引服务器取得关键字的数据服务器列表,然后根据列表,如果是读操作,则逐一去读取,直到读取成功,如果是写操作,则并行发起写入,并收集响应,可根据策略决定收集到多少响应才算成功。

技术选型:
访问服务器相对来说比较劲量,应该基本不存在借用开源方案了。

性能和可靠性:
a) 访问服务器是一层薄薄厚厚的代理,单点性能的好坏取决于工程的设计。
b) 访问服务器基本可以说无状态,单点性能不够时可以布署新的节点。
c) 单点故障基本对服务没有影响。

校正服务器:

校正服务器的主要功能为:
由于系统支持更新操作,同一个关键字的多个数据副本有一定概率存在版本不一致的问题,需要有一定的机制发现并修正这种错误。

性能和可靠性:
a) 性能的高低对系统没有绝对的影响,但更高的性能使系统的全库检查周期更短。
b) 可以单点,也可以集群工作,但集群工作时应避免重复对同一个关键字执行检查。

均衡服务器:

均衡服务器主要用于解决:
a) 各数据服务器的数据存储量由于数据的不均衡性或其他原因造成较大分布不均,需要有一种机制进行数据位置调度。
b) 当现有数据服务器的数据存储量超过警戒水位时,需要增加新的服务器,此时也要进行数据的迁移调度。

性能和可靠性:
a) 同样的,性能的高低对系统没有绝对的影响,但更高的性能使系统的迁移速度更快,更容易达到数量平衡。
b) 可以单点,也可以集群工作,但集群工作时应避免重复对同一个关键字执行迁移。

服务器设计:

索引服务器:

简单方案:
使用mysql做核心存储,数据备份依靠mysql主从同步。
索引服务器单点布署,故障发生时,人工介入切主从,切索引服务器。

复杂方案:
使用leveldb做核心存储,数据备份需要自研。
采用paxos主从选举方案,3点互备,主从同步,故障自动切换。

索引记录格式:
key -> dbnodeid1,datatime1,flag;dbnodeid2,datatime2,flag;dbnodeid3,datatime3,flag;
flag取值:
0 - 正常
1 - 删除中,即已发起删除操作,但未得到dbnodeid确认。

数据服务器内存负载表:
dbnodeid1 -> cpu, disk
dbnodeid2 -> cpu, disk
dbnodeid3 -> cpu, disk

接口列表:
设置索引
命令格式:setidx key dbnodeid
删除索引
命令格式:delidx key dbnodeid soft/hard
说明:soft软删除由访问服务器调用,表示打上删除标记。hard由数据服务器调用,表示完成删除。
获取(批量)索引
命令格式:getidx key1,key2,key3,......
说明:返回的dbnodeid列表,列表分两个子列表:存活的列表和非存活的列表,其中存活的列表按time降序、cpu负载升序。
安排索引
命令格式:planidx key 3
说明:什么时候出现安排索引?当新的kv写入系统时,get key返回空,这时需要plan key给到计划的dbnodeid节点列表。返回的列表应排除非存活的节点,并且按disk、cpu升序。
遍历索引
命令格式:seekidx key/null 10000
设置数据服务器负载
命令格式:setload dbnodeid cpu disk
获取数据服务器负载表
命令格式:getloads

设置索引流程:

  1. 仅写库。

删除索引流程:

  1. 对于软删除,写标志。
  2. 对于硬删除,删对应索引项,如果索引项全部为空,清除记录。

获取索引流程:

  1. 读库,组织响应。

安排索引流程:

  1. 对存活的dbnode节点按disk排序,返回指定的个数,并组织响应。

设置数据服务器负载流程:

  1. 仅更新内存表。

获取数据服务器负载表流程:

  1. 读取内存表,组织响应。

多点同步流程:

  1. 略。
数据服务器:

mysql存储方案:
使用mysql表来存储kv,不需要做主从同步。
数据的多副本由写流程完成。
数据的均衡迁移或故障备份由均衡流程完成。

redis存储方案:
使用redis来存储kv,需要做持久化,不需要做主从同步。
数据的多副本由写流程完成。
数据的均衡迁移或故障备份由均衡流程完成。

leveldb存储方案:
使用leveldb来存储kv,不需要做多点同步(因为不知道多点)。
数据的多副本由写流程完成。
数据的均衡迁移或故障备份由均衡流程完成。

binlog小文件存储方案:
使用自研的binlog来存储kv,不需要做多点同步(因为不知道多点)。
需要做本地索引,用来定位kv的位置,可能不如leveldb来得方便。
数据的多副本由写流程完成。
数据的均衡迁移或故障备份由均衡流程完成。

文件系统存储方案:
使用文件系统来存储kv,不需要做多点同步(因为不知道多点)。
使用文件系统路径如 xx/xx/xx/xx/key 来定位kv的位置。
数据的多副本由写流程完成。
数据的均衡迁移或故障备份由均衡流程完成。

接口列表:
写KV
命令格式:setkv key value [time]
说明:指定time表示复制数据
读KV
命令格式:getkv key
删KV
命令格式:delkv key

写KV流程:

  1. 写库。
  2. 同步向索引服务器写索引。
  3. 若索引写入失败,则删库。
  4. 返回结果。
    小概率异常:若第2步返回失败,第3步停电,数据有多。
    解决办法:建一个binlog文件,处理原子异常。

读KV流程:

  1. 读库。
  2. 如果读取失败,同步向索引服务器删除索引。
  3. 返回结果。

删KV流程:

  1. 删库。
  2. 同步向索引服务器删除索引。
  3. 返回结果。

上报负载流程:

  1. 定时上报disk、cpu负载。
访问服务器:

访问服务器不涉及存储,是纯逻辑服务器。

接口列表:
写KV
命令格式:setkv key value
读KV
命令格式:getkv key
删KV
命令格式:delkv key

写KV流程:

  1. 访问服务器接收setkv key value请求。
  2. 访问服务器向索引服务器同步询问getidx key。
  3. 索引服务器查询key的索引记录,取出dbnodeid列表,列表按time降序、cpu升序,排除非存活的节点。
  4. 访问服务器接收返回结果,如果列表为空,或者列表数量太小如1,访问服务器可以再向索引服务器同步询问安排索引planidx key 3。
  5. 索引服务器考虑dbnode节点的disk、cpu,给出建议的dbnodeid列表,但排除非存活的节点。
  6. 访问服务器接收返回结果,向这些数据服务器列表发起setkv key value请求(建议并发,如果不好弄,串行也可以)。
  7. 每个数据服务器接收setkv key value请求,先写存储,再写索引,串行进行,若失败则回滚,返回操作结果。
  8. 访问服务器接收返回结果,统计写入成功的数量,并根据策略决定最终操作结果,响应给请求者。

读KV流程:

  1. 访问服务器接收getkv key请求。
  2. 访问服务器向索引服务器同步询问getidx key。
  3. 索引服务器查询key的索引记录,取出dbnodeid列表,列表按time降序、cpu升序,排除非存活的节点。
  4. 访问服务器接收返回结果,如果列表为空,返回失败给请求者。
  5. 非空,访问服务器继续向列表顺序串行向数据服务器发起getkv key请求。
  6. 数据服务器接收getkv key请求,读存储,可能找不到或异常,返回操作结果。
  7. 访问服务器接收返回结果,并转发给请求者。

删KV流程:

  1. 访问服务器接收delkv key请求。
  2. 访问服务器向索引服务器同步询问getidx key。
  3. 索引服务器查询key的索引记录,取出dbnodeid列表,列表按time降序、cpu升序,但不排除非存活的节点。
  4. 访问服务器接收返回结果,如果列表为空,返回成功给请求者。
  5. 非空,访问服务器继续按列表顺序串行向数据服务器串行发起delkv key请求。
  6. 每个数据服务器接收delkv key请求,删存储,再硬删索引,串行进行,若失败则回滚,返回操作结果。
  7. 访问服务器收集结果,对操作失败的节点执行软删除。
  8. 若执行到这一步,访问服务器总是返回成功。
校正服务器:

两个功能:

  1. 对同一个key的多个dbnodeid更新时间差值太大的进行校验。
  2. 对key下dbnodeid有软删除标志的进行索引进行数据清理。

版本差值校正流程:

  1. 校正服务器向索引服务器同步询问seekidx ""/key 10000
  2. 索引服务器执行开头(区间)遍历,返回指定的个数索引信息。
  3. 校正服务器检查每个key的索引时差,如果超过一定时间(如60S),则判定为需要校正。
  4. 取出最新的那个dbnodeid和time,向有差值的数据服务器发起setkv key value time调用。

软删除清理流程:

  1. 校正服务器向索引服务器同步询问seekidx ""/key 10000
  2. 索引服务器执行开头(区间)遍历,返回指定的个数索引信息。
  3. 校正服务器检查每个key的软删除标志,如果有,则向对应的数据服务器发起delkv key调用。
均衡服务器:

主要用于自动迁移数据,触发条件:

  1. 某个数据服务器的disk占比超过平均值的20%,需要往disk低于平均值最多的节点迁移。
  2. 某个数据的副本数低于3个,需要往disk低于平均值最多的节点迁移。
  3. 某个数据的副本>=3个,但是其中部分节点已经明确不在系统中(即已被正式移除),需要往disk低于平均值最多的节点迁移。

自动迁移流程1:

  1. 均衡服务器向索引服务器同步询问seekidx ""/key 10000
  2. 均衡服务器向索引服务器同步询问负载情况
  3. 均衡服务器检查每个key的索引信息,检查是否有disk占比超过平均值20%的节点,如果有,找到disk占比最低的节点。
  4. 均衡服务器向待迁移的数据服务器读取KV,并向目的写入KV。
  5. 若上一步操作成功,继续删除源的KV,若删除失败,写soft删除标志。

自动迁移流程2:

  1. 均衡服务器向索引服务器同步询问seekidx ""/key 10000
  2. 均衡服务器检查每个key的索引信息,首先检查有效副本数,即如果副本的数据服务器已经被移除,则先擦除这部分索引。
  3. 继续检查副本数,如果<3,则需要执行迁移,即接下的两步。
  4. 均衡服务器向待迁移的数据服务器读取KV,并向目的写入KV。
  5. 若上一步操作成功,继续删除源的KV,若删除失败,写soft删除标志。

总而言之,设计一个完整的健全的分布式存储系统本身就是一个艰难的任务,由于我的时间和精力有限,以上写得比较基础,权当读者参考。

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容