Apache Pinot 实时 OLAP 深度实战:从 Segment 列存、Star-Tree 预聚合到 Upsert 与查询路由的工程全解
如果说 ClickHouse 解决的是"单机列存能跑多快",Trino 解决的是"跨源联邦能查多广",那么 Apache Pinot 从第一天起解决的问题是另一个维度:在每秒百万级事件持续写入的同时,让面向用户的分析查询仍然稳定返回在几十毫秒以内。
这不是一个"更快的 OLAP",而是一套围绕"实时可见性 + 高并发低延迟"重新设计的存储与查询体系。LinkedIn 内部用它支撑站内分析(谁看了我的主页),Uber 用它支撑数万 QPS 的看板查询。本文从存储格式、预聚合索引、实时摄取协议、Upsert 语义到查询路由五个层面拆开来看,并给出可直接落地的配置与调优建议。
一、设计原点:Scatter-Gather 与不可变 Segment
Pinot 的整个架构建立在一条铁律上:Segment 是不可变的(immutable)。
一条数据进入 Pinot 后,会被写入某个 Segment;Segment 写满(按行数或时间)即封存(seal),此后只读。这个约束看似笨拙,却是所有性能特性的源头:
- 索引可以无锁构建。Bloom Filter、倒排索引、范围索引、Star-Tree 全部在封存期一次性构建完毕,查询期只读,不需要并发控制,也不需要像 LSM 那样维护多层结构的读放大。
- 副本以文件为单位搬运。Segment 就是一个自包含的目录(若干列文件 + 元数据 + 索引),扩容、rebalance、深度存储上传下载都退化为文件操作。
- 实时与离线可以统一抽象。实时 Segment(realtime,正在写)与离线 Segment(offline,批处理产出)在查询层是同一种对象,只是来源不同。
查询执行采用经典的 Scatter-Gather + 两段式:
Broker 接收 SQL
→ 依据路由表选出 (server, segment) 集合
→ 并行下发到各 Server(每个 Segment 由单线程处理,天然并行)
→ 各 Server 本地过滤/聚合,返回中间结果
→ Broker 归并(merge)后返回客户端
关键点在于:Server 本地只做"能下推的部分",Broker 只做归并。因为 Segment 之间没有共享状态,聚合可以在每个 Segment 上独立完成——这让 Pinot 的聚合天然具备分布式可交换性。
二、存储层:列式布局与字典编码
一个 Segment 目录大致长这样:
myTable_OFFLINE_0/
├── creation.meta
├── metadata.properties # 段级元数据、时间范围、CRC
├── index.drain # 列级索引清单
├── column_A.dict # 字典(cardinality 低时启用)
├── column_A.fwd # 正向索引(字典ID 或原始值)
├── column_A.inv # 倒排索引
├── column_A.bloom # Bloom Filter
├── column_A.sorted / .range # 排序索引 / 范围索引
├── star_tree_index # Star-Tree(可选)
└── v3/ # 数据版本目录
字典编码是 Pinot 列存的默认选择。低基数列(城市、设备类型、状态码)会被编码为紧凑的字典 ID,正向索引存的是 ID 而不是原值。这带来两个直接收益:
- 存储压缩比显著下降(一个 4 字节 int 替代几十字节字符串);
- 谓词可以先在字典上求值。查询
WHERE city = 'Beijing'只需要在字典里二分查到 ID,然后正向索引上做 int 比较——字典 ID 数组天然适合 SIMD 批量比较,这就是 Pinot 谓词下推快的根本原因。
索引配置在 tableConfig 中逐列声明:
{
"tableName": "userEvents",
"tableType": "OFFLINE",
"segmentsConfig": {
"timeColumnName": "ts",
"replication": "3",
"segmentPushType": "APPEND",
"retentionTimeUnit": "DAYS",
"retentionTimeValue": "90"
},
"tableIndexConfig": {
"loadMode": "MMAP",
"invertedIndexColumns": ["city", "deviceType", "eventName"],
"bloomFilterColumns": ["userId"],
"rangeIndexColumns": ["ts"],
"sortedColumn": ["eventName"],
"noDictionaryColumns": ["rawPayload"]
},
"tenants": { "broker": "DefaultTenant", "server": "DefaultTenant" }
}
实战观点一:sortedColumn 是最被低估的优化项。 按某列物理排序后,该列的谓词退化为有序数组上的二分查找 + 区间扫描,且字典 ID 连续分布时压缩率极高。如果你的查询几乎总带某个过滤条件(例如租户 ID 或事件名),把它设为 sortedColumn 往往比加倒排索引收益更大——代价只是单个 Segment 内失去其他维度的聚簇。
实战观点二:高基数列要显式 noDictionaryColumns。 一个 UUID 列做字典编码,字典本身会比原始数据还大,且失去 raw 比较的 SIMD 友好性。经验阈值:cardinality / 行数 超过 1/1000 就该关闭字典。
三、Star-Tree:把预聚合做成索引
Star-Tree 是 Pinot 最具辨识度的设计。它解决的是"OLAP 立方体"的老问题:预计算所有维度组合会带来组合爆炸,而 Star-Tree 用按维度有序排列 + 层级前缀聚合把空间压到可控。
原理可以概括为三步:
- 维度按基数降序排列。高基维度在前,低基在后,形成一棵前缀树。
- 递归聚合。对每一层前缀做聚合,聚合结果作为父节点。因为维度有序,父节点覆盖的维度集合是子节点的前缀。
- **保留
*(星号)节点**,表示"该维度已被聚合掉"。
查询 SELECT city, SUM(clicks) FROM t GROUP BY city 时,优化器发现存在以 city 为前缀的聚合节点,直接读取该层预聚合结果,扫描行数从十亿级降到城市数量级。
配置方式:
"tableIndexConfig": {
"starTreeIndexConfigs": [{
"maxLeafRecords": 10000,
"dimensionsSplitOrder": ["country", "city", "deviceType"],
"skipStarNodeCreationForDimensions": ["deviceType"],
"functionColumnPairs": [
"SUM__clicks",
"COUNT__*",
"AVG__latencyMs",
"DISTINCTCOUNT__userId"
]
}]
}
-- 命中 Star-Tree:只扫描预聚合层,无需触碰明细行
SELECT city, SUM(clicks), COUNT(*)
FROM userEvents
WHERE country = 'CN'
GROUP BY city
ORDER BY 2 DESC
LIMIT 20;
实战观点三:Star-Tree 不是免费的。 它会显著拉长 Segment 构建时间(构建期需要多轮排序与聚合),并让 Segment 体积增加 20%~50%。判断标准很直接:只有当你的查询是"固定维度组合 + 固定聚合函数"的高频看板时才值得建。Ad-hoc 探索型查询命中率极低,白白付出写入成本。
另外一个常见坑:AVG 这类不可直接累加的度量,Pinot 会存成 SUM + COUNT 两个列,在查询期相除。所以 functionColumnPairs 里写 AVG__x 时,底层实际物化了两列——评估体积时别漏算。
四、实时摄取:Low Level Consumer 与 Segment Completion Protocol
Pinot 的实时表从 Kafka 消费,这里有一个精妙的工程取舍。
早期 Pinot 使用 High Level Consumer(HLC),依赖 Kafka 的消费者组做分区分配。问题是:Kafka 的 rebalance 会把分区从 A 实例挪到 B 实例,而 Pinot 正在写的 Segment 是内存状态,搬不走——只能丢弃重放,延迟抖动严重。
于是 Pinot 引入了 Low Level Consumer(LLC):由 Pinot Controller 自己做分区分配,每个 partition 的 consumption 状态记录在 ZooKeeper/Helix IdealState 里,Kafka 侧只保留一个最简单的 consumer,不加入消费者组。
这带来一个关键能力:Pinot 自己决定何时封存 Segment,从而可以实现精确一次语义的提交协议——Segment Completion Protocol。
流程(简化):
1. Server 写满 N 行 / 到达时间边界 → 停止消费该 partition
2. Server 构建索引、计算 CRC、把 Segment 上传至深度存储(S3/HDFS)
3. Server 向 Controller 发起 commit,携带 (segmentName, crc, rowCount)
4. Controller 在 ZK 上发起一次乐观锁写入:
- 第一个成功者成为 committer,返回 CONTINUE
- 其余副本收到 HOLD / DISCARD,从深度存储下载 committer 的 Segment
5. Controller 更新 IdealState,把该 partition 的 CONSUMING → ONLINE
6. Server 从 commit 成功的 offset 继续消费,开启下一个 Segment
为什么 CRC 比对是核心? 因为多个副本并行消费同一 partition,各自算出的结果可能不同(乱序、迟到数据、副本启动时间不同)。Controller 只允许一个副本提交,其余必须下载同一份文件——这保证了同一 partition 的所有副本 Segment 内容字节一致,查询才不会因为路由到不同副本而得到不同结果。
// 自定义部分更新合并器(Partial Upsert)
public class LatencyMerger implements PartialUpsertMerger {
@Override
public Object merge(Object incoming, Object current) {
if (current == null) return incoming;
// 保留最新的 p99,保留最早的首见时间
return Math.max(((Number) incoming).longValue(),
((Number) current).longValue());
}
}
五、Upsert:让事件流支持"按主键更新"
Pinot 本质上是 append-only 的不可变存储,但面向用户的分析经常需要"最新状态"——比如"每个订单当前状态"。Pinot 的 Upsert 通过主键路由 + 查询期归并实现:
"upsertConfig": {
"mode": "FULL",
"primaryKeyColumns": ["orderId"],
"comparisonColumn": "ts",
"defaultPartialUpsertStrategy": "OVERWRITE",
"metadataTTL": 1,
"metadataTTLTimeUnit": "DAYS"
}
约束很硬:同一主键的所有记录必须落在同一个 partition、同一台 Server 上。因此开启 Upsert 的表,Kafka 分区键必须等于主键(或主键的超集)。写入时 Pinot 维护一张 primaryKey → (segment, docId) 的路由表,查询期每个 Segment 内先按主键去重取最新版本,再执行过滤聚合。
实战观点四:Upsert 的代价在查询侧而非写入侧。 很多团队以为 Upsert 只影响写入吞吐,实际上每次查询都要做主键归并,且该过程无法下推到索引层(必须先解出 docId 才能比较版本)。实测在开启 Upsert 的表上,相同查询的 P99 通常比纯 append 表高出 30%~80%。因此:
- 能用 append 建模的(比如只关心事件计数)不要上 Upsert;
- 主键基数过高(亿级)时路由表本身会成为内存瓶颈,需要评估
metadataTTL与 Server 堆内存; - 只需要更新少量列时用
PARTIAL模式配自定义 merger,避免整行覆盖带来的版本冲突。
六、查询路由:三段裁剪
Broker 拿到 SQL 后,路由阶段决定了 90% 的性能上限。Pinot 的裁剪是分层的:
第一层:时间边界裁剪。 每个 Segment 在元数据里记录了时间列的最小/最大值,带时间范围的查询直接跳过无关 Segment。这也是为什么 timeColumnName 必须正确配置——配错了,这层裁剪完全失效。
第二层:Segment 级裁剪(Pruner)。 利用列级 min/max、Bloom Filter、倒排索引的存在性判断,在 Broker 侧读取 Segment 元数据(缓存在内存)过滤掉不可能命中的 Segment。
第三层:Segment 内裁剪。 落到 Server 后,先用倒排/Bloom/范围索引算出匹配的 docId 位图,再去做列扫描。
路由策略本身也可配:
"routing": {
"segmentPrunerTypes": ["time", "partition", "value"],
"instanceSelectorType": "ReplicaGroup",
"replicaGroupPartitionConfig": { "numInstancesPerReplicaGroup": 2 }
}
实战观点五:ReplicaGroup vs Balanced 的选型常被搞反。 Balanced 追求的是跨 Segment 的负载均衡,每个查询都能打满所有实例;ReplicaGroup 则把副本分组,让一个查询只命中一个组——牺牲集群总吞吐,换取单个查询的延迟稳定性和缓存局部性(同一组 Server 反复处理同一批 Segment,page cache 命中率高)。面向用户的低延迟看板应该选 ReplicaGroup;离线大批量扫描选 Balanced。
七、多阶段查询引擎 v2
Pinot 从 0.11 起引入 Multi-Stage Query Engine(v2),用 DAG 化的算子替换了单阶段的 Scatter-Gather,支持大表 Join、Window 函数、CTE。
-- v2 引擎:大表 Join + 窗口函数
SET useMultistageEngine = true;
WITH recent AS (
SELECT userId, SUM(amount) AS amt
FROM orders WHERE ts > now() - INTERVAL '7' DAY
GROUP BY userId
)
SELECT u.country, r.userId, r.amt,
RANK() OVER (PARTITION BY u.country ORDER BY r.amt DESC) AS rk
FROM recent r JOIN users u ON r.userId = u.userId
WHERE rk <= 10;
需要明确的是:v2 引擎目前不适合作为面向用户的低延迟查询路径。它引入了 shuffle、临时表、更重的调度开销,延迟通常在百毫秒到秒级。生产上的合理划分是:
- v1(单阶段):面向用户的高 QPS 看板、告警查询,毫秒级;
- v2(多阶段):分析师 Ad-hoc、ETL 补数、复杂 Join,秒级。
两者共存是常态,不要试图用 v2 替换 v1。
八、调优清单(可直接照做)
| 问题 | 检查项 | 建议 |
|---|---|---|
| 查询慢 | 是否命中 Star-Tree | 用 EXPLAIN PLAN 确认走的是聚合层 |
| 写入延迟抖动 | LLC 提交是否频繁 HOLD | 检查副本消费速度差,调大 segmentFlushThresholdRows |
| Segment 过多 | 时间粒度是否过细 | 用 APPEND 合并小 Segment,或调大 segmentPushFrequency |
| 内存吃紧 | 字典是否过大 | 高基数列加 noDictionaryColumns |
| Upsert 表 P99 高 | 是否真的需要 Upsert | 评估改为 append + 查询期取最新 |
| 实时数据不可见 | 时间边界配置 | 确认 timeColumnName 与事件时间一致,避免乱序导致边界偏移 |
结语
Pinot 的设计哲学可以浓缩为一句话:用不可变性换取确定性,用预计算换取低延迟。它放弃了通用 OLAP 的灵活性(不支持任意 UPDATE、Join 能力弱于 Trino/Spark),换来的是在"数据持续写入 + 用户持续查询"这个特定场景下几乎无法被替代的表现。
选型时的判断坐标很清晰:如果你的场景是面向终端用户的高并发实时看板(而不是分析师跑批),Pinot 是极少数能同时把写入延迟、查询延迟和数据新鲜度都压到秒级以内的选择。反之,如果是重度 Ad-hoc、复杂多表关联的分析负载,那么 Trino 或 Spark 仍然是更合适的工具。
理解 Star-Tree 的预聚合边界、Segment Completion Protocol 的提交语义、以及 Upsert 的查询期代价,是把 Pinot 用好和用崩的分水岭。

发表评论 取消回复