改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

大数据机器集群规划

ckckck2025年1月10日5 浏览

大数据集群节点规划

硬件规划

硬件规划决定集群将使用多少硬件资源,以及什么配置的硬件资源。可以从以下几个维度进行评估:

  • 数据现状
    • 盘点所有数据情况,包括数据源、数据量、数据大小、数据维度等信息
  • 工作负载
    • 评估在集群与数据之上将执行的任务类型
    • 如实时计算、离线计算、图像处理、关系网络等应用场景以及是否提供OLTP服务等
  • 未来数据量预估
    • 根据数据源与业务应用场景可以对未来衍生的数据总量与数据增量做大致评估
    • 评估的时间范围视业务场景而定,建议做不少于一年的规划
  • 硬件资源现状
    • 盘点目前可用的硬件资源,确认是否满足所评估的规模及要求
    • 机房机柜空间、电源(双)等是否充足(需考虑后续扩容问题)
    • 网络交换机性能是否满足要求(建议万兆网卡(双))
    • 查看服务器磁盘、内存、CPU等资源是否需要补充
  • 硬件选择
    • 现有硬件资源不满足的需求的情况下,结合运维建议提出需要增加或者新采购的硬件型号、配置等
    • 确认所需服务器数

数据存储总量

  • 日增量 * 周期 * 副本数
  • 或者 单条数据大小 * 单日的数据量 * 周期 * 副本数

所需存储服务器数量

  • 数据总量 / (单台服务器总磁盘大小 * 0.8)

QPS估算和峰值QPS(二八定律)

  • QPS(TPS)= 并发数 / 平均响应时间
  • 峰值时间:每天80%的访问集中在20%的时间里,这20%时间叫做峰值时间
  • 峰值时间QPS = (日总PV数 * 波峰数据量占比) / (60 * 60 * 波峰时间)

Kafka/Pulsar机器规划

以电商平台为例,Kafka集群每天需要承载10亿+请求流量数据,一天24小时,对于平台来说,晚上12点到凌晨8点这8个小时几乎没多少数据涌入的。这里我们使用「二八法则」来进行预估,也就是80%的数据(8亿)会在剩余的16个小时涌入,且8亿中的80%的数据(约6.4亿)会在这16个小时的20%时间(约3小时)涌入。

通过上面的场景分析,可以得出如下:

  • QPS计算公式 = 640000000 ÷ (3 * 60 * 60) = 6万
  • 也就是说高峰期集群需要扛住每秒6万的并发请求。假设每条数据平均按20kb(生产端有数据汇总)来算,那就是1000000000 * 20kb = 18T
  • 一般情况下我们都会设置3个副本,即54T
  • 另外Kafka数据是有保留时间周期的,一般情况下是保留最近3天的数据,即54T * 3 = 162T

总结:要搞定10亿+请求,高峰期要支撑6万QPS,需要大约162T的存储空间。

Kafka物理机数量

系统高峰期的时候要支撑6万QPS,如果公司资金和资源充足的情况下,我们一般会让高峰期的QPS控制在集群总承载QPS能力的30%左右,这样的话可以得出集群能承载的总QPS能力约为20万左右,这样系统才会是安全的。

根据经验可以得出每台物理机支撑4万QPS是没有问题的,从QPS角度分析,我们要支撑10亿+请求,大约需要5台物理机,考虑到消费者请求,需要增加约1.5倍机器,即7台物理机。

:20万 / 4万 * 1.5 = 7台

Kafka磁盘类型选择

Kafka写磁盘是顺序追加写的,所以对于Kafka集群来说,我们使用普通机械硬盘就可以了。(省钱)

Kafka磁盘选择

根据前面的分析,我们需算出需要7台物理机,一共需要存储162T数据,大约每台机器需要存储23T数据,根据以往经验一般服务器配置11块硬盘,这样每块硬盘大约存储2T的数据就可以了,另外为了服务器性能和稳定性,我们一般要保留一部分空间,保守按每块硬盘最大能存储3T数据。

:ceil((162T / 7) / 11) = 3T

Kafka内存选择

从上图可以得出Kafka读写数据的流程主要都是基于os cache,所以基本上Kafka都是基于内存来进行数据流转的,这样的话要分配尽可能多的内存资源给os cache。

