Apache Druid 物化视图(Materialized View):基于 derivativeDataSource 的查询加速实战指南
【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid6/druid
导读
Apache Druid 的物化视图(Materialized View)功能用于解决一类典型场景:当查询的 dataSource 包含大量维度,而实际查询只涉及其中少数几个维度时,每次查询仍须扫描整张表的全部列,造成不必要的开销。本文基于当前仓库中的官方文档 docs/development/extensions-contrib/materialized-view.md,系统讲解物化视图的两大组成部分——materialized-view-maintenance(物化视图维护)与materialized-view-selection(物化视图选择),并深入其源码实现:从derivativeDataSourcesupervisor 的提交与增量构建,到view查询类型的自动改写与多子查询合并。读完本文,你将掌握物化视图扩展的加载方式、supervisor spec 的完整配置、view查询的编写方法,以及命中率监控等运维要点。
功能概览与适用前提
物化视图的本质是"用空间换时间":为一个维度很多的基表(base-dataSource)预先构建若干个只包含查询所需维度子集与指标子集的派生表(derived-dataSource),查询时自动路由到更小的派生表上执行,从而显著降低扫描量与计算成本。
使用该功能需要满足两个前提(见 materialized-view.md):
- 加载两个扩展:
materialized-view-selection与materialized-view-maintenance; - 当前依赖 Hadoop 集群:派生数据源的构建走的是 Hadoop 批式索引任务(
HadoopIndexTask),因此本功能目前要求环境中存在可用的 Hadoop 集群。
从当前仓库的模块结构可以印证这一点:materialized-view-maintenance 负责生成并提交HadoopIndexTask,而 materialized-view-selection 负责在 Broker 端改写查询。
扩展加载方式遵循 Druid 的标准机制:将两个扩展目录放入 Druid 的extensions目录,并在common.runtime.properties中配置druid.extensions.loadList,使其包含druid-materialized-view-maintenance与druid-materialized-view-selection。
Materialized-view-maintenance:派生数据源的创建与维护
基本概念:base-dataSource 与 derived-dataSource
在materialized-view-maintenance中,用户摄入的数据源被称为base-dataSource。对于每一个 base-dataSource,可以提交一个derivativeDataSourcesupervisor,用于创建并持续维护其他数据源,这些数据源被称为derived-dataSource。派生数据源的维度(dimensions)和指标(metrics)必须是 base-dataSource 相应集合的子集。
每个derivativeDataSourcesupervisor 只负责一个派生数据源,其核心职责是保持派生数据源的时间线(timeline)与 base-dataSource 保持一致:当基表新增、更新或删除了某个时间区间的数据时,supervisor 会调度 Hadoop 索引任务重建派生数据源对应区间的 segment。
derivativeDataSource supervisor 配置示例
以下是一个完整的derivativeDataSourcesupervisor spec 示例(来自官方文档,基于 wikiticker 数据):
{ "type": "derivativeDataSource", "baseDataSource": "wikiticker", "dimensionsSpec": { "dimensions": [ "isUnpatrolled", "metroCode", "namespace", "page", "regionIsoCode", "regionName", "user" ] }, "metricsSpec": [ { "name": "count", "type": "count" }, { "name": "added", "type": "longSum", "fieldName": "added" } ], "tuningConfig": { "type": "hadoop" } }在 MaterializedViewMaintenanceDruidModule.java 中,MaterializedViewSupervisorSpec以derivativeDataSource类型名注册到 Jackson 序列化体系,同时DerivativeDataSourceMetadata也注册为同一类型名;该模块带有@LoadScope(roles = NodeRole.OVERLORD_JSON_NAME)注解,说明 supervisor 相关逻辑只在 Overlord 节点加载。
Supervisor 配置字段说明
| 字段 | 描述 | 是否必填 |
|---|---|---|
| type | supervisor 类型,固定为derivativeDataSource。 | 是 |
| baseDataSource | base-dataSource 名称。该数据源的数据必须已经存在于 Druid 中,并将作为输入数据。 | 是 |
| dimensionsSpec | 指定数据的维度,必须是 base-dataSource 维度的子集。 | 是 |
| metricsSpec | 聚合器列表,必须是 base-dataSource 指标的子集。详见 聚合(aggregations)。 | 是 |
| tuningConfig | 必须是 HadoopTuningConfig。详见 Hadoop tuning config。 | 是 |
| dataSource | 该派生数据源的名称。 | 否(默认 =baseDataSource+ supervisor 的 hashCode) |
| hadoopDependencyCoordinates | Hadoop 依赖坐标的 JSON 数组,Druid 将使用它覆盖默认 Hadoop 坐标;一旦指定,Druid 会从druid.extensions.hadoopDependenciesDir指定的位置查找这些 Hadoop 依赖。 | 否 |
| classpathPrefix | 为 Peon 进程前置追加的 classpath。 | 否 |
| context | 见下文。 | 否 |
关于默认 dataSource 名称的源码细节
官方文档描述默认名称为"baseDataSource-hashCode of supervisor",从源码看更精确的生成逻辑位于 MaterializedViewSupervisorSpec.java:
this.dataSourceName = dataSourceName == null ? StringUtils.format( "%s-%s", baseDataSource, DigestUtils.sha1Hex(dimensionsSpec.toString()).substring(0, 8) ) : dataSourceName;即未显式指定dataSource时,默认名称为{baseDataSource}-{dimensionsSpec.toString() 的 SHA-1 哈希前 8 位}。这意味着维度组合不同的 supervisor 会得到不同的派生数据源名称,而维度组合完全相同的 supervisor 会产生相同的默认名称,多个 supervisor 提交到同一派生数据源时需注意避免冲突。
此外,源码中supervisorId的格式为MaterializedViewSupervisor-{dataSourceName}(见 MaterializedViewSupervisor.java),提交的任务 ID 前缀为index_materialized_view。
Context 配置字段
| 字段 | 描述 | 是否必填 |
|---|---|---|
| maxTaskCount | supervisor 同时提交的最大任务数。 | 否(默认 = 1) |
源码中除maxTaskCount外还支持另一个 context 参数minDataLagMs(MaterializedViewSupervisor.java):
maxTaskCount默认值为1(DEFAULT_MAX_TASK_COUNT);minDataLagMs默认值为1 天(DEFAULT_MIN_DATA_LAG_MS = TimeUnit.DAYS.toMillis(1)),其作用是:派生数据源与基表之间存在滞后,为防止延迟数据导致反复重建,时间区间起点距最新区间不足minDataLagMs的区间不会被构建(对应hasEnoughLag方法)。
派生数据源的构建流程(源码级解析)
在 MaterializedViewSupervisor.java 中,supervisor 通过scheduleWithFixedDelay周期性执行run(),间隔由druid.materialized.view.task.taskCheckDuration控制(默认PT1M,见 MaterializedViewTaskConfig.java)。
每次执行的完整链路为:
- 校验元数据:从元数据存储读取派生数据源的
DataSourceMetadata,必须是DerivativeDataSourceMetadata且其 baseDataSource、dimensions、metrics 与 spec 完全一致(run())。首次启动时,supervisor 会调用insertDataSourceMetadata写入该元数据。 - 清理已结束任务:遍历
runningTasks,将状态非 runnable 的任务从运行表中移除(checkSegmentsAndSubmitTasks())。 - 计算缺失区间(
checkSegments):对比基表与派生表的 segment 时间线:- 基表有数据而派生表没有的区间 → 需要构建(
entriesOnlyOnLeft); - 派生表 segment 的 version 不是基表 segment 最大
created_date的区间 → 需要重建(entriesDiffering且 base version 更大); - 派生表有数据而基表已无数据的区间 → 将对应 segment 标记为 unused(
markSegmentAsUnused)。 - 构建顺序按区间开始时间倒序排列,即最新区间的数据最先构建。
- 基表有数据而派生表没有的区间 → 需要构建(
- 提交 Hadoop 任务(
submitTasks):在runningTasks.size() < maxTaskCount的前提下,为每个待构建区间调用 createTask() 生成HadoopIndexTask并加入任务队列。从createTask的实现可以看到任务的内部构造:- 解析格式为
timeAndDims的map类型 parser; - 使用
ArbitraryGranularitySpec+Granularities.NONE,即以单一 interval 为粒度构建 segment; - 输入为
dataSource类型的HadoopIOConfig,DatasourceIngestionSpec指向 base-dataSource 及对应区间的 segment 列表; - 派生 segment 的 version 取基表该区间所有 segment 中最大的
created_date。
- 解析格式为
supervisor 的状态与报告可通过 Overlord 的 supervisor API 获取:/druid/indexer/v1/supervisor/{supervisorId},报告(MaterializedViewSupervisorReport.java)中包含 base-dataSource、dimensions、metrics、缺失区间(condensed intervals)以及健康状态等字段。
Materialized-view-selection:view 查询类型与自动优化
基本概念
在materialized-view-selection中实现了一种新的查询类型view。当发起一个 view 查询时,Druid 会基于查询的 dataSource 与 intervals 尽力优化:将原始查询改写为一个或多个子查询,用更小的派生数据源替换基表。
该模块通过 MaterializedViewSelectionDruidModule.java 注册,注解@LoadScope(roles = NodeRole.BROKER_JSON_NAME)表明该逻辑只在Broker 节点加载。模块内部:
- 注册
MaterializedViewQuery的 Jackson 子类型(类型名为view); - 将
MaterializedViewQuery绑定到 MaterializedViewQueryQueryToolChest; - 以生命周期方式注册 DerivativeDataSourceManager,并绑定单例 DataSourceOptimizer;
- 注册
DataSourceOptimizerMonitor用于指标上报; - 配置项绑定前缀为
druid.manager.derivatives。
view 查询示例
以下是一个完整的 view 查询示例(来自官方文档,内层为一个 groupBy 查询,取added最大的 1 个user):
{ "queryType": "view", "query": { "queryType": "groupBy", "dataSource": "wikiticker", "granularity": "all", "dimensions": [ "user" ], "limitSpec": { "type": "default", "limit": 1, "columns": [ { "dimension": "added", "direction": "descending", "dimensionOrder": "numeric" } ] }, "aggregations": [ { "type": "longSum", "name": "added", "fieldName": "added" } ], "intervals": [ "2015-09-12/2015-09-13" ] } }view 查询的字段说明
view 查询由两部分组成:
| 字段 | 描述 | 是否必填 |
|---|---|---|
| queryType | 查询类型,固定为view。 | 是 |
| query | view 查询内包裹的真实查询。真实查询必须是 groupBy、topN 或 timeseries 类型之一。 | 是 |
这一点在源码中得到强制校验:MaterializedViewQuery.java 的构造函数通过Preconditions.checkArgument仅接受TopNQuery、TimeseriesQuery、GroupByQuery三种类型,其他查询类型会直接抛出 "Only topN/timeseries/groupby query are supported" 异常。因此,view 查询内层不支持 Scan、Search 等其他查询类型。
优化执行流程(源码级解析)
view 查询的执行分为两层:外层由MaterializedViewQueryQueryToolChest委托给真实查询的 ToolChest 处理(merge、排序、指标等均复用真实查询的实现);内层则由MaterializedViewQueryRunner调用DataSourceOptimizer.optimize()完成改写。
DataSourceOptimizer.optimize() 的核心逻辑如下:
- 类型与数据源检查:仅当查询是 TopN/Timeseries/GroupBy 且 dataSource 为
TableDataSource时才可能被优化,否则原样返回。 - 获取候选派生表:从
DerivativeDataSourceManager.getDerivatives(baseName)取该基表的所有派生数据源集合。该集合是有序的——DerivativeDataSource实现了Comparable,按"每个 segment 粒度上的平均数据大小(avgSizeBasedGranularity)"升序排列(DerivativeDataSource.java),平均数据量最小的派生表优先级最高,因为它的扫描成本最低。 - 提取必需字段:通过 MaterializedViewUtils.getRequiredFields() 从查询中提取所有依赖的列:包括 filter 所需列、聚合器所需字段(
requiredFields(),若是FilteredAggregatorFactory还会计入其过滤列)、topN 的维度列、groupBy 的维度列。 - 筛选可用派生表:仅当派生表
getColumns()(dimensions ∪ metrics,见 DerivativeDataSourceMetadata.java)包含全部必需字段时,该派生表才可作为候选。若没有候选,则记录missFields统计并原样返回查询。 - 按区间切分与改写:按优先级遍历候选派生表,对查询的每个 interval,通过 Broker 的
TimelineServerView查找派生表在该区间的实际覆盖(serverView.getTimeline(...).lookup(interval));将覆盖到的子区间改写为指向该派生表的新子查询(query.withDataSource(...).withQuerySegmentSpec(...)),并从剩余区间集合中减去(MaterializedViewUtils.minus,实现为两个区间列表的差集运算,支持不连续区间的拆分与合并)。 - 兜底基表:所有派生表都覆盖完后若仍有剩余区间,则追加一个仍指向基表、仅包含剩余区间的子查询,保证结果完整。
- 结果合并:
MaterializedViewQueryRunner使用MergeSequence按原始查询的结果排序规则将所有子查询结果合并返回(MaterializedViewQueryRunner.java)。
典型场景示例:基表base覆盖2011-04-01/2011-04-06,派生表derivative只覆盖2011-04-01/2011-04-04。当用户对base发起查询2011-04-01/2011-04-06时,优化器会将其拆分为两个子查询:一个在derivative上执行2011-04-01/2011-04-04,另一个在base上执行2011-04-04/2011-04-06,最终合并结果。这一行为在测试 DatasourceOptimizerTest.testOptimize 中有完整的断言验证。
派生表信息的同步机制
Broker 侧需要知道当前存在哪些派生表及其字段集合,这一信息由 DerivativeDataSourceManager 负责:
- 它以固定延迟周期性地扫描元数据库的
dataSource表(SELECT DISTINCT dataSource, commit_metadata_payload FROM ...),将commit_metadata_payload反序列化为DataSourceMetadata,仅保留DerivativeDataSourceMetadata类型的记录; - 对每个派生表调用
getAvgSizePerGranularity计算其"每个时间粒度区间的平均 segment 大小"(查询segments表中used = true的记录,用总大小除以不同 interval 数),并过滤掉大小为 0 的派生表; - 结果以
baseDataSource → SortedSet<DerivativeDataSource>的映射保存在静态引用DERIVATIVES_REF中,供优化器实时读取; - 轮询周期由配置项
druid.manager.derivatives.pollDuration控制,默认PT1M(见 MaterializedViewConfig.java)。
这意味着:新建/删除派生数据源后,Broker 最多需要等待一个轮询周期(默认 1 分钟)才能感知变化。
命中率与性能监控
Broker 的DataSourceOptimizerMonitor(DataSourceOptimizerMonitor.java)定期调用optimizer.getAndResetStats()并上报以下指标:
| 指标 | 含义 |
|---|---|
/materialized/view/query/totalNum | 可优化查询(针对有派生表的基表)的总次数 |
/materialized/view/query/hits | 成功命中派生表的查询次数 |
/materialized/view/query/hitRate | 命中率 = hits / totalNum |
/materialized/view/select/avgCostMS | 优化器平均耗时(毫秒) |
/materialized/view/derivative/numSelected | 每个派生表被选中的次数(维度为 derivative 名称) |
/materialized/view/missNum | 因缺少必需字段而无法命中时,缺失字段组合的出现次数(维度为 fields) |
这些指标可以帮助判断物化视图配置是否合理:若hitRate长期偏低,且missNum显示某些字段组合频繁缺失,说明现有派生表的维度/指标子集与实际查询负载不匹配,应调整派生表的字段设计。
使用流程小结
综合以上两部分,在一个已具备 Hadoop 集群的 Druid 环境中,启用物化视图的完整流程为:
- 加载扩展:在
common.runtime.properties的druid.extensions.loadList中加入druid-materialized-view-maintenance和druid-materialized-view-selection,重启相关服务(Overlord 加载维护模块,Broker 加载选择模块)。 - 确保基表就绪:先完成 base-dataSource 的摄入,确保其 segment 已发布。
- 提交 supervisor:通过 Overlord 的 Supervisor API(
POST /druid/indexer/v1/supervisor)提交derivativeDataSource类型的 spec,指定维度子集、指标子集与 Hadoop tuningConfig。supervisor 会周期性地对比基表与派生表时间线,自动补齐缺失区间并随基表更新而重建。 - 发起 view 查询:通过 Broker 的查询 API(
POST /druid/v2)提交queryType: view的查询,内层为 groupBy/topN/timeseries 之一。Broker 自动完成派生表选择、区间切分与结果合并,应用层无需改动。 - 监控与调优:通过
hitRate、avgCostMS、missNum等指标评估物化视图收益,按需调整派生表的字段集合与数量。
注意事项与限制
官方文档在结尾给出明确提醒:Materialized View 目前被标记为 experimental(实验特性)。使用时请务必保证所有进程(Overlord、Broker、Historical、Peon 等)的系统时间一致且单调递增,否则查询结果可能出现意外错误。这一要求与实现机制直接相关:派生表 segment 的 version 取自基表 segment 的created_date,时间不一致会导致版本比较失效,进而引发错误的区间重建或查询结果错乱。
其他限制总结如下:
- 依赖 Hadoop 集群,构建任务为
HadoopIndexTask; - 派生表字段必须是基表字段的子集,且至少覆盖查询实际用到的全部字段才能命中;
- 内层查询仅支持 groupBy、topN、timeseries 三种类型;
- 同一时间最多并行构建的任务数默认只有 1 个(可通过
context.maxTaskCount调大); - 受
minDataLagMs(默认 1 天)影响,最新区间存在滞后,派生表数据可能落后于基表; - Broker 感知派生表变更存在最长一个
pollDuration(默认 1 分钟)的延迟; - 派生表数据是基表数据的预聚合子集,若查询需要基表独有的列,则无法命中物化视图,仍会回退到基表执行。
相关仓库资源
- 官方文档:docs/development/extensions-contrib/materialized-view.md
- 维护模块主类:MaterializedViewSupervisorSpec.java、MaterializedViewSupervisor.java、DerivativeDataSourceMetadata.java
- 选择模块主类:MaterializedViewQuery.java、DataSourceOptimizer.java、DerivativeDataSourceManager.java、MaterializedViewUtils.java
- 相关测试:MaterializedViewSupervisorTest.java、DatasourceOptimizerTest.java
- 扩展模块定义:materialized-view-maintenance/pom.xml、materialized-view-selection/pom.xml
【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址: https://gitcode.com/gh_mirrors/druid6/druid
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考