改版通知

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

StarRocks高效应用与最佳实践指南

ckckck2025年1月10日11 浏览

Pipeline 参数相关优化

Profile 分析(非 Pipeline 版本)

物化视图及分析优化

BACKUP/RESTORE 操作流程案例文档

实时导入数据 too many tablet version 解决办法

Profile 分析及优化指南(Pipeline 版本 StarRocks 2.3+)

Pipeline 参数优化

关闭/开启 Pipeline

sql 复制代码
SET enable_pipeline_engine = true;

并行度相关参数

sql 复制代码
-- 简称 dop
SET pipeline_dop = 0;

-- INT 含义为可以设置为整数值,简称 instance_num
SET GLOBAL parallel_fragment_exec_instance_num = INT;

并行度参数调整说明

  • StarRocks version <= 2.3 版本(包括 2.3 版本):

    • dop != 0 时,instance_numdop 都会生效;be_num * instance_num * dop = 总并行度
    • dop = 0 时,会忽略 session 变量 instance_num(单 BE instance 个数),会自动调整 dopinstance_num,最终的 dop * instance_num = BE 节点核数一半be_num * instance_num * dop = 总并行度
  • 2.3 版本之后:

    • Pipeline 的 instance_num 永远是 1,会忽略 session 变量的 instance_num,通过 dop 参数进行并行度的调节。

注意:

  • 实际场景中,一个 Fragment 实例的并行数量存在上限,为一张表在一个 BE 中的 Tablet 数量。
  • 在高并发场景下,CPU 资源往往已充分利用,因此建议设置 Fragment 实例的并行数量为 1,以减少不同查询间资源竞争,从而提高整体查询效率。

Profile 分析

背景

我们时常遇到 SQL 执行时间不及预期的情况,为了优化 SQL 达到预期查询时延,我们能够做哪些优化。本文旨在分析查询 Profile 各阶段耗时是否合理以及对应优化方式。

准备

打开 Profile 分析上报。

sql 复制代码
mysql -h ip -P9030 -u root -p xxx
## 该参数开启的是 session 变量,若想开启全局变量可以 set global is_report_success=true; 一般不建议全局开启,会略微影响查询性能
mysql> set is_report_success=true;

该参数会打开 Profile 上报,后续可以查看 SQL 对应的 Profile,从而分析 SQL 瓶颈在哪,如何进一步优化。

如何获取 Profile?

