第 03 关 · ★★

架构总览:FE、BE 与 MPP

FE/BE 职责边界、MPP 并行、列存、向量化与 Pipeline 如何共同撑起「快」。

已点亮 · 最佳 100 分

Day 3|StarRocks 架构总览:FE、BE、MPP、列存、向量化与 Pipeline

前两天已经建立了数据系统的坐标系,也明确了 StarRocks 的定位、场景和边界。今天开始进入 StarRocks 内部,沿着一条 SQL 的执行链路,认识 FE、BE、MPP、Fragment、Instance、Exchange、列式存储、向量化和 Pipeline。

学完这一篇,你应当能够把“StarRocks 为什么快”拆成若干可验证的环节,也能够在查询变慢时先判断问题位于规划、扫描、计算、网络交换,还是结果返回。文中的计划图属于教学示意,实验脚本尚待在隔离环境运行,不作为实测输出。

本日学习目标

完成今天的学习后,你应当能够:

  1. 说明 FE 与 BE 的职责边界,以及 Leader、Follower、Observer 的作用;
  2. 描述 SQL 从客户端进入 FE,到多个 BE 并行执行,再返回结果的完整链路;
  3. 解释 Fragment、Instance、Exchange 之间的关系;
  4. 区分 MPP、向量化执行和 Pipeline,各自解决哪一层问题;
  5. 理解列存、编码压缩、CPU Cache、SIMD 与多核调度如何协同;
  6. 使用 EXPLAIN 和 Query Profile 对慢查询做第一轮分层定位;
  7. 完成一条 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 和分阶段聚合的假设下,可以观察下面这些依赖:

  1. 扫描维表,把 Build 侧数据广播到参与 Join 的执行位置;
  2. 接收端构建 Hash Table,并在条件允许时形成 Runtime Filter;
  3. 事实表扫描应用能够生效的静态谓词和运行时过滤;
  4. 订单行在对应位置 Probe,得到匹配的区域;
  5. 各执行位置完成局部聚合;
  6. 局部状态按 region 分发,合并同一区域的结果;
  7. 最终排序与结果输出完成全局有序的返回。

如果维表实际上很大,广播会让参与 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

这种方式带来几类收益:

  1. 函数调用成本被一批数据分摊;
  2. 顺序访问更容易利用 CPU Cache 与预取;
  3. 基础类型列可以紧凑存放,减少逐行对象管理;
  4. 编译器和手写实现更容易使用 SIMD;
  5. 算子之间传递 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

  1. 哪些部分属于 Scan,实际扫描范围是什么;
  2. Join 使用了什么分布方式,选择依据是什么;
  3. 是否生成 Runtime Filter,运行时是否起作用;
  4. 有哪些跨片段或本地 Exchange;
  5. 聚合分成了哪些阶段,如何合并结果;
  6. 全局顺序与 Result Sink 在哪里形成;
  7. 哪些决定由 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。官网资料按本次核对快照留档,具体实验结果仍待独立运行。