Kafka的核心源码基本都是用scala和java(客户端)写的,底层都是基于JVM来运行的,所以要分配一定的内存给JVM以保证服务的稳定性。对于Kafka的设计,并没有把很多的数据结构存储到JVM中,所以根据经验,给JVM分配6~10G就足够了。

从上图可以看出一个Topic会对于多个partition,一个partition会对应多个segment,一个segment会对应磁盘上4个log文件。假设我们这个平台总共100个Topic,那么总共有100 Topic * 5 partition * 3 副本 = 1500 partition。对于partition来说实际上就是物理机上一个文件目录,.log就是存储数据文件的,默认情况下一个.log日志文件大小为1G。

如果要保证这1500个partition的最新的.log文件的数据都在内存中,这样性能当然是最好的,需要1500 * 1G = 1500G内存,但是我们没有必要所有的数据都驻留到内存中,我们只保证25%左右的数据在内存中就可以了,这样大概需要1500 * 250M = 1500 * 0.25G = 375G内存,通过第二步分析结果,我们总共需要7台物理机,这样的话每台服务器只需要约54G内存,外加上面分析的JVM的10G,总共需要64G内存。还要保留一部分内存给操作系统使用,故我们选择128G内存的服务器是非常够用了。

Kafka CPU核数选择

  • Accept线程(1个)
  • Process线程(默认3个,一般设置9个)
  • RequestHandle线程(默认8个,一般设置32个)

估算下来Kafka内部有100多个线程。因此CPU一般建议在16核-32核。

总结:要搞定10亿+请求,需要7台物理机,每台物理机内存选择128G内存为主,需要16个CPU core(32个性能更好)。

Kafka网卡选择

千兆网卡和万兆网卡的区别最大之处在于网口的传输速率的不同,千兆网卡的传输速率是1000Mbps,万兆网卡的则是10Gbps万兆网卡是千兆网卡传输速率的10倍。高峰期的时候,每秒会有大约6万请求涌入,即每台机器约1万请求涌入(60000 / 7),每秒要接收的数据大小为:10000 * 20kb = 184M/s,外加上数据副本的同步网络请求,总共需要184 * 3 = 552 M/s。

一般情况下,网卡带宽是不会达到上限的,对于千兆网卡,我们能用的基本在700M左右,通过上面计算结果,千兆网卡基本可以满足,万兆网卡更好。

Starrocks集群规划

最低要求

为了实现集群高可用,建议集群最低3个节点,FE和BE分开部署也可以混合部署。

单节点配置要求

  • BE推荐16核64GB内存以上,FE推荐8核16GB内存以上。
  • 磁盘可以使用HDD或者SSD。
  • CPU必须支持AVX2指令集,cat /proc/cpuinfo | grep avx2确认有输出即可,如果没有支持,建议更换机器,StarRocks的向量化技术需要CPU指令集支持才能发挥更好的效果。
  • 网络需要万兆网卡和万兆交换机。

环境支持要求

  • Linux (Centos 7+)
  • Java 1.8+

通常,FE服务不会消耗大量的CPU和内存资源。建议您为每个FE节点分配8个CPU内核和16 GB RAM。与FE服务不同,如果您的应用程序需要在大型数据集上处理高度并发或复杂的查询,BE服务可能会使用大量CPU和内存资源。因此,建议您为每个BE节点分配16个CPU内核和64 GB RAM。由于FE节点仅在其存储中维护StarRocks的元数据,因此在大多数场景下,每个FE节点只需要100 GB的HDD存储。推荐BE节点CPU:内存=1:4,BE:FE=2:1

假定内存、磁盘都不会拖后腿的情况下,分析/查询的性能瓶颈在CPU的处理能力。所以通过对CPU的算力要求,来预估集群的数量。

集群需要的总CPU资源

e_core = scan_rows / cal_rows / e_rt * e_qps

场景样例

  1. 数据量:事实表一年3.6亿行数据,大约100万行/天;
  2. 典型查询场景:一个月的事实表数据(3000万)和比较小的的维度表(万级别)做关联,再进行group by、sum等聚合计算;
  3. 期望:响应时间在300ms以内,业务的峰值QPS达180左右。

