跳到主要内容

关于使用ELK的记录

· 阅读需 8 分钟

最近因为公司系统项目统计压力大 。mysql集群查询进入瓶颈,索引优化也不能快速响应统计结果,于是接入 Elasticsearch、Logstash、Kibana 三大开源套餐。

先交代一下背景。统计类查询和普通的业务查询不是一回事:业务查询大多是按主键或索引取少量行,MySQL 很擅长;而统计查询往往要扫描大量数据做分组、聚合,B+ 树索引在这种场景下帮助有限。数据量上来之后,一条统计 SQL 跑几秒甚至几十秒,还会拖累同库的业务请求。这时候常见的思路就是把统计负载从 OLTP 数据库里剥离出去,交给专门做聚合分析的引擎。Elasticsearch 底层是倒排索引加列式的 doc values,天生适合过滤和聚合,ELK 这套组合也是社区里最容易上手的方案之一,所以值得记一笔。

https://www.elastic.co/cn/

选型与分工

关于 Logstash 其实还可以使用ali DataX 也是可以的 功能比 Logstash 还要更上一层楼。我采用的还是Logstash,因为ali DataX 我是后面才知道的,所以就没有去替换掉了。

两者定位有些差别:Logstash 是 Elastic 家族的数据管道,靠 jdbc input 插件定时拉取数据,配置文件写好就能跑,和 Elasticsearch 的输出对接是现成的;DataX 是阿里开源的异构数据源同步工具,插件覆盖的数据源更多,批量吞吐也更强。对于"MySQL 定时增量同步到 ES"这种单一链路,Logstash 已经够用,没必要为了换而换。

三个组件的分工很清晰:

通过 Logstash抓取过滤mysql统计数据,增量同步到Elasticsearch中。

再通过项目java api 调用 Elasticsearch查询;

Kibana 可以web可视化集群,索引具体情况。

也就是说 Logstash 管数据进,Java 应用管数据出,Kibana 负责让人看得见集群和索引的健康状况,三者互不干扰。

Elasticsearch查询语句挺简单,但是感觉官网的例子不是很多,很多比较复杂的聚合需要自己摸索。比如多层嵌套的 aggregation、聚合结果再排序这类写法,文档里往往只给最基础的示例,实际业务里得靠 Kibana 的 Dev Tools 一点点调出来。

整个学习难度不是很高。入手很快

Java 端的接入

java api jar 使用的是 elasticsearch-rest-high-level-client

之后再根据业务场景自己封装了一下工厂,抽象了几层代码给开发人员使用。封装的目的很简单:不希望每个开发都去直接拼 SearchSourceBuilder、解析 SearchResponse,把常用的条件过滤、分页、聚合模式收敛成几个方法,业务代码只关心传参和拿结果。

api的查询和得到结果集方法可能是操作起来有点麻烦,但其实只要对着查询语句写的话。其实还是很好理解的。high level client 的 builder 结构和 Query DSL 的 JSON 结构基本是一一对应的,先在 Kibana 里把 DSL 调通,再翻译成 Java 代码,几乎不会出错。

总的感觉还是比较友善。所有的软件都是安装即可用,不过还是需要改一下配置,这个就不多说了。ip 端口,中文英文,密码等等。

只同步统计需要的字段

公司主要是做统计,于是通过 Logstash 抓取关键统计数据,公司几百万数据抓取关键统计数据只有100m不到。这也可以大大提升es 聚合统计的速度,所以不建议把整张表都抓取进来,只拿最主要的统计字段。

这一点值得展开说。Elasticsearch 的聚合是在内存和 doc values 上做的,索引越瘦,段文件越小,能缓存的比例就越高,聚合自然越快。而且 ES 不是数据库,不需要承担"存全量明细"的职责,明细永远以 MySQL 为准,ES 里只放统计维度和指标字段,坏了随时可以重建,心理负担也小很多。

所有流程如图:

这也是我给我们公司员工写的流程图。

其实如果抓取数据量极大,可以在中间使用kafaka进行缓存再次筛选。由于我们公司数据量并没有这么高,也就不用过分去开销其他服务器资源。中间加一层消息队列的意义在于削峰和解耦:上游抓取和下游写入速率不一致时,队列能兜住突发流量,也方便在中间再挂一道清洗逻辑。但每多一个组件就多一份运维成本,量不到就不要上。

Kibana 现在已经很完善了,写官方的Query DSL 语句也可以,写SQL语句也可以,但官方还是推荐使用 Query DSL 语句进行聚合与查询

增量同步的实现

关于增量同步,更新同步数据,我采用的方案是在抓取的数据表上设置data_version字段(乐观锁原理,数据版本号),一旦某条数据进行了修改或者是新增的数据,data_version这个值将会设置,我设置的是时间戳(公司业务原因),建议给此字段设置bigint类型,timestamp只能使用到2028年。

这样就给了数据一个版本,在同步增量与更新数据时,只会抓取data_version 改变过的数据。大大降低mysql服务器,Logstash,Elasticsearch的压力。

Logstash 可以记录上次抓取数据最后一条数据的 data_version 值。这样在下一次抓取数据时,带上 上一次的data_version 值,就可以通过where条件过滤出来有效数据。当然sql里面是要对data_version 进行排序的,因为Logstash 只会记录上次抓取时最后一条数据的data_version 值。

这背后就是 Logstash jdbc 插件的 tracking column 机制:插件把上次运行时最后一条记录的追踪列值持久化到本地文件,下次执行 SQL 时以参数形式带入 where 条件。所以排序是必须的——如果结果集不按 data_version 升序排,记录下来的"最后一条"就不是最大值,下一轮同步会漏数据。

不然每一次同步数据都要全表同步吗?那压力可想而知。应该没人会采用,全表同步只是第一次同步时会全表同步。之后都是增量与更新数据同步。

更新的数据之所以也能被同步到,是因为每条数据在 ES 里以固定的文档 id(通常就是 MySQL 主键)写入,同一条数据版本变了会被再次抓到,写入 ES 时就是一次覆盖更新,不会产生重复文档。

踩坑与注意

1)追踪列的类型要想清楚。上面提到过,时间戳做版本号时字段建议用 bigint,避免 timestamp 类型的取值上限问题;另外时间戳精度不够时,同一秒内的多次修改可能在边界上漏掉,业务上要能接受或者改用自增序列。

2)物理删除同步不到。data_version 方案只能感知新增和修改,MySQL 里直接 delete 的行不会再出现在结果集里,ES 中会残留旧文档。要么业务上用逻辑删除标记,把删除也当成一次"修改"同步过去;要么定期重建索引兜底。

3)不要把 ES 当唯一存储。ES 里的统计索引应该随时可以从 MySQL 重放出来,这样映射设计错了、字段要加了,直接删索引重跑一遍 Logstash 即可,不用做复杂的在线迁移。

4)增量 SQL 一定要验证排序和边界条件。where 条件用大于还是大于等于、排序是否生效,建议先手工执行一遍 SQL,再对比 Logstash 两轮运行的记录值,确认没有漏抓和重复抓。

小结

这次接入 ELK,核心收益是把统计聚合从 MySQL 集群里剥离出来:Logstash 基于 data_version 做增量同步,只抓统计需要的字段;Elasticsearch 负责聚合查询,Java 端用 high level client 封装后给业务使用;Kibana 兼顾调试 DSL 和观察集群状态。整套方案没有引入多余的组件,按数据量的实际规模做取舍,够用就好。

评论 / COMMENTS