Apache Druid 数据 Rollup(预聚合)完全指南:配置、原理与最佳实践
2026/9/23 23:54:12 网站建设 项目流程
  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

导读

Rollup(滚动聚合)是 Apache Druid 在**摄入阶段(ingestion time)**对原始数据进行的一种摘要化、预聚合处理,它可以把大量原始事件压缩成显著更小的存储形态,是 Druid 在高性能实时分析场景下控制存储成本与提升查询性能的核心手段。本文基于 Druid 官方文档 rollup 指南 展开,结合 ingestion-spec、schema-model 与仓库中的GranularitySpec源码实现,系统讲解 rollup 的配置方法、工作原理、适用场景、压缩率评估与 Perfect/Best-effort 两种模式的差异,并给出可直接落地的 schema 设计建议。读完本文,你将掌握:如何通过granularitySpecrollup开关控制预聚合行为、如何评估并最大化 rollup 比率、以及在不同摄入方式下如何选择正确的 rollup 模式。

什么是 Rollup:摄入期的摘要化预聚合

Druid 可以在摄入时对数据进行 rollup,从而减少磁盘上需要存储的原始数据量。Rollup 本质上是一种摘要化(summarization)或预聚合(pre-aggregation)操作。开启 rollup 后,数据规模可以被显著压缩,行数甚至可能降低若干个数量级(orders of magnitude)。

作为 rollup 高效性的代价,你会失去查询单个原始事件(individual events)的能力——原始记录一旦被聚合,就再也无法按事件级别还原。因此在决定是否开启 rollup 之前,需要先回答一个问题:你的查询是否依赖逐行(row-level)的原始数据?

Rollup 的核心工作机制

在摄入时,rollup 行为由dataSchemagranularitySpec中的rollup设置控制,且默认开启true)。开启后,Druid 会把满足以下两个条件的多行输入数据合并为一行

  1. 维度(dimension)值完全相同——参见 schema-model 中的 dimensions 说明;
  2. 主时间戳(primary timestamp)在经过queryGranularity截断后相同——参见 ingestion-spec 的 granularitySpec 章节与 schema-model 中的 primary timestamp 说明。

被合并到同一行的各行的指标(metric)值,会按照metricsSpec中定义的聚合函数(如sumcountlongSum等)逐列聚合,得到最终的指标值。这就是为什么 rollup 属于“预聚合”:聚合发生在数据写入磁盘之前,而不是查询时。

关闭 Rollup 的含义

当你把rollup设为false时,Druid原样加载每一行,不进行任何形式的预聚合。这种模式与不支持 rollup 特性的传统数据库行为类似。如果希望 Druid 将每条记录按原样存储(不经过任何摘要化处理),就应设置rollup: false

需要特别注意:即使queryGranularity设置为none(时间戳不做任何截断),rollup 依然生效——只要两条记录的时间戳完全相同、维度值完全相同,它们仍会被合并。这一点在 ingestion-spec 的 granularitySpec 字段说明中也有明确说明。

如何配置 Rollup:granularitySpec 全解

Rollup 开关位于摄入规范(ingestion spec)的dataSchemagranularitySpec。下面是一份完整的granularitySpec示例,出自 ingestion-spec:

"granularitySpec": { "segmentGranularity": "day", "queryGranularity": "none", "intervals": [ "2013-08-31/2013-09-01" ], "rollup": true }

granularitySpec负责四项工作:

  1. 通过segmentGranularity将数据源划分成时间块(time chunks);
  2. 通过queryGranularity对时间戳做截断(如需要);
  3. 通过intervals指定批量摄入要创建的 segment 时间块;
  4. 通过rollup指定是否启用摄入期 rollup。

rollup外,其余三项操作都基于主时间戳。各字段说明如下(摘自 ingestion-spec 的 granularitySpec 表格):