估算解释

  1. StarRocks的处理能力在“单核1000万~1亿/秒”,此场景有「多表join」和「group by」以及一些表达式函数,相对复杂,所以按照「3000万/s的计算能力」估算,需要3个vCPU:3000万 / 3000万/s / 300ms = 3c。
  2. 并发峰值为180qps,因此需要3 * 180 = 540c,即总共需要540个vCPU。按单台物理机48虚拟核(vCPU)算,理论计算大约需要12台物理机。
  3. 实际POC过程中,用3台物理机16虚拟核进行压力测试,能够在40qps下满足300-500ms的响应时间。最终,线上确定用7台48虚拟核的物理机。所以,还是建议用户要根据实际的业务场景做一下POC测试。

综上:根据POC的测试结果,建议用户搭建3个FE节点每个节点16核64GB内存、7个BE节点每个节点48核152GB内存。

其他说明

  1. 计算业务越复杂、处理中的一行的列数量越多越复杂,每秒能处理的行数就会越少;
  2. 计算中「条件过滤」的效果越好(能过滤掉很多数据),则能处理的行数就会越多(因为内部有一些索引结构,能更快地帮助处理数据);
  3. 不同「表模型」会对处理能力有很大影响,上面是按照「明细模型」估算。其他模型,内部会有一些特殊处理,真实的数据量行数会和用户理解的数据量行数有一些差异;同时,分区/分桶,也会对查询性能有很大影响;(我们有其他相关文档来指导用户如何使用以达到最佳性能)
  4. 对于一些需要扫描大量数据的场景,磁盘的性能也会影响处理能力。需要时,可以使用SSD来加速。

BE存储空间

  • BE节点所需的总存储空间 = 原始数据大小 * 数据副本数 / 数据压缩算法压缩比
  • 原始数据大小 = 单行数据大小 * 总数据行数

主键索引占用内存空间

假设存在主键模型,主键为dt、id,数据类型为DATE(4个字节)、BIGINT(8个字节)。则主键占12个字节。

假设该表的热数据有1000万行,存储为三个副本。

则内存占用的计算方式:(12 + 9(每行固定开销)) * 1000W * 3 * 1.5(哈希表平均额外开销)= 945 M

集群架构

混合型集群

指由一个统一的大集群提供所有大数据服务,所有组件集中安装在同一个集群中,有部署简单、运维方便、易于使用等优点。

但是由于混合型集群集群承载了所有功能,职能繁多,网络带宽、磁盘IO等为集群共享,会因大型离线任务占用大量网络或磁盘IO峰值,对线上业务会造成短暂延迟。

且集群环境较为复杂,有较多对线上业务造成影响的风险。

专用型集群

专用型集群指根据不同的需求与功能职责对集群进行划分,由多个职责不同、硬件隔离的集群组成集群组环境提供服务。

子集群各司其职,根据自身业务最大化利用硬件资源,互相独立互不影响。部署较为复杂,运维难度增加。

专用型集群根据业务与应用场景可以划分如下:

  • 离线计算集群
  • 实时计算集群
  • 数据服务集群
  • GPU深度学习集群
  • 图数据库集群
  • 等等。

节点规划

进行节点角色划分时尽可能遵守以下原则

  • CM监控服务在小集群下可以部署在同一主机,大集群下需要独立部署
  • 集群主节点与子节点独立部署(HDFS/HBase/Yarn),且各自子节点部署在相同主机上
    • 独立部署可以避免子节点大量读写、计算引起IO、CPU、网络等资源阻塞而导致主节点异常甚至宕机
    • HDFS/HBase/Yarn部署在相同主机上可以最大化利用数据本地化特性
    • 如果数据量巨大,而集群存储空间不足的时候忽视以上两点,满足业务需求放在第一位
  • Hive、Hue、Impala、Sentry等服务/元数据服务需要部署在同一主机
    • Hue+Sentry对Hive与Impala进行权限控制的时候需要读取Linux主机的用户与用户组进行判别,如果部署在不同主机上则需要在每个主机上创建相同的用户与用户组
    • 或者使用LDAP进行账号管理
    • 条件允许情况下,独立主机或者压力小的主机
  • Zookeeper尽量使用5个节点,且条件允许下最好在不同的物理主机上
    • 5个节点的zk可以保证leader的快速选举
    • 在不同的物理主机上可以最大限度保证安全
  • 集群高可用保证,针对关键性管理节点,尽可能散落在不同的机器节点上,避免集中在某一两台机器上。

