37DATA

ElasticSearch的框架与应用

简介:

     ElasticSearch 新简介
     github: https://github.com/elastic/elasticsearch

背后故事:

        许多年前,一个刚结婚的名叫 Shay Banon 的失业开发者,跟着他的妻子去了伦敦,他的妻子在那里学习厨师。在寻找一个赚钱的工作的时候,为了给他的妻子做一个食谱搜索引擎,他开始使用 Lucene 的一个早期版本。直接使用 Lucene 是很难的,因此 Shay 开始做一个抽象层,Java 开发者使用它可以很简单的给他们的程序添加搜索功能。他发布了他的第一个开源项目 Compass。后来 Shay 获得了一份工作,主要是高性能,分布式环境下的内存数据网格。这个对于高性能,实时,分布式搜索引擎的需求尤为突出, 他决定重写 Compass,把它变为一个独立的服务并取名 Elasticsearch。第一个公开版本在2010年2月发布,从此以后,Elasticsearch 已经成为了 Github 上最活跃的项目之一,他拥有超过300名 contributors(目前736名 contributors )。一家公司已经开始围绕 Elasticsearch 提供商业服务,并开发新的特性,但是,Elasticsearch 将永远开源并对所有人可用。据说,Shay 的妻子还在等着她的食谱搜索引擎…

认识ES:

  • 一个分布式实时文档存储,每个字段可以被索引与搜索

  • 一个分布式实时分析搜索引擎

  • 能胜任上百个服务节点的扩展,并支持PB级别的结构化或非结构化数

ES应用场景:

  • 提示与纠错

Image

  • 搜索

Image

  • Feed流

Image

  • ELK日志

    Image

  • 文章实时热度

    Image

    ES架构:

    Image

    Image

    Mapping: 分词器, 和格式相关 Gateway代表ElasticSearch索引的持久化存储方式。      

    在Gateway中,ElasticSearch默认先把索引存储在内存中,然后当内存满的时候,再持久化到Gateway里。当ES集群关闭或重启的时候,它就会从Gateway里去读取索引数据。比如LocalFileSystem和HDFS、AS3等。  

    DistributedLucene Directory,它是Lucene里的一些列索引文件组成的目录。它负责管理这些索引文件。包括数据的读取、写入,以及索引的添加和合并等。  

    River,代表是数据源。是以插件的形式存在于ElasticSearch中。    

    Mapping,映射的意思,非常类似于静态语言中的数据类型。比如我们声明一个int类型的变量,那以后这个变量只能存储int类型的数据。比如我们声明一个double类型的mapping字段,则只能存储double类型的数据。          

    Mapping不仅是告诉ElasticSearch,哪个字段是哪种类型。还能告诉ElasticSearch如何来索引数据,以及数据是否被索引到等。  

    Search Moudle,这个很简单   

    Index Moudle,这个很简单   

    Disvcovery,主要是负责集群的master节点发现。比如某个节点突然离开或进来的情况,进行一个分片重新分片等。这里有个发现机制。         

    发现机制默认的实现方式是单播和多播的形式,即Zen,同时也支持点对点的实现。另外一种是以插件的形式,即EC2。  Scripting,即脚本语言。包括很多,这里不多赘述。如mvel、js、python等。      

    Transport,代表ElasticSearch内部节点,代表跟集群的客户端交互。包括 Thrift、Memcached、Http等协议   RESTful Style API,通过RESTful方式来实现API编程。  

    3rd plugins,代表第三方插件。  

    Java(Netty),是开发框架。  

    JMX,是监控。

    集群分片:

Image

启动ES:

docker run -d -e ES_JAVA_OPTS="-Xms256m -Xmx256m" \ -p 9200:9200 -p 9300:9300 \ -v ${elastcisearch}/config/elasticsearch.yml:${docker_es}/config/elasticsearch.yml \ -v ${elastcisearch}/data:${docker_es}/data \ -v ${elastcisearch}/logs:${docker_es}/logs \ -e "discovery.type=single-node" \ -h es-node0 --name es-00 elasticsearch:7.2.0

ES写入过程:

Image

translog

