Day 3|StarRocks 架构总览:FE、BE、MPP、列存、向量化与 Pipeline
前两天已经建立了数据系统的坐标系,也明确了 StarRocks 的定位、场景和边界。今天开始进入 StarRocks 内部,沿着一条 SQL 的执行链路,认识 FE、BE、MPP、Fragment、Instance、Exchange、列式存储、向量化和 Pipeline。
学完这一篇,你应当能够把“StarRocks 为什么快”拆成若干可验证的环节,也能够在查询变慢时先判断问题位于规划、扫描、计算、网络交换,还是结果返回。文中的计划图属于教学示意,实验脚本尚待在隔离环境运行,不作为实测输出。
本日学习目标
完成今天的学习后,你应当能够:
- 说明 FE 与 BE 的职责边界,以及 Leader、Follower、Observer 的作用;
- 描述 SQL 从客户端进入 FE,到多个 BE 并行执行,再返回结果的完整链路;
- 解释 Fragment、Instance、Exchange 之间的关系;
- 区分 MPP、向量化执行和 Pipeline,各自解决哪一层问题;
- 理解列存、编码压缩、CPU Cache、SIMD 与多核调度如何协同;
- 使用
EXPLAIN 和 Query Profile 对慢查询做第一轮分层定位;
- 完成一条 JOIN 加聚合 SQL 的执行链路注释图。
一、先建立一个正确的架构视角
学习数据库架构时,最容易陷入两种状态。
第一种状态是背组件名:FE 负责元数据,BE 负责存储和计算。结论没有错,但它解释不了查询为什么要经过多个阶段,也无法帮助我们定位生产问题。
第二种状态是过早钻进源码:先研究某个 Java 类、C++ 算子或线程池。实现细节会随版本变化,学习者很快失去全局方向。
更稳妥的方式是先建立三条主线。
控制路径回答“这条 SQL 应该怎样执行”。连接、鉴权、语义分析、优化、切分、调度都属于这条路径,核心工作主要发生在 FE。
数据路径回答“数据从哪里读取,经过哪些算子,如何在节点间流动”。扫描、过滤、Join、聚合、排序、Exchange 和结果输出主要发生在 BE。
资源路径回答“执行过程中消耗了什么”。磁盘与对象存储提供数据,内存承载中间状态,CPU 执行表达式和算子,网络承担 Shuffle,客户端与代理负责拉取结果。
观察一条 SQL 的耗时,可以沿着下面几个环节记录时间:
客户端建立连接、提交请求
→ FE 解析、语义分析、优化与计划生成
→ 查询准入、调度与任务下发
→ BE 扫描、计算与 Exchange 协作执行
→ 结果输出与客户端读取
阶段可能重叠,排队位置也随执行路径变化。多台 BE 的累计算子、I/O 和网络时间不能直接相加,作为客户端总耗时。应结合墙钟时间、依赖与等待原因诊断,再把后续的存储、优化器和资源治理知识放回这三条主线中。[8][13]
本文以 FE + BE 存算一体架构建立概念,关键源码核对到 StarRocks 4.1.4。存算分离由 FE、CN 与远端存储协作,CN 负责计算和缓存,持久数据位于对象存储或 HDFS。Day 4 会专门讨论两种架构的选择,今天先把查询引擎的共同核心讲清楚。[1]
二、FE 与 BE:控制面和数据面的职责边界
2.1 FE:理解 SQL、管理元数据、生成计划、协调执行
FE 是 Frontend 的缩写,主要使用 Java 实现。客户端通常通过 MySQL 协议连接 FE,MySQL Client、JDBC 和 BI 工具可以沿用熟悉的接入方式;具体 SQL、类型与驱动行为仍需核对。[1][5]
FE 主要承担五组工作。
1. 协议、连接与会话
FE 接收连接,完成认证,创建会话,管理数据库上下文、时区、SQL Mode、资源组、超时和各类 Session Variable。一个看似简单的查询,首先要在这里确认用户身份、会话状态和执行环境。
生产环境可以在多个 FE 前配置统一入口,帮助新连接选择可用节点。已有会话和在途查询不会自动迁移。连接池应配置重试与超时,并区分查询、写入和事务的重试语义,避免重连风暴。
2. 元数据管理
数据库、表、分区、Tablet、副本、用户、权限、Catalog 和作业状态都属于元数据。查询计划必须知道表有哪些列、分区在哪里、Tablet 分布在哪些 BE、哪些副本健康,才能决定扫描范围与执行位置。
FE 使用 BDB JE 管理元数据日志的复制与高可用,其他 FE 通过日志回放维护元数据。理解日常分工时,可以区分以下三种角色:[2]
| 角色 |
主要职责 |
是否参与选举 |
| Leader |
处理元数据写入,协调集群状态;由可选举节点中产生 |
是 |
| Follower |
回放元数据日志,可参与选主,也可接收查询 |
是 |
| Observer |
回放元数据日志,扩展查询接入与规划能力 |
否 |
在正常的多数派确认配置下,选主和元数据写入依赖有选举资格节点的有效多数,Leader 计入其中,Observer 不增加投票数。FE 元数据副本与 BE 业务数据副本各有职责;查询可由其他 FE 接收,元数据写入按协调路径交给 Leader 处理。[2]
3. 解析、绑定与语义检查
SQL 文本进入 FE 后,会先被解析成语法结构。随后进行表列绑定、类型推导、函数解析、权限检查和语义校验。
例如,假定下面涉及的表和列存在,并且当前会话启用了 ONLY_FULL_GROUP_BY:
SELECT region, SUM(amount)
FROM sales.orders
GROUP BY city;
在上述条件下,region 未出现在 GROUP BY 中,也未用于聚合表达式,语义检查会拒绝查询。关闭相应 SQL Mode 后,行为可能不同。FE 还要检查权限;权限不足的请求不会作为正常查询下发。本节保留这个错误示例,不提供未经运行的错误输出。[3]
这一阶段的产物已经具备明确的表、列、类型和权限含义,可以继续进入优化器。
4. 查询优化
StarRocks 使用自己的 CBO,采用 Cascades 风格的优化框架。优化过程包含规则改写和基于代价的选择,关键路径可以在 QueryOptimizer 的逻辑改写、Memo 优化和物理计划提取中对应起来。[4]
规则改写适合处理确定性的变换,例如常量折叠、谓词下推、投影裁剪、子查询改写和无效表达式消除。基于代价的优化依赖统计信息,评估不同 Join 顺序、Join 方法、数据分布和并行方式的成本。
统计信息通常包含行数、列大小、Min/Max、NULL 比例、NDV,也可以包含直方图等分布信息。NDV 表示不同值数量,对估算过滤选择率和 Join 基数很重要。统计信息过期时,优化器可能把大表当成小表,选择不合适的 Broadcast,或低估 Hash Join Build 侧的内存需求。[4]
5. Coordinator 与调度
FE 在计划生成与调度阶段结合数据分布、可用副本、并行设置和节点状态,组织并下发执行任务。承担本次查询协调工作的 Coordinator 跟踪实例状态、处理错误,并从结果端获取数据交给客户端;普通扫描、Join、全局聚合与排序仍在 BE 执行。[5]
这里要分清两个概念:FE 是集群组件,Coordinator 是某次查询中的协调角色。一台 FE 可以同时协调很多查询;同一个查询的主要计算则分散在多台 BE。
2.2 BE:保存数据、执行算子、交换数据、产生结果
BE 是 Backend 的缩写,使用 C++ 实现。在存算一体模式中,每台 BE 同时拥有本地存储、内存和 CPU,承担数据副本与查询执行。
BE 的工作可以拆成六组。
1. Tablet 与副本存储
StarRocks 将原生表的数据划分为 Tablet。在存算一体中,同一逻辑 Tablet 可以有多个副本;每个副本维护相关 Rowset 与 Segment 文件。Rowset 表示一批数据形成的文件集合,Segment 按列组织数据与索引。[6]
这套层次会在后续存储引擎课程中详细展开。今天只需要记住:FE 保存“数据分布地图”,BE 保存实际数据文件。一次普通扫描会选择适合的副本读取,同一份业务数据不会因为存在多个副本而被重复计入结果。[1][6]
2. Scan 扫描
Scan 读取 Segment,按可用的键范围、摘要和索引缩小候选数据,读取所需列,完成必要的解压、解码和谓词过滤。Primary Key 表还要依据相应版本的 DelVector 排除失效行。具体处理路径随表模型、谓词和扫描方式变化。[6]
扫描速度取决于读取字节数、命中索引、磁盘吞吐、文件数量、缓存状态和数据分布。很多慢查询表面上是 CPU 高,根源却是扫描范围过大。
3. Join、聚合、排序与表达式计算
过滤后的列数据以 Chunk 为主要批次结构,进入 Join、Aggregation、Sort、Window、TopN 等算子。Hash Join 需要构建哈希表,Aggregation 需要维护分组状态,排序可能需要较大内存。这些算子会形成查询的主要 CPU 与内存消耗。[7]
4. Exchange 与网络传输
当数据需要跨 Instance 或跨节点重新分布时,Exchange 的发送端和接收端负责传递 Chunk。不同计划可能采用 Broadcast、按键分发或 Gather;同节点之间还可以走本地通路,Instance 内也有 Local Exchange。看到 Exchange 不能直接认定发生了同等规模的跨机网络传输。[8][13]
Exchange 的成本包括序列化、网络传输、接收队列、反序列化和下游等待。数据倾斜会让某几个接收端承担过多数据,整个查询最终被最慢实例拖住。
5. Pipeline 调度
BE 通过 Pipeline Driver 推进可执行的算子链。Driver 在输入可用、输出可接收时运行;输入为空、输出阻塞或前置依赖未满足时,交由调度器等待或重新调度。执行线程、扫描 I/O 任务与网络处理有各自的职责,不能把整个 BE 简化为一个固定线程池。[7]
6. 局部结果与返回
BE 可以先完成局部聚合、局部 TopN 或局部过滤,再把较小的结果发送到后续 Fragment。普通查询的最终结果由结果端输出,FE 的 Coordinator 通过 ResultReceiver 获取批次,再沿客户端连接返回。[5]
局部计算是分布式分析性能的重要来源。假设三台 BE 各处理互不重复的一亿行,并有效汇总成几十个分组,网络传输就能明显缩小。实际输出还受聚合策略与批次影响,需要检查计划和运行状态。
三、一条 SQL 进入 FE 后经历什么
为了让流程更具体,使用下面这条查询作为主线。假定 dim_city.city_id 在业务上唯一,order_date 为 DATE,amount 为精确数值类型;维表出现重复键会放大 Join 结果,必须单独检查。
SELECT
c.region,
SUM(o.amount) AS total_amount
FROM fact_orders o
JOIN dim_city c
ON o.city_id = c.city_id
WHERE o.order_date >= '2026-08-01'
GROUP BY c.region
ORDER BY total_amount DESC;
3.1 连接与鉴权
客户端通过 MySQL 协议发送 SQL。FE 结合当前用户、默认 Catalog、Database、会话变量、Resource Group 和超时设置完成后续分析与调度。查询队列中的准入等待也需要单独观察,不能全部计入解析耗时。[1][13]
这一步出现问题时,常见现象包括连接建立慢、连接池耗尽、认证失败、代理转发异常和请求排队。此时 BE 可能完全没有收到任务,扩容 BE 也不会改善情况。
3.2 解析与绑定
解析器识别 SELECT、JOIN、WHERE、GROUP BY 和 ORDER BY 的语法结构。绑定阶段将 fact_orders、dim_city 和各个列名映射到真实元数据,完成别名解析、类型检查、函数解析和权限校验。
o.city_id = c.city_id 两侧类型需要兼容,SUM(o.amount) 需要找到匹配的聚合函数签名,日期字符串也需要按规则转换为日期类型。
3.3 逻辑计划与规则优化
从关系语义上,查询大致可以理解为下面的树。下图是教学示意,实际计划另行获取:
Sort(total_amount DESC)
└─ Aggregate(region, SUM(amount))
└─ Join(o.city_id = c.city_id)
├─ Filter(order_date >= '2026-08-01')
│ └─ Scan fact_orders
└─ Scan dim_city
规则优化会下推适用的过滤,裁剪无关列。日期分区可以帮助排除无关分区,分桶裁剪也有对应条件。第九节未创建日期分区,且没有 city_id 的固定等值过滤,不能据此认定这两类裁剪已经发生。[8]
3.4 CBO 选择 Join 与数据分布
优化器需要估算过滤后的行数、构建侧大小及聚合基数,再决定 dim_city 是否适合广播到参与 Join 的执行位置。广播范围随实例部署确定,不等于集群中所有节点。
在采用 Broadcast Join 和分阶段聚合的假设下,可以观察下面这些依赖:
- 扫描维表,把 Build 侧数据广播到参与 Join 的执行位置;
- 接收端构建 Hash Table,并在条件允许时形成 Runtime Filter;
- 事实表扫描应用能够生效的静态谓词和运行时过滤;
- 订单行在对应位置 Probe,得到匹配的区域;
- 各执行位置完成局部聚合;
- 局部状态按
region 分发,合并同一区域的结果;
- 最终排序与结果输出完成全局有序的返回。
如果维表实际上很大,广播会让参与 Join 的执行位置接收重复的 Build 数据,增加网络与内存成本。统计信息有助于评估这种代价,但仍需要用实际计划和 Profile 检查选择是否合适。[4][9]
3.5 物理计划、Fragment 与调度
物理计划明确具体算子和数据分布,Fragment 则组织分布式执行边界。FE 根据扫描范围与调度策略部署 Fragment Instance;Instance 内的 Pipeline 还可以由多个 Driver 并行推进。这些层次不能直接与线程数一一对应。[7][8]
计划描述“做什么”,Instance 表示“这份工作在哪台 BE 上运行”。一个 Fragment 可以在多台 BE 上拥有多个 Instance,所有实例执行同一类算子,处理各自的数据分片。
四、MPP:把一条查询变成集群级并行任务
4.1 Shared-nothing 的核心含义
在 StarRocks 存算一体集群中,每台 BE 使用自己的 CPU、内存和本地存储。查询扫描可利用本地副本,节点间的数据交换则通过通信接口完成。Shared-nothing 描述资源组织方式,不代表节点彼此没有数据传输。[1]
这种结构便于水平扩展。增加 BE 后,集群获得更多 CPU、内存、磁盘和并行执行槽位;副本均衡完成后,新节点也会承担数据扫描。
扩容收益受数据分布和查询形态约束。单分区只有少量 Tablet、查询总是命中一个热点 Bucket、Join Key 严重倾斜时,新增节点无法自动形成线性加速。架构提供并行能力,表设计和数据分布决定能力能否被充分利用。
4.2 Fragment:分布式执行阶段
Fragment 可以理解为分布式物理计划中的一个片段,包含一组算子及其输出方式。Fragment 之间常通过 Exchange 连接;片段内仍可能有 Pipeline 依赖和 Local Exchange,不能把所有 Exchange 都当成跨 Fragment 的边界。[8][13]
下面用四个概念片段表示一种 Broadcast Join、按区域合并及全局排序的组织方式。编号仅用于示意,不对应实验应当出现的固定 Fragment ID:
Fragment 0:最终输出
Result Sink
└─ 全局有序汇聚(例如归并各路有序结果)
↑ Gather / Merging Exchange
Fragment 1:合并区域结果
发送局部有序结果
└─ Sort(total_amount DESC)
└─ Merge Aggregate(region)
↑ 按 region 接收局部聚合状态
Fragment 2:接收端 Join 与局部聚合
Exchange Sink(按 region 分发)
└─ Local Aggregate(region, SUM(amount))
└─ Hash Join
├─ Probe:Scan fact_orders
└─ Build:接收广播数据 → 构建 Hash Table
↑ Broadcast Exchange
Fragment 3:提供维表行
Broadcast Exchange Sink
└─ Scan dim_city
实际计划可以合并、增加或调整片段。维表行先送到 Join 执行位置,Build 与 Probe 在那里配合;局部聚合之后合并同组状态,多路有序结果最后形成全局顺序。各个步骤的职责与数据依赖要分别标清。[8][9]
4.3 Instance:Fragment 的运行实例
Fragment 是计划片段,Instance 是运行时的执行实体。假设事实表有 32 个 Tablet、集群有 4 台 BE,扫描片段可以部署到多台 BE,每个 Instance 分配一组 Scan Range。数字仅用于说明层次,实际映射由调度决定。
一个 Instance 运行在一台 BE 上,可以包含多条 Pipeline;同一 Pipeline 可以有多个 Driver。Instance 与 Tablet 存储副本属于不同层次,数据分片和处理并行度也应分别确认。[7][8]
查询并行度受到多种因素影响:
- 参与执行的 BE 数量;
- 命中的 Tablet 与可用副本分布;
- Scan Range 和可拆分的扫描工作量;
- Bucket 数与分区设计;
pipeline_dop 等并行设置;
- Resource Group 的资源配额与查询准入限制;
- Join 或聚合后的数据分布;
- 热点 Key 和数据倾斜;
- 节点 CPU、内存与 I/O 饱和程度。
并行度过低可能闲置资源,过高又会增加调度、缓冲和交换开销。pipeline_dop 控制 Pipeline 并行度,不等于总线程数。调优时应比较 Instance 和 Driver 的工作量、活动时间与等待时间。[7][13]
4.4 Exchange:连接各个执行阶段的数据通道
Exchange 负责把数据从发送端交给接收端。下表同时列出常见数据分发方式,以及减少分发的 Join 策略,注意它们所属层次不同:
| 方式或策略 |
数据如何流动 |
常见用途 |
主要风险 |
| Broadcast |
把 Build 侧数据送到参与 Join 的接收位置 |
大表 JOIN 小表 |
大小估算错误造成网络与内存放大 |
| Hash Shuffle |
按 Join Key 或 Group Key 重新分发 |
Join、分组聚合 |
传输量大、热点 Key、带宽瓶颈 |
| Gather |
多路结果汇聚到最终接收位置 |
最终输出、全局排序 |
汇聚端或客户端成为瓶颈 |
| Bucket Shuffle Join |
一侧按另一侧既有桶分布发送 |
减少 Join 中的数据移动 |
分布及 Join 条件必须满足要求 |
| Colocate Join |
利用匹配的数据放置,省去本次 Join 所需的跨节点重分布 |
分布稳定的关联查询 |
分桶条件与副本放置需匹配,其他阶段仍可能交换数据 |
Exchange 也解释了为什么网络是 MPP 系统的重要资源。查询在单机上很快,扩展到多表跨节点 Join 后变慢,常见原因正是 Shuffle 数据量和倾斜。
4.5 MPP 与 MapReduce 的差异
两者都能把任务分到多台机器。主要差异体现在执行模型和服务目标。
对包含 Reduce 阶段的 MapReduce 作业,Map 输出会经过缓冲、排序和本地落盘,再供 Reduce 拉取处理。这种阶段模型适合大规模离线作业,与数据库中的 Pipeline 有不同的组织方式。[19]
MPP 数据库是常驻服务。查询计划被切分成多个并行阶段,数据通过网络 Exchange 以流式方式在算子之间传递,优化器会根据表统计和分布选择 Join 与聚合策略,目标是较低交互延迟和较高并发。
StarRocks 支持部分算子的 Spill,开启并正确配置后,可以把部分中间状态写入磁盘。MPP 查询同样可能落盘,Spill 也不能解决所有内存不足。关键在于整个链路围绕常驻数据库服务、流水线算子和代价优化器组织。[12]
五、BE 内部:列存、向量化、SIMD 与 Pipeline
5.1 列式存储先解决“少读数据”
分析查询经常面对宽表。假设表有 100 列,某条 SQL 只使用日期、地区和金额三列,列式布局就可以按列读取。行式存储的实际 I/O 还取决于数据页、覆盖索引和是否回表,不能仅按列数比例计算性能收益。11
同一列具有明确的类型,也可以按取值分布选择编码,再配合通用压缩:
- 重复值可使用适用的游程编码;
- 低基数字符串可考虑字典编码;
- 固定宽度值可以利用位级组织改善存储;
- LZ4 常用于兼顾压缩和解压效率;
- ZSTD 可通过配置和数据分布取得不同的空间、CPU 权衡。
压缩可以减少存储体积和读取字节,解压会消耗 CPU;Exchange 的传输压缩又是另一条路径。需要按冷热访问、数据分布和算子负载权衡空间与处理成本,不能用磁盘压缩比直接推算网络或整体查询速度。11
列存仍然需要配合数据裁剪。读取三列却扫描全年数据,成本依旧可能很高。分区、分桶、ZoneMap、Bloom Filter、GIN 和 Runtime Filter 各有适用条件,只有实际建成并被查询利用的能力,才能减少相应工作量。6
5.2 向量化把逐行处理改成批量处理
传统火山模型通常由上层算子反复调用下层算子的 next(),每次处理一行或很小的数据单元。大量虚函数调用、分支判断和对象访问会消耗 CPU,也不利于缓存命中。
StarRocks 通过 Chunk 承载一批列数据,批量过滤、计算表达式与更新聚合状态。变长、可空和复杂类型各有布局,应结合实际列类型理解数据访问方式。7
这种方式带来几类收益:
- 函数调用成本被一批数据分摊;
- 顺序访问更容易利用 CPU Cache 与预取;
- 基础类型列可以紧凑存放,减少逐行对象管理;
- 编译器和手写实现更容易使用 SIMD;
- 算子之间传递 Chunk,减少不必要的逐行格式转换。
向量化适合扫描、过滤、聚合和批量表达式计算。单行主键点查中,规划、网络和定位成本的占比可能更高,因此 StarRocks 也提供 Prepared Statement 和满足条件的短路查询。行列混存还要求 Primary Key 表,且不适用于本章对照的存算分离形态。17
5.3 SIMD 让一条指令处理多个值
SIMD 是 Single Instruction Multiple Data 的缩写。CPU 可以用一条向量指令同时处理多个整数或浮点数。列式连续数组天然适合这种执行方式。
例如,过滤 amount > 100 时,可以批量比较并生成过滤掩码。部分实现能利用向量指令,具体路径取决于类型、算子和硬件;后面的 DECIMAL 示例也需要单独核对。
SIMD 收益取决于数据类型、算子实现、CPU 指令集和分支复杂度。字符串解析、复杂 UDF、频繁类型转换和数据依赖较强的逻辑,通常难以获得同等幅度的加速。硬件升级也无法替代合理建模和数据裁剪。
5.4 CPU Cache 亲和性
CPU 访问 L1、L2、L3 Cache 的延迟远低于访问主内存。连续读取列式数组更容易预取,缓存行利用率也更高。
大量独立分配、依靠指针连接的行对象可能增加缓存失效。StarRocks 的列式结构能减少这类访问:NullableColumn 分开保存数据列与 NULL 标记,BinaryColumn 使用字节区与偏移组织变长值。顺序处理通常更利于局部性,但 Hash Join 等随机访问仍可能成为瓶颈。11
这解释了列存与向量化经常同时出现:列式磁盘布局减少无关 I/O,列式内存布局改善批量处理效率,两者通过 Chunk 在扫描与计算之间衔接。7
5.5 Pipeline 解决多查询与多核调度
假设一个查询包含 Scan、Filter、Join、Aggregation 和 Exchange。某个 Scan 正在等磁盘,Exchange 正在等网络,Join 的 Build 侧尚未完成。若每个执行链长期占用独立线程,大量线程会处于等待状态,线程切换和栈内存也会不断增加。
Pipeline 把可串联的算子组织成流水线,由 Driver 在执行线程上推进。输入未就绪、输出无法接收、前置依赖未完成或运行时间片用尽时,Driver 进入相应等待或重新调度路径。扫描 I/O 可以由另外的任务承担,结果就绪后再供算子消费。7
Pipeline 有助于多个执行链共享 CPU。并发哈希表、聚合与排序仍会占用内存,磁盘和网络饱和也会造成等待。后续还要结合 Resource Group、准入控制、内存限制和 Spill 治理这些压力。12
5.6 三层技术栈的关系
可以用一句话记住:
MPP 决定工作怎样分到多台 BE;
Pipeline 决定一台 BE 内的任务怎样共享多核;
向量化决定一个算子怎样高效处理一批列式数据。
三者层次不同,组合后形成从集群、节点到 CPU 指令的完整并行体系。
六、把 JOIN 与聚合的数据流完整串起来
继续使用区域销售额查询。下面沿着前述 Broadcast Join 和分阶段聚合的假设串起七个步骤;小样例的真实计划可能采用更短的路径。
步骤 1:FE 裁剪扫描范围
order_date >= '2026-08-01' 可以作为扫描过滤条件。如果表另外定义了日期分区,才可讨论相应分区裁剪;分桶裁剪也需要匹配的条件。第九节的未分区表不能用于证明日期分区已被裁掉。
步骤 2:小表构建 Hash Table
在本例假设中,dim_city 是 Build 侧。维表扫描端发送数据,参与 Join 的接收位置构建 Hash Table。同一 Instance 内的 Probe Driver 可以共享相应构建结果,不能按 Driver 数机械计算完整维表的副本数量。9
步骤 3:生成 Runtime Filter
Build 侧可以形成 Runtime Filter,供允许下推的 Probe 侧算子或扫描路径使用;具体类型由 Join 条件、类型和版本决定。它可以提前排除部分不匹配数据,仍须由 Join 验证匹配关系。Bloom 类过滤可能有假阳性,不能替代 Join 本身。9
Runtime Filter 到达需要时间。扫描端可以等待一小段时间,也可以先开始工作。等待过久会增加延迟,到达过晚又可能错过大量裁剪机会,具体策略由执行器和参数共同决定。
步骤 4:多台 BE 并行扫描事实表
每个 Instance 处理分配到的扫描工作。SegmentIterator 根据谓词和已有索引选择裁剪路径,运行时过滤也可参与执行,处理顺序随路径变化。实验表未创建 Bloom 或 GIN,不能标注为已命中。6
步骤 5:本地 Join 与局部聚合
订单行 Probe 对应的 Hash Table 后得到区域,在采用局部聚合的计划中,各执行位置先计算分组状态。业务维表必须保持预期的键唯一性,否则同一笔订单会因多次匹配被重复汇总。
假设每个执行位置处理一亿行,只有几十个区域,并且这些行被有效汇总为少量状态,发送量就能大幅下降。实际输出还受聚合策略与中间刷新影响,应查看聚合输出和 Exchange 字节,而不只看原始行数。
步骤 6:按分组键 Shuffle
局部聚合状态按 region 分发,同一区域进入对应的合并位置。这里传递的是能够继续合并的状态或局部结果,不能把所有聚合都当作简单数字相加。例如平均值需要对应的和与计数,去重计数还要维护相应去重信息。
步骤 7:全局聚合与排序
后续 Fragment 合并各区域的局部和。若结果分散在多个实例中,还需要 Gather 或有序归并得到全局排序,再由结果端输出。FE 的 Coordinator 获取结果批次并返回客户端,不负责重新执行普通 BE 聚合算子。5
这条链路中,任何环节都可能成为瓶颈:
- 应有的分区或数据裁剪没有生效,Scan 读取过多;
- 小表估算错误,Broadcast 过大;
- Runtime Filter 到达过晚或过滤作用很小;
- Join Key 或扫描分布不均,少数实例拖尾;
- 局部聚合分组过多,状态和输出数据变大;
- 排序或结果返回产生额外等待;
- 客户端读取速度跟不上结果生成。
架构知识的价值就在这里:看到“SQL 运行 30 秒”,我们能够提出可验证的问题,而不会直接把原因归结为机器不够。
七、存储、内存、CPU 与网络怎样协同
查询引擎的整体效率取决于四类资源之间的配合。
7.1 存储负责提供有效数据
本地 SSD、HDD 或远端对象存储决定原始读取能力。分区、Tablet、Segment、Page 和索引裁剪决定真正需要读取多少数据。
最佳优化通常从“减少读取”开始。读取量下降,可能连带减少解压和过滤工作;网络传输量还取决于下游 Join、聚合和分发方式,内存峰值也取决于需要保留的状态,几项资源不会总是同比例下降。
7.2 内存承载执行状态
Hash Join 的 Build 表、Aggregation 的分组状态、Sort 缓冲区、Exchange 队列和缓存都会占用内存。并发查询数量增加时,单条查询内存看似合理,集群总内存仍可能触顶。
StarRocks 使用内存跟踪与限制机制,查询、实例、算子和进程指标各有口径。查询达到自身或资源组限制时,机器仍有空闲内存也可能失败。Spill 仅能缓解部分算子的压力。12
7.3 CPU 完成解码与算子计算
解压、表达式求值、Hash、比较、聚合和序列化都会消耗 CPU。向量化、SIMD 和 Cache 亲和性提高单位 CPU 的处理效率。
CPU 使用率高未必意味着异常。高吞吐分析查询理应充分使用 CPU。需要关注的是:处理了多少有效数据,是否存在大量无效扫描,单核是否被热点 Instance 占满,算子是否因数据倾斜拖尾。
7.4 网络连接分布式阶段
网络负责计划下发、Exchange、Runtime Filter、心跳与结果传输。大规模 Shuffle 对带宽和延迟很敏感,跨可用区或跨机房部署会放大影响。
网络瓶颈常见于大表 Join、全局高基数聚合、Broadcast 误判、热点 Key 和超大结果集。表分桶、Colocate、局部聚合和合理 Join 顺序都能减少网络流量。
7.5 从资源开销理解查询成本
观察查询代价时,可以分别记录下面几类资源工作与等待:
读取工作:扫描多少有效数据,缓存命中情况怎样
CPU 工作:解压、解码、表达式和算子计算
内存状态:哈希表、聚合、排序与缓冲的峰值
交换工作:序列化、传输以及接收端处理
等待时间:准入、调度、I/O 与上下游依赖
分区裁剪减少扫描,物化视图复用计算,Colocate 可省去特定 Join 的重分布,Resource Group 管理资源竞争。优化器用自己的代价模型比较候选计划,实际耗时还受并发和运行条件影响;内核估算值与用户等待时间需要分别解释。418
八、慢查询分层定位:先找阶段,再找资源
遇到慢查询时,建议按照固定顺序排查。
8.1 先确认问题是否真实发生变化
记录 SQL 文本、用户、数据库、版本、执行时间、并发、数据量、导入任务和集群变更。确认下列问题:
- SQL 是否完全相同;
- 参数和时间范围是否变化;
- 数据是否刚完成大批量导入;
- 统计信息是否更新;
- 集群是否扩缩容、升级或发生节点异常;
- 同期是否有 Compaction、Schema Change 或大导入;
- 冷缓存与热缓存是否被混在一起比较。
没有可比条件的“昨天 2 秒、今天 10 秒”很难直接归因。
8.2 查看 EXPLAIN
EXPLAIN 用来观察计划形态。重点关注:
- 分区命中数量;
- Tablet 扫描数量;
- Join 顺序与 Join 分布;
- Broadcast、Shuffle、Colocate 等选择;
- Runtime Filter;
- 物化视图改写;
- 聚合阶段与 Exchange 边界。
普通 EXPLAIN 展示计划与估算,无法给出本次查询的真实算子耗时;EXPLAIN ANALYZE 会执行查询并报告运行信息,不能将两者混用。计划看起来合理,运行时仍可能受数据倾斜、缓存、I/O 和并发影响。8
8.3 读取 Query Profile
Profile 需要按配置实际采集。可在测试会话启用 enable_profile,执行查询后通过 FE 页面、SHOW PROFILELIST 或 get_query_profile 获取记录。合并视图可能隐藏 Instance 和 Driver 差异,必要时用 pipeline_profile_level = 2 保留原始层次。常用观察项如下,具体路径以目标版本为准。13
| 层次 |
重点观察 |
可能问题 |
| FE |
Parser、Analyzer、Optimizer、Pending、Prepare、Deploy 等阶段 |
元数据、优化器、准入与调度等待 |
| Scan |
RawRowsRead、RowsRead、BytesRead、ScanTime、I/O 等待 |
裁剪效果差、文件过多、存储或扫描调度慢 |
| Join |
BuildHashTableTime、HashTableMemoryUsage、Build/Probe 工作量 |
构建侧过大、估算偏差、数据倾斜 |
| Agg / Sort |
分组状态、输入输出行数、内存与 Spill |
高基数、排序量大、内存压力 |
| Exchange |
BytesSent、BytesReceived、序列化与等待 |
传输量大、接收端慢、网络受限 |
| Instance / Driver |
最大与最小工作量、Active/Pending/Schedule 时间 |
节点差异、热点、调度或依赖等待 |
| Result |
返回行数、结果端缓冲与客户端读取 |
结果过大、客户端或代理限制 |
8.4 再决定优化动作
诊断结果不同,动作也应不同:
- 扫描过大:检查谓词、分区、排序键和适用索引;
- 计划不合理:检查统计信息、估算行数与 Join 分布;
- 交换量大:评估局部聚合、Colocate、Bucket Shuffle 或数据模型;
- 数据倾斜:核验热点和结果多重性,再评估拆分、预聚合或改写;
- 内存过高:控制并发与 Build 侧大小,检查适用的 Spill 和资源限制;
- I/O 受限:核对文件、Compaction、缓存与存储介质;
- 返回过慢:减少不必要的结果,按稳定排序设计分页,并核对驱动读取方式。
参数调优应放在明确瓶颈之后。盲目提高线程数、并行度和内存限制,可能让资源竞争更加严重。
九、动手实验:标注一条 JOIN 与聚合 SQL 的执行链路
下面保留单副本、小数据的存算一体练习,用于认识计划结构,不代表生产拓扑或性能。仅在获准的独立实验环境执行;sr_day03 必须是本次专用的新数据库,若已存在应先停止,另选空命名空间,不清理或覆盖其他课程数据。单副本不提供数据副本冗余。14
9.1 创建维表
CREATE DATABASE sr_day03;
USE sr_day03;
CREATE TABLE dim_city (
city_id INT,
city_name VARCHAR(64),
region VARCHAR(32)
)
DUPLICATE KEY(city_id)
DISTRIBUTED BY HASH(city_id) BUCKETS 4
PROPERTIES (
"replication_num" = "1"
);
9.2 创建事实表
CREATE TABLE fact_orders (
order_date DATE,
city_id INT,
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(18, 2)
)
DUPLICATE KEY(order_date, city_id, order_id)
DISTRIBUTED BY HASH(city_id) BUCKETS 4
PROPERTIES (
"replication_num" = "1"
);
这里没有 PARTITION BY,日期列只是排序键的一部分。两表同样按 city_id 分桶,也没有自动加入同一 Colocation Group,不能仅据此认定会执行 Colocate Join。10
9.3 插入少量样例数据
INSERT INTO dim_city VALUES
(1, '北京', '华北'),
(2, '上海', '华东'),
(3, '广州', '华南'),
(4, '成都', '西南');
INSERT INTO fact_orders VALUES
('2026-08-01', 1, 1001, 501, 120.50),
('2026-08-01', 1, 1002, 502, 230.00),
('2026-08-01', 2, 1003, 503, 80.00),
('2026-08-02', 3, 1004, 504, 310.00),
('2026-08-02', 4, 1005, 505, 160.00);
两张表均为 Duplicate Key,重复执行 INSERT 会保留重复行。只导入一次,并先检查维表 4 行、事实表 5 行及维表键唯一性。随后可按五笔金额手工推导各区域销售额,作为结果对照;这个推导不代表已经在 StarRocks 中运行成功。14
9.4 查看执行计划
EXPLAIN VERBOSE
SELECT
c.region,
SUM(o.amount) AS total_amount
FROM fact_orders o
JOIN dim_city c
ON o.city_id = c.city_id
WHERE o.order_date >= '2026-08-01'
GROUP BY c.region
ORDER BY total_amount DESC;
小数据、单机、统计信息和优化器选择都可能改变计划形态;没有出现示意图里的四个 Fragment,也不构成失败。先保留实际 EXPLAIN,再执行第三节 SELECT 核对结果。需要 Profile 时,在当前测试会话启用采集后重新执行查询,不把 EXPLAIN 当作实际扫描。8
- 哪些部分属于 Scan,实际扫描范围是什么;
- Join 使用了什么分布方式,选择依据是什么;
- 是否生成 Runtime Filter,运行时是否起作用;
- 有哪些跨片段或本地 Exchange;
- 聚合分成了哪些阶段,如何合并结果;
- 全局顺序与 Result Sink 在哪里形成;
- 哪些决定由 FE 作出,哪些算子在 BE 运行。
9.5 完成交付物
复制实际 EXPLAIN,再配控制与数据路径图,用四种标记完成注释。解析、鉴权和优化通常不会作为 BE 算子出现在计划中,应单独标在 FE 控制路径上:
- 蓝色:FE 解析、优化与调度;
- 绿色:BE Scan、Join、Agg、Sort;
- 橙色:Exchange 与数据分布;
- 紫色:结果合并与返回。
最终产出一张《一条 SQL 的端到端链路注释图》。当后续学习 Profile、Join 调优和存储引擎时,可以继续在这张图上补充指标。
十、版本与表达边界
本章沿用 FE、BE、Fragment、Instance、Exchange 等架构概念,关键实现以 StarRocks 4.1.4 的固定源码为证据。官网会持续更新,文档中的新功能或新字段不能自动归入这个补丁。阅读其他版本资料时,尤其注意以下事项。
第一,MPP 不代表所有中间数据始终留在内存。StarRocks 的部分算子可以在配置允许时 Spill,但表达式求值等开销仍可能受内存限制。落盘可以缓解特定压力,也会增加 I/O 与延迟。12
第二,Fragment、Instance、Pipeline、Driver 和 Operator 是不同层次,Profile 还可能合并其中的实例信息。抓住“解析—分析—优化—计划—调度—执行—交换—返回”的流程,再用对应版本的计划与 Profile 验证,避免把算子、数据批次和线程混为一谈。713
第三,存算分离中由 CN 承担计算与缓存,持久数据位于对象存储或 HDFS。它与存算一体共享许多查询概念,但扫描、缓存、扩缩容和可用性条件不同。部署与数据路径会在 Day 4 单独说明。1
第四,性能数字必须放在明确条件下理解。宽表聚合、点查、大表 Join 和日志搜索的资源路径不同,无法用一个 QPS 或一条 benchmark 代表全部业务。
十一、Knowledge Check
题目 1
FE 与 BE 各自承担哪些核心职责?
答案要点: FE 负责连接、鉴权、元数据、解析、绑定、优化、计划和调度;BE 负责数据存储、扫描、向量化算子、Pipeline 调度、Exchange 和结果生成。
题目 2
Fragment 与 Instance 有什么区别?
答案要点: Fragment 是分布式物理计划片段,Instance 是它在某台计算节点上的运行实体。一个 Fragment 可以有多个 Instance,Instance 内还有 Pipeline 和 Driver;它与 Tablet 的存储副本不同。8
题目 3
Exchange 为什么重要?
答案要点: Exchange 连接执行阶段,承载 Broadcast、Shuffle、Gather 等分发;Local Exchange 还可在 Instance 内交换数据。应区分本地传递与跨节点传输,结合实际字节和等待判断成本。8
题目 4
向量化为什么能降低函数调用成本?
答案要点: 算子通过 Chunk 批量接收和处理列数据,使调用成本由一批行分摊;合适的数据布局有利于 CPU Cache 和 SIMD,但具体收益受类型与算子实现影响。7
题目 5
Pipeline 解决的主要问题是什么?
答案要点: Pipeline Driver 推进可执行的算子链,遇到依赖、输入或输出阻塞时等待调度,条件满足后继续运行。它帮助多个执行链共享 CPU,扫描 I/O 等任务另有执行路径,不能据此推导查询线程总数。7
题目 6
MPP 的并行度受哪些因素影响?
答案要点: 参与节点、Tablet 和扫描工作量、分布方式、可用副本、Pipeline DOP、资源配额、查询准入、数据倾斜和硬件状态共同影响有效并行。DOP、Instance 数和线程数应分别观察。718
题目 7
EXPLAIN 与 Profile 各自回答什么问题?
答案要点: 普通 EXPLAIN 描述计划及估算,已采集的 Profile 记录实际运行信息。分析累计值、峰值、等待和实例差异时要注意指标口径;EXPLAIN ANALYZE 会执行查询,不能当作普通计划查看。8
十二、本日总结
今天建立了 StarRocks 查询架构的第一张完整地图。
FE 管理控制路径:接收连接,理解 SQL,读取元数据,完成语义检查,通过自身的优化器和计划构建生成分布式计划,再协调 Fragment Instance 的部署与执行。4
BE 承担数据路径:读取 Tablet 与 Segment,完成裁剪、解码和过滤,通过向量化算子执行 Join、聚合与排序,使用 Pipeline 调度多核任务,必要时通过 Exchange 在节点之间传输数据。
MPP、Pipeline 和向量化分别覆盖集群、节点和算子三个层次。列式存储减少无关 I/O,编码压缩减少数据体积,Chunk 组织批量计算,CPU Cache 与 SIMD 在适用路径提高效率;局部聚合和 Runtime Filter 则有助于减少交换与后续处理。
当查询变慢时,可以沿着同一张地图定位:先区分 FE 规划、BE 执行、网络交换和结果返回,再使用 EXPLAIN 与 Profile 查找扫描量、算子时间、内存、Shuffle 和数据倾斜。
明天进入 Day 4:存算一体、存算分离与版本选择。届时会讨论本地存储、远端存储成本、计算扩缩容、缓存冷启动,以及版本与已知问题如何进入生产治理。
官方资料
本章关键源码固定在 StarRocks 4.1.4 的 commit 4a9848edf03f5c936dac664b2d52527f48e72eb0。官网资料按本次核对快照留档,具体实验结果仍待独立运行。