节点规划建议

  1. 分别对增量数据、全量数据进行规模测算。在实际业务中,对于任何一个企业其业务条线都会有其自己的业务增长趋势,这样随着业务规模的增长企业数据量也会不断增长,所以我们在规划集群时可以按照短期业务(1~2年)、中长期(3~5年)业务增长进行规划,而不是在集群初始化时一次性导入的初始数据量。

    比如:我们假设企业的客户量x,所有客户产生的所有业务数据量为y,在构建集群时的初始数据量为c,那么可能的一种企业数据量增长模型为

    y = a * f(x) + b * g(x) + c

    其中,在大部分企业中如果客户量增加一个量级dx,那么其所对应的日志和订单业务数据量可能是客户量的线性模型f(x),而对于交易类业务的数据量可能是客户量的非线性增长模型g(x),比如笛卡尔积模型。

    所以不能通过简单的节点动态增加来调整集群规模和存储计算能力,最好还是通过构建短期业务和中长期业务增长趋势模型来规划集群。

  2. 针对离线和实时场景,需要做好压测(QPS(每秒查询率)=并发数/平均响应时间),评估任务量、任务数以及各个组件在不同配置下的处理性能。

  3. 集群高可用保证,针对关键性管理节点,尽可能散落在不同的机器节点上,避免集中在某一两台机器上。

  4. 对于管理节点与NameNode, ResourceManager, Master, JobManager...等,存储和内存,尽量分配小一点。

  5. 对于存储如HDFS DataNode尽量磁盘分配比较大的机器,选择合适的文件格式,并且进行数据压缩。而对于计算节点,如spark、flink、doris等尽量分配比较高CPU和内存的机器。

  6. 任务提交参数合理优化:消耗内存的分批次提交。堆内存进行估算,设置合理并行度。

    对于计算资源内存的估算没有绝对的标准,需要根据公司使用的计算系统来区分,需根据实际使用的组件(MRFlinkSpark),执行了多少任务,实时任务数量,离线任务数量、算法模型等进行实际估算。

    一般实时任务占用的资源都是固定的,可以根据业务个数估算。离线任务可以根据ETL任务数和任务资源配置情况估算,计算资源离线和实时同时启用的时候不能超过资源90%。实时任务资源占用需要小于50%,如果数据增量为570G,实时任务7407/s的QPS,一分钟窗口,如全量计算444420 * (2 / 1024 / 1024) = 0.85G, 有的设置5分钟窗口,如全量计算大概是4.25G。

    如果离线任务,可以把1/4的数据放到内存,那就就需要47.5G的内存。如果二者同时运行,按照不超过90%来计算,需要57.5G。上面只是单纯的从实时和离线单任务来计算,可以选择64G内存(一般机器比较新的情况下)。

  7. Kafka、ZK、Flume传输数据比较紧密的放在一起。

  8. 客户端尽量放在1到2台服务器上,一是风险隔离,导致集群内部受到不必要的干扰。二是作为跳板机,方便工程师外部访问。

  9. 有依赖关系的尽量放到同一台服务器(例如:HIVE和DS)

  10. CPU和内存的比例一般为1:2, 1:3或者1:4三种,具体分配要重点看有多少线程。CDH集群中的MR任务一般采用1:3,实时接收数据的kafka端一般配比高一些,对于spark、kudu这类需要先缓存到内存再保存到磁盘的组件,一般内存需要设置大一些。

  11. 对于集群中产生的数据可以按照业务中间数据、临时数据、集群的系统日志、集群的预留空间安全系数等来进行规划。业务中间数据和临时数据会分配一定的空间比例,对于集群的预留空间安全系数可以按照当集群的总体规模使用达到80%就需要进行横向扩容,等等。

    笔者曾经在实践中遇到过如下情形,原始的业务数据大概有15T左右,通过多副本存储策略、数据处理过程中产生的大量中间和临时数据、再加上集群需要有预留空间的安全系数等,当时整个集群120T的总空间尽然都不够用,也就是说现实中的业务数据在使用中总体上可能会膨胀好多倍。
    end