1.保证在filesystem cache中的数据不会因为elasticsearch重启或是发生意外故障的时候丢失。

2.当系统重启时会从translog中恢复之前记录的操作。

3.当对elasticsearch进行CRUD操作的时候,会先到translog之中进行查找,因为tranlog之中保存的是最新的数据。

4.translog的清除时间时进行flush操作之后(将数据从filesystem cache刷入disk之中)。flush 1.es的各个shard会每个30分钟进行一次flush操作。2.当translog的数据达到某个上限的时候会进行一次flush操作。

缓存层级

Image

Buffer IO 适用于普通类型的文件读写,性能尚可,操作简单,无注意事项。

MMAP 小数据量读写性能高,但不灵活。

Direct IO 需要自己控制Cache时,可以适用Direct IO,例如数据库/中间件应用,可以避免文件的读写还经过一层Page Cache,造成额外开销。

Lucene索引:

Image

最终形成了两个新文件:.cfs:用于保存数据 .cfe:保存了前者的一个Entry Table

offset方式读取文件:fseek 删除过程会有遗留的.liv文件标记 LUKE读取工具

排序:

Image

分词器:

  • Standard Analyzer - 默认分词器,按词切分,小写处理 

  • Simple Analyzer - 按照非字母切分(符号被过滤),小写处理 

  • Stop Analyzer - 小写处理,停用词过滤(the ,a,is) 

  • Whitespace Analyzer - 按照空格切分,不转小写 

  • Keyword Analyzer - 不分词,直接将输入当做输出

  • Pattern Analyzer - 正则表达式,默认 \W+ Language - 提供了 30 多种常见语言的分词器 

  • Customer Analyzer - 自定义分词器

ES写入过程:

Image

实战:

  • 索引文章

    ImageImage

    Image

    设置刷新内存时间:refresh_interval

    Image

    PUT /paraper/essay/6?refresh=wait_for POST /paraper/_refresh

    Image

    POST paraper/_flush?wait_if_ongoing

    Image

    POST paraper/_flush?_forcemerge?max_num_segments=1

Image

自增ID:体积小,消耗小,易保存。吞吐量锁限制,爬虫安全,数据迁移合并操作为 UUID:吞吐量大,数据合并迁移方便。体积大,存储消耗性能(索引大),排序慢

  • 查看索引

ImageImage

  • 搜索索引

    ImageImage

    ImageImage

Image

Image

  • 多条件搜索

Image

  • 先查询再聚合

    Image

    Lucene过程图:

    Image

  • https://www.shenyanchao.cn/blog/2018/12/04/lucene-index-files/

  • .fdx:Field Index,.fdt文件的索引/指针。通过该文件可以快速从.fdt文件中读取field数据

  • .nvd, .nvm:Norms,这两个都是用来存储Norms信息的,前者用于存储norms的数据,后者用于存储norms的元数据。

  • .dvd, .dvm:Per-Document Values,这两个都是用来存储DocValues信息的,前者用于数据,后者用于存储元数据。

  • dim,.dii:Point values,这两个文件用于记录indexing的Point信息,前者保存数据,后者保存索引/指针,用于快速访问前者。

  • .tvd:Term Vector Data,用于存储term vector数据。

  • .tvx:Term Vector Index,用于存储Term Vector Data的索引数据。

  • .fnm:Fields,用于记录fields设置类信息,比如字段的index option信息,是否存储了norm信息、DocValue等。

  • .tim:Term Dictionary,存储所有文档analyze出来的term信息。同时还包含term对应的document number以及若干指向.doc, .pos, .pay的指针,从而可以快速获取term的term vector信息。。

  • .tip:Term Index,该文件保存了Term Dictionary的索引信息,使得可以对Term Dictionary进行随机访问。

  • .doc:Frequencies,存储了一个documents列表,以及它们的term frequency信息。

  • .pos:Positions,和.doc类似,但保存的是position信息。

  • .pay:Payloads,和.doc类似,但保存的是payloads和offset信息。

    Image

  • refresh: 这些在内存缓冲区的文档被写入到一个新的段中,且没有进行 fsync 操作。这个段被打开,使其可被搜索。内存缓冲区被清空。

  • fsync: buffer->oscache(segment)

  • flush: 所有在内存缓冲区的文档都被写入一个新的段。缓冲区被清空。一个提交点被写入硬盘。文件系统缓存通过 fsync 被刷新(flush)。老的 translog 被删除。