如上设置打开 Profile 上报后,打开 FE 的 HTTP 界面(http://ip:8030),如下点击 queries 后,点击相应 SQL 后的 Profile 即可查看对应信息。

注: 此处需要进 master 的 HTTP 页面。如不确定集群哪台是 master,可以 show frontends 查看 IsMaster 值为 true 的 IP。

Explain 分析

分区分桶

  • partitions 字段 x/xx 表示查询分区/总分区。
  • tabletRatio 字段 x/xx 表示查询分桶/总分桶。

查看对应查询 SQL 是否包含分区字段,是否正确裁剪。如未正常裁剪,确认是否有以下问题:

  • 字段类型不一致
  • 字段有函数,例如:date_format('2009-10-04 22:23:00', '%W %M %Y')

存储层聚合

何时需要存储层聚合?

  • 聚合表的聚合发生在导入、Compaction、查询时。

PREAGGREGATIONOn 表示存储层可以直接返回数据,存储层无需进行聚合。
PREAGGREGATIONOff 表示存储层必须聚合,可以关注下 Off 的原因是否符合预期。

是否命中物化视图

通过执行计划可以看到命中的物化视图的名称以及是否进行了预聚合。

SCAN 环节分析及优化

在 Profile 中搜索 OLAP_SCAN_NODE,会有很多个结果形如 OLAP_SCAN_NODE (id=0),其中 id=x 有多个,表示同一个表的 scan 信息。如下是一个典型的 scan 慢节点。

plaintext 复制代码
OLAP_SCAN_NODE (id=0):(Active: 56s208ms[56208256470ns], % non-child: 0.00%)
    - Table: xxxx
    - Rollup: xxxx
    - Predicates: 3: svrIp = 'xxx.xxx.xxx.xxx'
    - BytesRead: 1.21 GB
    - NumDiskAccess: 0
    - PeakMemoryUsage: 1.24 MB
    - PerReadThreadRawHdfsThroughput: 0.0 /sec
    - RowsRead: 0
    - RowsReturned: 96
    - RowsReturnedRate: 1
    - ScanTime: 56s206ms
    - ScannerThreadsInvoluntaryContextSwitches: 0
    - ScannerThreadsTotalWallClockTime: 0ns
    - MaterializeTupleTime(*): 0ns
    - ScannerThreadsSysTime: 0ns
    - ScannerThreadsUserTime: 0ns
    - ScannerThreadsVoluntaryContextSwitches: 0
    - TabletCount : 1
    - TotalRawReadTime(*): 0ns
    - TotalReadThroughput: 0.0 /sec

数据倾斜问题

查询某张表的 scan 信息,比如上述 test 表对应的 OLAP_SCAN_NODE (id=0),分别检索查看多个 Active: xxxms 信息,观察是否差距很大,如果存在个别节点耗时是其他节点数据量倍数,则有数据倾斜的问题。

关键指标解读

以下对 scan 的关键指标做一些解读,如果对应字段的值占比总查询时间很高,可以针对该阶段进行分析。

  • BytesRead: 读取 tablet 数据量大小,该值太大或者太小表示 tablet 设置的均不合理。
  • RowsReturned: 扫描返回符合要求的行数,如果 BytesRead 很大,而 RowsReturned 很小,但是 scan 占比挺久,可以考虑将过滤条件中的字段建表时设置为 key 列,对点查效果有很好的加速效果。
  • RowsReturnedRate: 结果集返回速率,如果有个别节点返回比较慢,可以查看磁盘读写是否异常,或者 CPU、内存资源是否负载很高,导致系统调度时间增加。
  • TabletCount: tablet 数量,关注此处是否太多或者太少,此处和 bucket 设置息息相关。
  • MERGE: 如果 merge 中的 aggr/union/sort 耗时特别久,则整体瓶颈在底层 rowset 的 merge 上。

Aggregate 环节分析及优化

关键指标解读

  • AggComputeTime: 构建 Hash 表和计算聚合函数的时间。
  • ExprComputeTime: 计算聚合函数内部标量函数的时间。
  • ExprReleaseTime: 内存释放的时间。
  • GetResultsTime: 将 Hash 表的数据转换成 Chunk 的时间。
  • HashTableSize: Hash 表大小。
  • InputRowCount: 聚合前行数。
  • PassThroughRowCount: Streaming 聚合时,没有经过 Hash 表,直接输出的行数。
  • PeakMemoryUsage: 内存使用。
  • StreamingTime: Streaming 聚合时,聚合函数列格式转换的耗时。

JOIN 环节分析及优化

名词解释

关于 HashJoin: 两个阶段,Build + Probe。

Build 阶段: 将其中一个表(一般是较小的那个)中的每一个条经过 Hash 函数的计算都放到不同的 Hash Bucket 中。

Probe 阶段: 对于另外一个表,经过 Hash 函数,确定其所在的 Hash Bucket,然后和上一步构建的 Hash Bucket 中的每一行进行匹配,如果匹配到就返回对应的行。

Join 左右表调整: StarRocks 是用右表构建 Hash 表,所以右表应该是小表,StarRocks 可以基于 cost 自动调整左右表顺序,也会自动把 left join 转 right join。

Join 多表 Reorder: 多表 Join 如何选择出正确的 Join 顺序,是 CBO 优化器的核心,当 Join 表的数量小于等于 5 时,StarRocks 会基于 Join 交换律和结合律进行 Join Reorder,大于 5 时,StarRocks 会基于贪心算法和动态规划进行 Join Reorder。

Join 分布式执行选择:

  • BroadCast Join: 将右表全量发送到左表的 HashJoinNode。
  • Shuffle Join: 将左右表的数据根据哈希计算分散到集群的节点之中。
  • Colocate Join: 两个表的数据分布都是一样的,只需要本地 join 即可,没有网络传输开销。
  • Bucket Shuffle Join: join 的列是左表的数据分布列(分桶键),所以相比于 shuffle join 只需要将右表的数据发送到左表数据存储计算节点。
  • Replicated Join: 右表的全量数据是分布在每个节点上的(也就是副本个数和 BE 节点数量一致),不管左表怎么分布,都是走本地 Join。没有网络传输开销。

RuntimeFilter: 基本原理是通过在 join 操作之前提前过滤掉那些不会命中 join 的输入数据来大幅减少 join 中的数据传输和计算,从而减少整体的执行时间。

分析和优化

是否有收集统计信息?

1.19.0 及其之后的版本默认开启了 CBO,之前的版本如果没有开启 CBO,可能就会没有统计信息。

如果下面的 SQL 有查询结果,表示有统计信息收集(table_name 为参与 join 的表名)。如果没有查询结果,可参考 analyze 命令手动触发统计信息收集。

sql 复制代码
select * from _statistics_.table_statistic_v1 where table_name like '%table_name';

怎么判断瓶颈点?

下图表示的是合理的方式,右表是小表,采用的 broadcast join 方式。

下图表示的是不合理的方式,右表是大表,采用的 broadcast join 方式,会将右表的数据拷贝 BE 数量 * parallel_fragment_exec_instance_num(并行度)份,导致 JOIN 节点的右子节点的 EXCHANGE 节点花费很多的执行时间。

下图表示的是不合理的方式,两个表数据量相差比较大,现在采用的是 shuffle join(两个孩子节点都是 EXCHANGE NODE),这种情况下建议可以尝试采用 broadcast join,调整下左右表顺序,小表在左边。例如加 hint 方式:

sql 复制代码
select 右表.x, 左表.y from 右表 join [broadcast] 左表 on 左表.x1 = 右表.x1

常见优化方法

当前启用了 CBO 优化器,一般情况下不需要人为触发优化,不过在一些场景下可以采用下面的方法尝试优化下:

  • join condition 的列,更应该使用 int、DATE 等简单类型。
  • 在 join 之前尽量添加一些 where 条件,能够充分发挥谓词下推,减少后续的数据 shuffle 和 join 节点处理的数据量。
  • 大表 join,能够使用 colocate join 的尽量使用,能够减少网络传输,极大的提升性能。
  • 大小表 join,左右表顺序有问题,可以通过 [broadcast] hint 方式调整小表为右表方式。例如:select a.x, b.y from a join [broadcast] b on a.x1 = b.x1
  • 两个相差不多的表(一般几百 k 行)join,有些情况下默认会选用 broadcast join,这个时候可以尝试采用 [shuffle] hint 的方式强制走 shuffle join。例如:select a.x, b.y from a join [shuffle] b on a.x1 = b.x1

物化视图及分析优化

简介

物化视图采用空间换时间的设计思路,一张表可以创建多个物化视图,查询会自动命中最优的物化视图。物化视图不能通过名称直接查询,但是其在底层存储时与一般表无异,创建物化视图后基表中的数据会以异步的方式同步到其所有物化视图中。

目前物化视图支持的两种场景,当然也支持两种方式混合:

  • 预聚合:对明细表的任意维度组合进行预先聚合。
  • 维度列变序:采用新的维度列的排序方式,以便命中前缀查询条件。

注意:

  • 创建了大量的物化视图,会导致数据导入速度过慢,并且部分物化视图的相互重复,查询频率极低,会有较高的查询延迟。
  • 只有明细模型和聚合模型支持创建物化视图,主键模型和更新模型不支持创建物化视图。

案例

以下 SQL 以 SSB 1T 测试数据集(lineorder_flat 474G)测试 SQL 为例子创建物化视图优化。

以 sum 函数聚合为例

sql 复制代码
-- 原始 SQL
select LO_ORDERDATE, sum(LO_QUANTITY) from lineorder_flat group by LO_ORDERDATE;

-- 创建物化视图
CREATE MATERIALIZED VIEW sum_mv as select LO_ORDERDATE, sum(LO_QUANTITY) from lineorder_flat group by LO_ORDERDATE;

原始 SQL 的执行时间为:27.61 sec。

创建物化视图后执行时间为:0.96 sec。

通过创建物化视图可以减少数据扫描量实现对查询的加速。

注意: 物化视图创建过程为异步过程,数据量越大耗时越久,通过命令可以查看创建进度:SHOW ALTER MATERIALIZED VIEW FROM databaseName;

BACKUP/RESTORE 操作文档

说明

本文介绍如何使用 BACKUP/RESTORE 功能来进行备份及恢复操作,举例流程可以作为辅助参考,详情请参考备份与恢复的文档。

总体流程

  1. 先创建云端仓库用于备份与恢复(新老集群都要创建云端仓库,REPOSITORY 名字要相同,BROKER Name 需要对应集群的 broker 名称)。
  2. 在老集群准备好需要进行迁移备份的表,Backup 到云端仓库。
  3. 再从云端仓库 Restore 到新集群。

新集群当中不用事先创建好需要备份恢复的表,因为在进行 Restore 操作会自动创建。

具体步骤

创建 REPOSITORY 远端仓库

通过 CREATE REPOSITORY 创建远端仓库,该语句用于创建仓库。仓库用于属于备份或恢复。仅 root 或 superuser 用户可以创建仓库。

sql 复制代码
CREATE [READ ONLY] REPOSITORY `repo_name`
WITH [BROKER `broker_name`]
ON LOCATION `repo_location`
PROPERTIES ("key"="value", ...);

说明:

  • 仓库的创建,依赖于已存在的 broker 或者其他协议访问云存储。
  • 根据对象存储的类型不同,PROPERTIES 有所不同,具体见示例。
  • BROKER 名称可以通过 show broker 来进行查看集群的 Broker 名称。

根据情况可以创建不同类型的云端数据仓库,用于备份或恢复,目前支持的数据源有:

  1. TenCent Cos: 腾讯云对象存储。
  2. Apache HDFS: 社区版本 hdfs。
  3. Amazon S3: Amazon 对象存储。
  4. Aliyun Oss: 阿里云对象存储。

因为我们会选择不同的第三方云端来进行创建所需的备份仓库,所以需要第三方库的连接认证,这里需要我们修改 PROPERTIES 里的配置参数来进行连接,参见其中 properties 的部分,具体见示例:

sql 复制代码
## Cos 形式建库语句:
CREATE REPOSITORY `Cos_RepositoryName`
WITH BROKER `brokerName`
ON LOCATION "cosn://*****************"
PROPERTIES
(
    "fs.cosn.userinfo.secretId" = "*************************",
    "fs.cosn.userinfo.secretKey" = "******************************",
    "fs.cosn.bucket.endpoint_suffix" = "cos.ap-beijing.myqcloud.com"
);

## HDFS 形式建库语句:
CREATE REPOSITORY `hdfs_RepositoryName`
WITH BROKER `hdfs_broker`
ON LOCATION "hdfs://hadoop-name-node:prot/*******/******/******/"
PROPERTIES
(
    "username" = "user",
    "password" = "password"
);

## S3 形式建库语句:
CREATE REPOSITORY `s3_RepositoryName`
WITH BROKER `hdfs_broker`
ON LOCATION "s3a://xxx"
PROPERTIES
(
    "fs.s3a.access.key" = "xxx",
    "fs.s3a.secret.key" = "yyy",
    "fs.s3a.endpoint" = "s3-ap-northeast-1.amazonaws.com"
);

## Oss 形式建库语句:
CREATE REPOSITORY `Oss_RepositoryName`
WITH BROKER `brokerName`
ON LOCATION "oss://**************"
PROPERTIES
(
    "fs.oss.accessKeyId" = "xxxxxxxxxxxxxxxxxxxxxxxxxx",
    "fs.oss.accessKeySecret" = "yyyyyyyyyyyyyyyyyyyy",
    "fs.oss.endpoint" = "oss-cn-shenzhen-internal.aliyuncs.com"
);

说明:

  • 一个集群可以创建多个仓库。只有拥有 ADMIN 权限的用户才能创建仓库。
  • 任何用户都可以通过 SHOW REPOSITORIES; 命令查看已经创建的仓库。
  • 在做数据迁移操作时,需要在源集群和目的集群创建完全相同的仓库,以便目的集群可以通过这个仓库,查看到源集群备份的数据快照。
  • 如果需要删除已经创建好的仓库,可以参考下方的文档指令删除云端仓库。

备份数据到云端仓库

1. 该语句用于备份指定数据库下的数据。 该命令为异步操作。提交成功后,需通过 SHOW BACKUP 命令查看进度。仅支持备份 OLAP 类型的表。

sql 复制代码
BACKUP SNAPSHOT [db_name].{snapshot_name}
TO `repository_name`
[ON|EXCLUDE] (
    `table_name` [PARTITION (`p1`, ...)],
    ...
)
PROPERTIES ("key"="value", ...);

说明:

  • 同一数据库下只能有一个正在执行的 BACKUP 或 RESTORE 任务。
  • 备份操作会备份指定表或分区的基础表及物化视图,并且仅备份一副本。
  • 备份操作的效率取决于数据量、Compute Node 节点数量以及文件数量。备份数据分片所在的每个 Compute Node 都会参与备份操作的上传阶段。节点数量越多,上传的效率越高,文件数据量只涉及到的分片数,以及每个分片中文件的数量。如果分片非常多,或者分片内的小文件较多,都可能增加备份操作的时间。
  • ON 子句中标识需要备份的表和分区。如果不指定分区,则默认备份该表的所有分区。
  • EXCLUDE 子句中标识不需要备份的表和分区。备份除了指定的表或分区之外这个数据库中所有表的所有分区数据。
  • PROPERTIES 目前支持以下属性:
    • "type" = "full": 表示这是一次全量更新(默认,当前仅支持 full)。
    • "timeout" = "3600": 任务超时时间,默认为一天。单位秒,最小只能调整为 10min。

2. 两种备份的案例为:

  1. 全量备份 example_db 下的表 example_tbl 到仓库 example_repo 中:
sql 复制代码
BACKUP SNAPSHOT example_db.snapshot_label1
TO example_repo
ON (example_tbl)
PROPERTIES ("type" = "full");
  1. 全量备份 example_db 下,表 example_tbl 的 p1, p2 分区,以及表 example_tbl2 到仓库 example_repo 中:
sql 复制代码
BACKUP SNAPSHOT example_db.snapshot_label2
TO example_repo
ON
(
    example_tbl PARTITION (p1,p2),
    example_tbl2
);

说明:

  • 如果需要取消正在执行的 BACKUP 任务,可以参考下方的文档指令取消 BACKUP 任务。

查看云端仓库任务

可以通过 SHOW BACKUP FROM example_db; 获得 example_db 下最后一次 BACKUP 任务,查看云端仓库中已有的备份,该记录中仅显示最近一次的 BACKUP 任务。

恢复云端数据到集群

确认完快照信息后到新集群进行恢复操作,该语句用于将之前通过 BACKUP 命令备份的数据,恢复到指定数据库下。该命令为异步操作。提交成功后,需通过 SHOW RESTORE 命令查看进度。仅支持恢复 OLAP 类型的表。

sql 复制代码
RESTORE SNAPSHOT [db_name].{snapshot_name}
FROM `repository_name`
[ON|EXCLUDE] (
    `table_name` [PARTITION (`p1`, ...)] [AS `tbl_alias`],
    ...
)
PROPERTIES ("key"="value", ...);

说明:

  • 同一数据库下只能有一个正在执行的 BACKUP 或 RESTORE 任务。
  • ON 子句中标识需要恢复的表和分区。如果不指定分区,则默认恢复该表的所有分区。所指定的表和分区必须已存在于仓库备份中。
  • EXCLUDE 子句中标识不需要恢复的表和分区。除了所指定的表或分区之外仓库中所有其他表的所有分区将被恢复。
  • 可以通过 AS 语句将仓库中备份的表名恢复为新的表。但新表名不能已存在于数据库中。分区名称不能修改。
  • 可以将仓库中备份的表恢复替换数据库中已有的同名表,但须保证两张表的表结构完全一致。表结构包括:表名、列、分区、Rollup 等等。
  • 可以指定恢复表的部分分区,系统会检查分区 Range 或者 List 是否能够匹配。
  • PROPERTIES 目前支持以下属性:
    • "backup_timestamp" = "2022-06-01-12-09-14": 指定了恢复对应备份的哪个时间版本,必填。查看 backup_timestamp 语句为: SHOW SNAPSHOT ON 云端仓库名;
    • "replication_num" = "3": 指定恢复的表或分区的副本数。默认为 3。若恢复已存在的表或分区,则副本数必须和已存在表或分区的副本数相同。同时,必须有足够的 host 容纳多个副本。
    • "timeout" = "3600": 任务超时时间,默认为一天。单位秒。

两种恢复的案例为:

  1. 从 example_repo 中恢复备份 snapshot_1 中的表 backup_tbl 到数据库 example_db1,时间版本为 "2022-06-01-12-09-14"。恢复为 1 个副本:
sql 复制代码
RESTORE SNAPSHOT example_db1.`snapshot_1`
FROM `example_repo`
ON (`backup_tbl`)
PROPERTIES
(
    "backup_timestamp"="2022-06-01-12-09-14",
    "replication_num" = "1"
);
  1. 从 example_repo 中恢复备份 snapshot_2 中的表 backup_tbl 的分区 p1,p2,以及表 backup_tbl2 到数据库 example_db1,并重命名为 new_tbl,时间版本为 "2022-06-01-12-09-14"。默认恢复为 3 个副本:
sql 复制代码
RESTORE SNAPSHOT example_db1.`snapshot_2`
FROM `example_repo`
ON
(
    `backup_tbl` PARTITION (`p1`, `p2`),
    `backup_tbl2` AS `new_tbl`
)
PROPERTIES
(
    "backup_timestamp"="2022-06-01-12-09-14"
);

说明:

  • 同一数据库下只能有一个正在执行的恢复操作。
  • 可以将仓库中备份的表恢复替换数据库中已有的同名表,但须保证两张表的表结构完全一致。表结构包括:表名、列、分区、物化视图等等。
  • 当指定恢复表的部分分区时,系统会检查分区范围是否能够匹配。
  • 恢复操作的效率:在集群规模相同的情况下,恢复操作的耗时基本等同于备份操作的耗时。如果想加速恢复操作,可以先通过设置 replication_num 参数,仅恢复一个副本,之后在通过调整副本数:ALTER TABLE 将副本补齐。
  • 如果需要取消正在执行的 RESTORE 任务,可以参考下方的文档指令取消 RESTORE 任务。

查看最近恢复的 job

sql 复制代码
SHOW RESTORE FROM 数据库名;

可简单执行 count 语句进行新老集群表数据条数对比,或者运用 sum 函数相加下后面的值看是否相等,来进行校验恢复是否成功。

其他

  1. 删除云端仓库
sql 复制代码
DROP REPOSITORY `repo_name`;

说明:

  • 删除仓库,仅仅是删除该仓库在 StarRocks 中的映射,不会删除实际的仓库数据。删除后,可以再次通过指定相同的 broker 和 LOCATION 映射到该仓库。
  1. 取消 BACKUP 任务
sql 复制代码
CANCEL BACKUP FROM db_name;
  1. 取消 RESTORE 任务
sql 复制代码
CANCEL RESTORE FROM db_name;

注意:

  • 当取消处于 COMMIT 或之后阶段的恢复左右时,可能导致被恢复的表无法访问。此时只能通过再次执行恢复作业进行数据恢复。

导入时报错 too many tablet version 解决方式

背景

实时导入过程中经常遇到导入失败报错 too many tablet version 的情况,此原因是因为 StarRocks 内部,表级别 version 数量有 1000 的限制,通常见于数据写入频次过高,数据合并来不及从而堆积大量的 version。本文主要介绍此种情况下如何恢复导入任务。

现象

导入失败时检索 BE 日志:

bash 复制代码
grep 'too many tablet versions' /data/xxx/starrocks/be/log/be.WARNING

恢复步骤

在日志里搜索所示字样:

bash 复制代码
grep 'Fail to init delta writer' /data/xxx/starrocks/be/log/be.WARNING

检索结果会有 tablet=xxx 的信息,这里以 tablet 716097 为例。

连接 StarRocks 然后查询 tablet 所属表信息:

sql 复制代码
show tablet 716097;

如下所示,DbNameTableName 分别表示对应 database 和所属表。

此时需要调整写入该表的任务间隔、并发、单次写入量:核心在于通过增大单次写入数据量,减少任务提交次数以及写入并发。

以 RoutineLoad 为例:

参考文档:RoutineLoad

sql 复制代码
max_batch_interval: 每个子任务最大执行时间,单位是「秒」。范围为 5 到 60。默认为 10。1.15 版本后: 该参数是子任务的调度时间,即任务多久执行一次,任务的消费数据时间为 fe.conf 中的 routine_load_task_consume_second,默认为 3s,任务的执行超时时间为 fe.conf 中的 routine_load_task_timeout_second,默认为 15s。
max_batch_rows: 每个子任务最多读取的行数。必须大于等于 200000。默认是 200000。1.15 版本后: 该参数只用于定义错误检测窗口范围,窗口的范围是 10 * max-batch-rows。
max_batch_size: 每个子任务最多读取的字节数。单位是「字节」,范围是 100MB 到 1GB。默认是 100MB。1.15 版本后: 废弃该参数,任务消费数据的时间为 fe.conf 中的 routine_load_task_consume_second,默认为 3s。

DataX 相应参数:

参考文档:DataX-starrocks-writer

sql 复制代码
maxBatchRows: 描述:单次 StreamLoad 导入的最大行数。必选:否。默认值:500000 (50W)。
maxBatchSize: 描述:单次 StreamLoad 导入的最大字节数。必选:否。默认值:104857600 (100M)。
flushInterval: 描述:上一次 StreamLoad 结束至下一次开始的时间间隔(单位:ms)。必选:否。默认值:300000 (ms)。

Flink 相应参数:

参考文档:Flink-connector-starrocks

sql 复制代码
sink.buffer-flush.max-rows: 单次刷新的最大行数。默认值:100000。
sink.buffer-flush.interval: 刷新间隔时间,单位毫秒。默认值:300000。
sink.max-retries: 最大重试次数。默认值:3。

Profile 分析优化指南

背景

我们时常遇到 SQL 执行时间不及预期的情况,为了优化 SQL 达到预期查询时延,我们能够做哪些优化。本文旨在分析查询 Profile 各阶段耗时是否合理以及对应优化方式。

准备

打开 Profile 分析上报。

sql 复制代码
mysql -h ip -P9030 -u root -p xxx
## 该参数开启的是 session 变量,若想开启全局变量可以 set global is_report_success=true; 一般不建议全局开启,会略微影响查询性能
mysql> set is_report_success=true;

该参数会打开 Profile 上报,后续可以查看 SQL 对应的 Profile,从而分析 SQL 瓶颈在哪,如何进一步优化。

如何获取 Profile?

如上设置打开 Profile 上报后,打开 FE 的 HTTP 界面(http://ip:8030),如下点击 queries 后,点击相应 SQL 后的 Profile 即可查看对应信息。

注: 此处需要进 master 的 HTTP 页面。如不确定集群哪台是 master,可以 show frontends 查看 IsMaster 值为 true 的 IP。

Explain 分析

Explain SQL 获取执行计划,如下:

sql 复制代码
EXPLAIN SELECT * FROM table_name;

分区分桶

上图中 partitions 字段 x/xx 表示查询分区/总分区,tabletRatio 字段 x/xx 表示查询分桶/总分桶。

查看对应查询 SQL 是否包含分区字段,是否正确裁剪。如未正常裁剪,确认是否有以下问题:

  • 字段类型不一致
  • 字段有函数,例如:date_format('2009-10-04 22:23:00', '%W %M %Y')

存储层聚合

何时需要存储层聚合?

  • 聚合表的聚合发生在导入、Compaction、查询时。

PREAGGREGATIONOn 表示存储层可以直接返回数据,存储层无需进行聚合。
PREAGGREGATIONOff 表示存储层必须聚合,可以关注下 Off 的原因是否符合预期。

是否命中物化视图

通过执行计划可以看到命中的物化视图的名称以及是否进行了预聚合。

关于 Pipeline Profile

在 pipeline 执行引擎中,查询的 profile 的结构如下,总共有五个层级,分别是:

  • Fragment
  • FragmentInstance
  • Pipeline
  • PipelineDriver
  • Operator

Pipeline 的 dop 自适应策略会保证 FragmentInstance * Dop = 核数的一半,对于复杂查询,比如 TPC-DS 中的查询,其产生的 profile 多达几十万行,除非借助分析脚本,肉眼很难直接分析。

因此,profile 在 2.3 版本即 pipeline 默认打开的版本(StarRocks 2.3+)进行了简化处理:

  1. 压缩 profile 的整体大小,整个 profile 的行数最好不好超过 100 行。对于复杂 SQL 而言,例如 tcp-ds,尽量不要超过 500 行。简化 profile 的另一个好处是,可以减少 BE-FE 之间传输的 profile 的开销。
  2. 对指标进行分类,突出核心指标,方便进行问题排查。
  3. 给出 fragment 维度的时间线以及 pipeline 维度的时间线,便于分析 fragment 之间或者 pipeline 之间的依赖关系。

Profile 级别

提供新的 session 变量 pipeline_profile_level,总共包含 3 个层级:

  • Level 0: 合并同构 profile,只包含几个核心指标,包括算子维度和 Pipeline 维度。
  • Level 1: 合并同构 profile,保留所有指标。默认级别。
  • Level 2: 保留所有层级的 profile,不做任何简化。

核心指标关系图

plaintext 复制代码
Pipeline::Active = Σ Operator::OperatorTotalTime + Pipeline::OverheadTime
InputEmptyTime = FirstInputTime + FollowupInputEmptyTime
PendingTime = InputEmptyTime + OutputFullTime + PreconditionBlockTime
DriverTotalTime = ActiveTime + PendingTime

指标含义说明

对于 pipeline_profile_level=0pipeline_profile_level=1,指标会进行合并,合并方式取决于类型:

  • 时间类型:求平均。
  • 非时间类型:求和。

PS: 对于耗时类型的合并操作(求均值),若均值和最大最小值偏差太大,会额外给出该指标的最大值和最小值。

计算公式如下:

plaintext 复制代码
Pipeline::ActiveTime - ∑Operator::OperatorTotalTime
PendingTime = InputEmptyTime + OutputFullTime + PreconditionBlockTime
InputEmptyTime = FirstInputEmptyTime + FollowupInputEmptyTime

各时间关系如下:

  • PendingTime: Pipeline 在 Pending 队列中的时间。可以细分为 InputEmptyTime、OutputFullTime、PreconditionBlockTime。
  • InputEmptyTime: 由于输入队列为空导致的等待时间。
    • FirstInputEmptyTime: 第一次由于输入队列为空导致的等待的时间。单独把第一次提出来是因为,第一次等待,大概率是由于 Pipeline 的依赖关系产生的。
    • FollowupInputEmptyTime: 后续(第二次开始)所有因为输入队列为空导致的等待的时间。
  • OutputFullTime: 由于输出队列为空导致的等待时间。
  • PreconditionBlockTime: 由于 Pipeline 依赖关系导致的等待时间。
  • ScheduleTime: 调度时间。从 pending 队列中移出,放入就绪队列,并被成功调度的这段时间。

SCAN 环节分析及优化

在 Profile 中搜索 OLAP_SCAN,如下是一个典型的 scan 慢节点。

plaintext 复制代码
OLAP_SCAN (plan_node_id=0):
    CommonMetrics:
        - CloseTime: 362.301us
        - OperatorTotalTime: 4s952ms
        - __MAX_OF_OperatorTotalTime: 7s914ms
        - __MIN_OF_OperatorTotalTime: 1s140ms
        - PeakMemoryUsage: 0.00
        - PullChunkNum: 1.466019M (1466019)
        - PullRowNum: 5.999989425B (5999989425)
        - PullTotalTime: 4s952ms
        - __MAX_OF_PullTotalTime: 7s914ms
        - __MIN_OF_PullTotalTime: 1s140ms
        - PushChunkNum: 0
        - PushRowNum: 0
        - PushTotalTime: 0ns
        - SetFinishedTime: 353ns
        - SetFinishingTime: 115ns
    UniqueMetrics:
        - Rollup: lineitem
        - Table: lineitem
        - BytesRead: 0.00
        - CachedPagesNum: 0
        - CompressedBytesRead: 2.25 GB
        - CreateSegmentIter: 44.855us
        - IOTime: 6s382ms
        - __MAX_OF_IOTime: 8s361ms
        - __MIN_OF_IOTime: 4s759ms
        - PushdownPredicates: 0
        - RawRowsRead: 5.999989425B (5999989425)
        - ReadPagesNum: 366.588K (366588)
        - RowsRead: 5.999989425B (5999989425)
        - ScanTime: 7s917ms
        - __MAX_OF_ScanTime: 10s725ms
        - __MIN_OF_ScanTime: 2s67ms
        - SegmentInit: 1s220ms
        - BitmapIndexFilter: 0ns
        - BitmapIndexFilterRows: 0
        - BloomFilterFilterRows: 0
        - ShortKeyFilterRows: 0
        - ZoneMapIndexFilterRows: 0
        - __MAX_OF_SegmentInit: 2s30ms
        - __MIN_OF_SegmentInit: 722.290ms
        - SegmentRead: 6s684ms
        - BlockFetch: 56.340ms
        - __MAX_OF_BlockFetch: 71.932ms
end