字段说明默认值
typegranularitySpec 的类型uniform
segmentGranularity数据源的时间块粒度。同一时间块内可以创建多个 segment。例如设为day时,同一天的事件落入同一时间块,并可基于其他配置与输入大小进一步分区。可使用任意粒度(granularity)。注意:同一时间块内的所有 segment 应具有相同的 segment 粒度。避免用WEEK做数据分区,因为周与月、年对不齐,难以按更粗粒度重新分区;建议使用DAYMONTHday
queryGranularitysegment 内时间戳的存储分辨率,必须等于或细于segmentGranularity。这是你能得到合理查询结果的最细粒度,但你可以按任何更粗的粒度查询。例如设为minute时,记录按分钟粒度存储,可以按分钟、5 分钟、小时等任意分钟的倍数合理查询。可使用任意粒度。设为none表示时间戳不做截断原样存储。注意:即使queryGranularitynone,rollup 依然生效——数据只要时间戳完全相同就会被 rollupnone
rollup是否使用摄入期 rollup。即使queryGranularitynone,rollup 依然有效——时间戳完全相同的行会被合并true
intervals定义 segment 时间块的 ISO8601 区间列表,如["2021-12-06T21:27:10+00:00/2021-12-07T00:00:00+00:00"]。省略时间部分时默认按00:00:00处理。Druid 会基于segmentGranularity对区间列表进行拆分与取整。若为null或未提供,批量摄入任务通常根据输入数据中的时间戳自行决定输出哪些时间块。若指定,批量摄入任务可能跳过“确定分区”阶段从而加快摄入,并可能一次性申请全部锁;区间外的记录会被丢弃。对任何流式摄入均忽略null

源码视角:rollup 如何被解析与传播

从源码层面看,rollup 开关最终落到GranularitySpec家族实现上。仓库中server/src/main/java/org/apache/druid/segment/indexing/granularity/目录下定义了:

  • GranularitySpec:接口,暴露getSegmentGranularity()getQueryGranularity()isRollup()等核心方法;
  • UniformGranularitySpec:默认实现,JSON 反序列化时直接接收segmentGranularityqueryGranularityrollupintervals四个字段(@JsonProperty);
  • ArbitraryGranularitySpec:支持不同区间使用不同粒度的实现;
  • BaseGranularitySpec:公共基类。

UniformGranularitySpec中,rollupnull时会被视为true(默认开启),queryGranularitysegmentGranularity为空时分别取默认值(DEFAULT_QUERY_GRANULARITYDEFAULT_SEGMENT_GRANULARITY,即noneday);intervals则被转换成IntervalsByGranularity,用于按 segment 粒度切分时间桶。这与文档中“rollup 默认开启”的描述完全一致。

在任务执行侧,批量摄入任务基类AbstractBatchIndexTask定义了抽象的isPerfectRollup()方法(见getPerfectRollup相关注释与第 290 行附近),用于标识任务是否处于完美(guaranteed)rollup 模式,并据此决定锁策略(如log.info("Using timeChunk lock for perfect rollup")所示,完美 rollup 模式使用 timeChunk 锁)。这意味着 rollup 模式不仅影响数据形态,还影响批量摄入任务的加锁与并发策略。

关闭 Rollup 时 metricsSpec 的建议

rollup关闭时,Druid 不做任何摄入期聚合,因此 ingestion-spec 建议metricsSpec置空——没有 rollup 就没有必要配置摄入期聚合器。不过也存在例外:如果你想通过定义 metric 来生成复杂列,以预计算部分近似聚合(如 approximate aggregations),此时即使关闭 rollup 也值得在metricsSpec中定义指标。

何时开启、何时关闭 Rollup

建议开启 Rollup 的场景

在创建表数据源(table datasource)时,如果同时满足以下两个条件,建议开启 rollup:

  • 你追求最优查询性能,或存在严格的存储空间约束
  • 不需要来自高基数维度的原始值(raw values)。

建议关闭 Rollup 的场景

反之,出现以下任一情况时建议关闭 rollup:

  • 你需要**逐行(individual rows)**的查询结果;
  • 你需要对任意列执行GROUP BYWHERE查询。

后者的原因在于:rollup 会丢失行级粒度,被聚合掉的原始行无法再被按行过滤或分组还原。

不同场景共存时的策略

如果不同用例对 rollup 的需求相互冲突(例如一部分查询需要明细、一部分查询只需要汇总),可以为不同的表分别创建不同的 rollup 配置。这是 Druid 支持多数据源的重要原因之一:明细表与汇总表各司其职。

如何评估 Rollup 效果:测量 Rollup 比率

要量化某个数据源的 rollup 收益,可以比较Druid 中的行数(COUNT(*)摄入的事件总数。后者可通过摄入时生成的count类型指标(num_rows)累加得到。

运行如下 Druid SQL 查询即可得到 rollup 比率:

SELECT SUM("num_rows") / (COUNT(*) * 1.0) FROM datasource

其中num_rows是摄入期生成的count类型指标。结果越大,说明 rollup 带来的收益越高。关于开启 rollup 时计数如何工作,可参考 schema-design 的 Counting 一节。

解读:若结果为 10,说明平均每个存储行由 10 条原始事件聚合而来,存储与扫描的行数都缩减为原来的 1/10。

如何最大化 Rollup 比率

Rollup 比率越高,存储与查询收益越大。以下策略来自 rollup 官方文档 的 "Maximizing rollup ratio" 章节:

  1. 精简 schema 设计:设计更少的维度、更低基数的维度,可以获得更好的 rollup 比率。维度越少、每个维度的取值种类越少,多行落入同一桶的概率越高。

  2. 用 Sketches 替代高基数维度:使用草图(sketches)(如 Datasketch 系列的近似聚合)避免存储高基数维度——高基数维度会显著拉低 rollup 比率。

  3. 调整queryGranularity:在摄入时调整queryGranularity,增加多行 Druid 记录时间戳匹配的概率。例如把分钟级PT1M改为五分钟级PT5M。时间戳截断得越粗,同一时间桶内的行越多,rollup 空间越大。

  4. 按需建立多个数据源:可以可选地把同一份数据加载到多个数据源:

    • 建一个“完整(full)”数据源:关闭 rollup,或开启 rollup 但保持极低的 rollup 比率;
    • 再建一个“精简(abbreviated)”数据源:维度更少、rollup 比率更高。

    当查询只涉及“精简”集合中的维度时,使用第二个数据源可以显著降低查询耗时。通常这种方法只需很小的额外存储开销,因为精简数据源往往比完整数据源小得多。

  5. 针对 best-effort rollup 的补救:如果使用无法保证完美 rollup 的摄入配置(详见下一节),可以尝试:

    • 切换到能保证完美 rollup 的方案;
    • 在初始摄入完成后,于后台对数据进行重新索引(reindex)或压缩(compaction)。

Perfect Rollup 与 Best-effort Rollup

根据摄入方式的不同,Druid 提供两种 rollup 模式:

  • 完美 rollup(Guaranteed perfect rollup):Druid 在摄入时完美地聚合输入数据,保证任意相同(时间戳、维度值)组合只在一个 segment 中出现一次。
  • 尽力而为 rollup(Best-effort rollup):Druid可能无法完美聚合输入数据,因此多个 segment 中可能残留具有相同时间戳与维度值的行。

为什么会出现 best-effort rollup

通常,提供 best-effort rollup 的摄入方式出于以下原因之一:

  • 并行化摄入但没有洗牌(shuffling)步骤:完美 rollup 需要把可聚合的行先分到同一分区,这依赖排序/洗牌预处理;省去该步骤的并行摄入无法保证跨分区的完美聚合。
  • 使用增量发布(incremental publishing):即在收到某个时间块的全部数据之前就先定稿并发布 segment,导致理论上可 rollup 的记录被拆散到不同 segment。

所有类型的流式摄入都运行在 best-effort 模式下

完美 rollup 的摄入方式则会在摄入前增加一个预处理步骤,先扫描整个输入数据集来确定区间与分区方案。虽然这增加了摄入耗时,但它提供了实现完美 rollup 所必需的分区信息。从源码看,这一“预处理确定分区”的阶段正是批量任务可以跳过determining-partitions阶段来加速摄入的前提——在ingestion-specintervals字段说明中也有对应描述。

各摄入方式的 rollup 行为对照

以下表格(摘自 rollup 文档)总结了各摄入方式对 rollup 的处理:

摄入方式处理方式
Native batchindex_parallelindex类型基于配置可为 perfect 或 best-effort
SQL-based batch(MSQ)始终 perfect
Hadoop始终 perfect
Kafka indexing service始终 best-effort
Kinesis indexing service始终 best-effort

实践提示

  • 如果你使用Kafka / Kinesis 流式摄入并希望获得完美的 rollup 效果,不要只依赖摄入时的 best-effort rollup,应定期对历史 segment 执行 compaction,让后台任务把分散在多个 segment 中的可聚合行重新合并。
  • 如果你使用native batch且对聚合质量有严格要求,请在任务配置中启用完美 rollup 模式(对应源码中isPerfectRollup()返回 true 的分支,见 AbstractBatchIndexTask)。

Rollup 实战演示:一个可复现的完整示例

为了让上面的概念落地,这里复现 Rollup 教程 中的完整演示。使用SQL-based ingestion(MSQ 任务引擎),通过 web console 的Query视图执行。示例数据为 9 条网络流量事件,字段包含timestampsrcIPdstIPpacketsbytes

第一步:加载示例数据

INSERT INTO "rollup_tutorial" WITH "inline_data" AS ( SELECT * FROM TABLE(EXTERN('{ "type":"inline", "data":"{\"timestamp\":\"2018-01-01T01:01:35Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":20,\"bytes\":9024}\n{\"timestamp\":\"2018-01-01T01:02:14Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":38,\"bytes\":6289}\n{\"timestamp\":\"2018-01-01T01:01:59Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":11,\"bytes\":5780}\n{\"timestamp\":\"2018-01-01T01:01:51Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":255,\"bytes\":21133}\n{\"timestamp\":\"2018-01-01T01:02:29Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":377,\"bytes\":359971}\n{\"timestamp\":\"2018-01-01T01:03:29Z\",\"srcIP\":\"1.1.1.1\",\"dstIP\":\"2.2.2.2\",\"packets\":49,\"bytes\":10204}\n{\"timestamp\":\"2018-01-02T21:33:14Z\",\"srcIP\":\"7.7.7.7\",\"dstIP\":\"8.8.8.8\",\"packets\":38,\"bytes\":6289}\n{\"timestamp\":\"2018-01-02T21:33:45Z\",\"srcIP\":\"7.7.7.7\",\"dstIP\":\"8.8.8.8\",\"packets\":123,\"bytes\":93999}\n{\"timestamp\":\"2018-01-02T21:35:45Z\",\"srcIP\":\"7.7.7.7\",\"dstIP\":\"8.8.8.8\",\"packets\":12,\"bytes\":2818}"}', '{"type":"json"}')) EXTEND ("timestamp" VARCHAR, "srcIP" VARCHAR, "dstIP" VARCHAR, "packets" BIGINT, "bytes" BIGINT) ) SELECT FLOOR(TIME_PARSE("timestamp") TO MINUTE) AS __time, "srcIP", "dstIP", SUM("bytes") AS "bytes", SUM("packets") AS "packets", COUNT(*) AS "count" FROM "inline_data" GROUP BY 1, 2, 3 PARTITIONED BY DAY

这条语句的关键点:

  • FLOOR(TIME_PARSE("timestamp") TO MINUTE)把时间戳向下取整到分钟(对应queryGranularity: minute的效果);
  • timestampsrcIPdstIP分组(这三列成为维度);
  • bytespackets作为指标,按SUM聚合;
  • 额外生成count指标,记录每次 rollup 合并了多少条原始行。

第二步:查询结果

SELECT * FROM "rollup_tutorial"

返回结果:

|__time|srcIP|dstIP|bytes|count|packets| | -- | -- | -- | -- | -- | -- | |2018-01-01T01:01:00.000Z|1.1.1.1|2.2.2.2|35,937|3|286| |2018-01-01T01:02:00.000Z|1.1.1.1|2.2.2.2|366,260|2|415| |2018-01-01T01:03:00.000Z|1.1.1.1|2.2.2.2|10,204|1|49| |2018-01-02T21:33:00.000Z|7.7.7.7|8.8.8.8|100,288|2|161| |2018-01-02T21:35:00.000Z|7.7.7.7|8.8.8.8|2,818|1|12|

原始 9 行数据被压缩为5 行

第三步:逐分钟拆解 rollup 过程

2018-01-01T01:01分钟内的 3 条原始事件为例:

{"timestamp":"2018-01-01T01:01:35Z","srcIP":"1.1.1.1", "dstIP":"2.2.2.2","packets":20,"bytes":9024} {"timestamp":"2018-01-01T01:01:51Z","srcIP":"1.1.1.1", "dstIP":"2.2.2.2","packets":255,"bytes":21133} {"timestamp":"2018-01-01T01:01:59Z","srcIP":"1.1.1.1", "dstIP":"2.2.2.2","packets":11,"bytes":5780}

时间戳先被FLOOR到分钟(01:01:00),三条记录的维度值{srcIP, dstIP}完全相同,因此被合并为一行:packets = 20+255+11 = 286bytes = 9024+21133+5780 = 35937count = 3,即上文结果表的首行。

再比如01:02分钟内的 2 条事件合并后bytes = 6289+359971 = 366260packets = 38+377 = 415count = 2;而01:03分钟只有 1 条事件,无可合并对象,count保持为1。可以看到:rollup 的实际行为正是“截断时间戳 + 维度分组 + 指标聚合”三者的组合,与 tutorial-rollup 中的演示完全一致。

与其它文档的衔接

  • 关于 rollup 与计数(Counting)的关系,以及 sketches 在降维中的作用,见 schema-design。
  • 关于granularitySpec全部字段的权威说明,见 ingestion-spec。
  • 关于 SQL-based ingestion 中 rollup 的概念,见 MSQ 概念文档。
  • 关于压缩(compaction)如何修复 best-effort rollup 留下的“重复行”,见 compaction 文档。
  • 关于PT5M等粒度的写法与含义,见 granularities 文档。
  • 数据库
  • OLAP
  • 大数据
  • 后端

【免费下载链接】druid

Apache Druid: a high performance real-time analytics database.

项目地址:https://gitcode.com/gh_mirrors/druid6/druid
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询