高可用:

  • Bully选举

    Image

    首先是ClusterState版本号的比较,版本号越大优先级越高,然后是节点id的比较,id越小优先级越高。ClusterState是Master向集群中各个节点发送的集群状态,这个状态有一个版本号码,如果集群状态发生了变化,比如集群新增了节点成员或者有节点成员退出了,那么这个版本号就会加一,比对这个版本号的目的是让拥有最新状态的节点成为Master的优先级最高。P3节点发现Master P6对自己长时间不作出响应,P3节点会请求其它节点判断P6节点是否存活,如果有1/2以上节点都认定P6存活,那么P3就会放弃发起选举

  • 脑裂处理

    法定人数/多数机制(Quorum) m = n / 2 + 1  (ES选用这种)

    隔离机制(Fencing) 

    隔离复活的节点 冗余通信机制(Redundant communication) 增加心跳线

    Image

  • 宕机问题 

    方法一,Quorums(法定人数)方式 比如3个节点的集群,Quorums = 2,也就是说集群可以容忍1个节点失效,这时候还能选举出1个lead,集群还可用。比如4个节点的集群,它的Quorums = 3,Quorums要超过3,相当于集群的容忍度还是1,如果2个节点失效,那么整个集群还是无效的。这是ZooKeeper防止“脑裂”默认采用的方法。 

    方法二,添加心跳线 集群中采用多种通信方式,防止一种通信方式失效导致集群中的节点无法通信。比如,添加心跳线。原来只有一条心跳线路,此时若断开,则接收不到心跳报告,判断对方已经死亡。若有2条心跳线路,一条断开,另一条仍然能够接收心跳报告,能保证集群服务正常运行。心跳线路之间也可以 HA(高可用),这两条心跳线路之间也可以互相检测,若一条断开,则另一条马上起作用。正常情况下,则不起作用,节约资源。 

    方法三,启动磁盘锁定方式。使用磁盘锁的形式,保证集群中只能有一个Leader获取磁盘锁,对外提供服务,避免数据错乱发生。但是,也会存在一个问题,若该Leader节点宕机,则不能主动释放锁,那么其他的Follower就永远获取不了共享资源。于是有人在HA中设计了"智能"锁。正在服务的一方只有在发现心跳线全部断开(察觉不到对端)时才启用磁盘锁。平时就不上锁了 

    方法四,仲裁机制方式。脑裂导致的后果是从节点不知道该连接哪一台Leader,此时有一个仲裁方就可以解决此问题。比如提供一个参考的IP地址,心跳机制断开时,节点各自ping一下参考IP,如果ping不通,那么表示该节点网络已经出现问题,则该节点需要自行退出争抢资源,释放占有的共享资源,将服务的提供功能让给功能更全面的节点。以上方式可以同时使用,可以减少集群中脑裂情况的发生,但不能完全保证,比如仲裁机制中2台机器同时宕机,那么此时集群中没有Leader 可以使用。此时就需要人工干预了。

  • 一致性问题

    Image

    一致性问题 Master有两种指令,一种是send指令,另一种是commit指令,Master将最新集群状态推送给其它节点的时候(这是send指令),Master节点进入等待响应状态,其它节点并不会立刻应用该集群状态,而是首先会响应Master节点表示它已经收到集群状态更新,同时等待Master节点的commit指令。 

    quorum, 有一台机器失效,对于集群来说就有可能丢失数据 基本想法:a、多数派:每次写都保证写入大于N/2个节点,每次读保证从大于N/2个节点中读。比如5个节点,每次写大于3个节点才算成功;读也是大于3个节点才算成功。

     Paxos: Acceptor(master), proposer(slave) , Client Client向proposer提交任务, proposer通过不断地2PC提交N给Acceptor,如果Acceptor在处理其他请求,则N+1后再次请求 

    准备阶段:一个节点被选为领导者并且选择序列号x和值v创建提议P1(x,v)。领导者把P1发送给接收者并等待大部分节点响应。接收者一旦接收到提议P1(x,v)会做下面的事:如果是接收者第一次收到提议而且它选择赞同,回复“赞同”-这是接收者的承诺,它将承诺拒绝将来所有小于x的提议请求。- 如果接收者已经赞同了提议:比较x和接收者赞同的提议的最高序列号,称为P2(y,v2) 若x<y,回应“拒绝”以及y的值 - 若x>y,回应“赞同”以及P2(y,v2) 

    接受阶段 - 如果大部分接收节点未能回应或者回应“拒绝”,领导者放弃这次协议并重新开始。- 如果大部分接收节点回应“赞同”,领导者也会接受大部分节点接受的协议的值。领导者选择这些值的任意一个并发送接受请求以及提议序列号和值。- 当接收者收到接受请求消息,它只在下面两种情况符合时发送“接受”信息,否则发送“拒绝”:- 值和之前接受的提议中的任一值相同 - 序列号和接收者赞同的最高提议序列号相同 - 如果领导者没有从大部分节点那接收到“接受”消息,它会放弃这次提议重新开始。但是如果接收到了大部分的“接受”消息,深思熟虑后协议也可能被终止。作为优化,领导者要给其它节点发送“提交”信息。


     raft: http://thesecretlivesofdata.com/raft/

  • 扩容

    水平扩容:增加节点数量,扩容了之后shard会rebalance,每个节点上的shard数量会变少,这样就会每个shard能使用的内存、磁盘空间更大,整个系统性能会更好。

     垂直扩容:使用更强大的服务器替换老服务器或者给服务器增加内存、硬盘空间、换SSD等硬件。价格昂贵,需要做数据迁移。

  • 负载均衡

    1、轮询法 2、随机法 3、源地址哈希法 4、加权轮询法 5、加权随机法 6、最小连接数法

    1、轮询法 将请求按顺序轮流地分配到后端服务器上,它均衡地对待后端的每一台服务器,而不关心服务器实际的连接数和当前的系统负载。 

    2、随机法 通过系统的随机算法,根据后端服务器的列表大小值来随机选取其中的一台服务器进行访问。由概率统计理论可以得知,随着客户端调用服务端的次数增多, 其实际效果越来越接近于平均分配调用量到后端的每一台服务器,也就是轮询的结果。 

    3、源地址哈希法 源地址哈希的思想是根据获取客户端的IP地址,通过哈希函数计算得到的一个数值,用该数值对服务器列表的大小进行取模运算,得到的结果便是客服端要访问服务器的序号。采用源地址哈希法进行负载均衡,同一IP地址的客户端,当后端服务器列表不变时,它每次都会映射到同一台后端服务器进行访问。

     4、加权轮询法 不同的后端服务器可能机器的配置和当前系统的负载并不相同,因此它们的抗压能力也不相同。给配置高、负载低的机器配置更高的权重,让其处理更多的请;而配置低、负载高的机器,给其分配较低的权重,降低其系统负载,加权轮询能很好地处理这一问题,并将请求顺序且按照权重分配到后端。

     5、加权随机法 与加权轮询法一样,加权随机法也根据后端机器的配置,系统的负载分配不同的权重。不同的是,它是按照权重随机请求后端服务器,而非顺序。6、最小连接数法 最小连接数算法比较灵活和智能,由于后端服务器的配置不尽相同,对于请求的处理有快有慢,它是根据后端服务器当前的连接情况,动态地选取其中当前 积压连接数最少的一台服务器来处理当前的请求,尽可能地提高后端服务的利用效率,将负责合理地分流到每一台服务器。

  • 自动发现注册机制

    Azure classic discovery 插件方式,多播 

    EC2 discovery 插件方式,多播 

    Google Compute Engine (GCE) discovery 插件方式,多播 

    Zen discovery 默认实现,多播/单播

     多播:同时把自身信息发给所有节点 单播:发送给集群内的一台节点,并从中获取所有其他节点信息