达达集团技术

StarRocks在到家的运维实践

1.引入 StarRocks

2.StarRocks的特点

3.部署中的问题

4.数据备份和同步

5.集群配置和升级

6.监控和报警方案

7.实践总结

1.引入StarRocks

作为DBA团队,我们负责运维到家MySQL,Redis,Es相关数据库,去年下半年新业务希望数据库支撑复杂的OLAP需求,进行千亿量级的数据分析,满足秒级多维度和条件的组合查询,实时聚合,改造成本要求尽量低。

对于我们来说:MySQL原生不支持跨实例查询,分库分表后业务逻辑复杂,空间使用超过4T,复杂查询耗时长。

Redis由于数据结构单一,查询方式相对简单,不满足复杂条件查询且内存成本昂贵。

Es业务改造成本高,业务逻辑需要重构。

根据市面上的相关产品,按照查询性能、数据规模、整体架构、扩展能力、开发成本、运维效率等原则进行评估,评估结果如下:

经过对比分析,StarRocks的OLAP性能高,硬件成本低(可混合部署),改造成本低(兼容SQL协议),运维成本低(开源且组件较少)。我们将此方案推荐给业务,在进行改造调研后,最终选择StarRocks作为新业务数据库解决方案。

2.StarRocks的特点

StarRocks是面向多种数据分析场景,高性能的分布式数据库, 架构设计融合了MPP数据库,以及分布式系统的设计思想 。其 架构精简,采用全面向量化引擎 ,支持 智能查询优化和物化视图 , 支持联邦查询,高效数据更新,以及标准SQL , 支持流式/批处理方式处理数据,原生支持高可用且节点易扩展。

什么是MPP数据库,向量化引擎是什么原理,分布式系统解决哪些痛点?

2.1 MPP数据库:是针对分析工作负载进行了优化的数据库,用于聚合和处理大型数据集。MPP数据库往往是列式的,因为MPP数据库通常将每一列存储为一个对象,而不是将表中的每一行存储为一个对象。
2.2 向量化引擎: 向量化计算是一种特殊的并行计算的方式,它可以在同一时间执行多次操作,通常是对不同的数据执行同样的一个或一批指令,或者说把指令应用于一个数组/向量。 向量化计算用到了CPU的SIMD的特性,SIMD将原本单指令单数据的处理方式改进为 指令译码后几个执行部件同时访问 内存 ,一次性获得所有操作数进行运算。这个特点使SIMD特别适合于多媒体应用等数据密集型运算。
2.3 分布式系统: 用于解决海量数据存储问题,节点可根据负载进行动态水平扩容/缩容。
2.4 StarRocks集群架构图如下:

StarRocks集群及组件介绍:

FE: 前端节点,负责管理元数据,管理客户端连接,进行查询规划,查询调度等工作。

BE: 后端节点,负责数据存储,计算执行,以及compaction,副本管理等工作。

Broker:提供 和外部HDFS/对象存储等数据源对接的中转服务,辅助提供导入导出功能。

3.部署中的问题

官方文档给出了手动部署的流程,但是部署验证过程中,我们遇到了一些问题。

3.1 权限管理:StarRocks 手动部署的组件有FE,BE,Broker三个,不要用root权限启动Broker组件,否则会有HDFS权限不足的报错,例如 Permissiondenied: user=root, access=READ,inode="/user/hive/xx/xx/xxx/xxx=xx/part-00005-875acf23-bc34-4986-81aa-6aa3b7a5841c.c000":hiv e:hive:-rwxrwx—x

报错信息显示是root对Hive权限不足,其根本原因是不能以root权限启动,原因是远端Hadoop创建的任务会通过Broker从Hive拉取数据,权限继承启动账户,故使用普通账户重启后问题解决,而FE/BE组件无此限制。

3.2 参数优化:默认的配置文件参数可能不满足复杂场景下的需要,而我们在MySQL上的经验积累并不能作为参考,所以需要根据报错信息调整参数,重启相关组件使之生效。 需要注意的是有些参数虽然支持set global xxx=xxx的形式(类似于MySQL全局变量),但实际是不生效的,因此对于此类参数仍然需要重启组件的方式进行调整。

3.2.1 FE组件中有两个控制连接数的参数: qe_max_connection(单节点最大连接数,默认1024)和max_conn_per_user (单用户最大连接数,默认100)。在MySQL的使用经验中,应用通过连接池进行访问,机器数量越多,消耗连接越多,实际上我们也遇到了连接数不足的情况,所以对应参数我们调整为1W和2K。

3.2.2 业务侧有很多通过stream load方式进行数据写入的场景,而默认配置下频繁报错,其信息如下:

参考了相关资料后,我们认为是单次写入量过大和频次过多导致,据此调整参数:

streaming_load_rpc_max_alive_time_sec=4800(默认1200)

tablet_writer_open_rpc_timeout_sec=240(默认60)

max_routine_load_task_num_per_be=9(默认5)

优化参数后数据导入报错大量减少。

3.3 语法兼容性:StarRocks虽然兼容大部分MySQL语法,但在上线的审核过程中,由于和MySQL的差异使得在MySQL上的经验积累并不能照搬套用。

3.3.1 DDL建表语句中需要指定副本数,而不是通常在部署过程中由配置文件的参数指定(副本数典型值n,集群节点2n-1,如果副本数大于可用的BE节点数,则在写入时提示无法写入),这会导致默认运维配置的参数变成实际由业务侧管理,这增加了业务和运维间的耦合性。假设水平扩容期间,一旦业务出现写入失败问题,需要双方同时进行排查,即问题定位时会造成困扰。而通常的互联网设计都是尽量避免不同系统间的耦合,数据库尽量避免使用存储过程也是同样的道理。

3.3.2 StarRocks索引类型(Bitmap , BloomFilter , 稀疏索引)多于 MySQL(Btree+) ,索引类型和字段选择不合理 会造成查询效率低下,所以需要根据业务场景和查询条件进行定制,开发人员需要对业务逻辑和索引类型了然于心。作为回报相同业务场景下MySQL的查询耗时在分钟级,而在StarRocks进行索引重构后耗时在秒级。

3.3.3 OLAP分析时常用到分区表, RANGE分区用于将数据划分成不同区间, 逻辑上可以理解为将原始表划分成了多个子表。业务上,多数用户会选择采用按时间进行partition,分区表可区分冷热数据且按照分区删除数据更加迅速,但是StarRocks 分区表功能暂不支持根据时间维度自动创建分区,所以如果有时间维度的查询分析需求,需先创建好分区。

4.数据备份和同步

我们线上有两套集群,业务需要在不同集群间进行数据同步。但是StarRocks周边工具较少,并没有 类似Mysqldump这类工具。
官方推荐使用backup/restore进行集群间数据同步,同步方式基于snapshot,在备份期间的写入数据无法记录于快照,所以建议备份期间不写入数据或者在restore后对目的集群二次同步数据。backup目的源支持HDFS/AWS的S3以及其他云端产品,但由于条件限制(SDK不兼容),最终选择搭建本地HDFS临时服务进行本地存储,其成本最低(本地部署),效率最高(无需三方沟通)。
backup /restore的时间消耗依赖于备份集大小,虽然官方不支持db级别的备份,但是可以将db中所需的table列表拼接成字符串, backup命令不支持跨db ,如果是多个db则需要多次备份。

集群间数据同步流程图如下:

根据实践,亿级别的单表恢复时间在1分钟 , 对比Mysqldump时间减少70%以上。

5.集群配置和升级

基于实际的生产环境,正式环境规格为3节点 , FE/BE混部 , 16C/64G/300G , 后期如果业务增长负载增加 , 既能进行节点拆分也能进行水平扩容。
由于1.18版本StarRocks不能进行BE节点间数据balance,根据木桶效应,单BE数据空间会成为存储瓶颈,在实际生产环境中我们也遇到了磁盘空间不足的情况,因此我们考虑旧1.18集群升级到1.19(该版本解决了空间问题,BE节点磁盘空间会按照不同百分比存储数据)。
升级过程中的注意事项:升级方案和回滚策略。
准备事项:组件升级前备份BE/FE配置文件和meta数据(类似mysql系统库,存放了FE/BE节点的数据版本号,用于调度节点数据)。
升级方案:先BE后FE,FE组件先升级非master节点,最后master节点。
校验方案:运维侧在组件升级后需要校验版本号,且以节点的cmd执行结果为准,查看版本号语句如下:show variables like"%version%",注意StarRocks的select version()结果不准确;程序侧校验业务完整性并进行回归测试。
回滚方案:升级前务必做好降级方案和修复方案。如果升级失败后切记数据二次写入进而造成数据修复方案复杂。

6.监控和报警方案

官方文档中建议使用 Prometheus +Grafana的方式,不过官方提供的Grafana监控模板只有少数监控项支持添加告警,同时官方提供了Http API,如果有自研运维平台,可自行抓取StarRocks集群信息进行监控报警。
运维平台的适应性改造:由于到家运维平台底层数据库是MySQL,而StarRocks兼容SQL语法,我们将StarRocks集群作为基础服务接入现有运维平台,并进行流程改造。
例如我们将原有流程中对MySQL拓扑检查和SQL分析相关逻辑做特殊化处理,改造后使得现有流程兼容StarRocks, 例如数据查询,应用授权,上线流程等 , 拓展了运维部的基础服务能力。

7.实践总结: