Apache Druid 物化视图(Materialized View):基于 derivativeDataSource 的查询加速实战指南
2026/9/23 1:18:53 网站建设 项目流程

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):

  1. 加载两个扩展materialized-view-selectionmaterialized-view-maintenance
  2. 当前依赖 Hadoop 集群:派生数据源的构建走的是 Hadoop 批式索引任务(HadoopIndexTask),因此本功能目前要求环境中存在可用的 Hadoop 集群。

从当前仓库的模块结构可以印证这一点:materialized-view-maintenance 负责生成并提交HadoopIndexTask,而 materialized-view-selection 负责在 Broker 端改写查询。

扩展加载方式遵循 Druid 的标准机制:将两个扩展目录放入 Druid 的extensions目录,并在common.runtime.properties中配置druid.extensions.loadList,使其包含druid-materialized-view-maintenancedruid-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 中,MaterializedViewSupervisorSpecderivativeDataSource类型名注册到 Jackson 序列化体系,同时DerivativeDataSourceMetadata也注册为同一类型名;该模块带有@LoadScope(roles = NodeRole.OVERLORD_JSON_NAME)注解,说明 supervisor 相关逻辑只在 Overlord 节点加载。

Supervisor 配置字段说明

字段描述是否必填
typesupervisor 类型,固定为derivativeDataSource
baseDataSourcebase-dataSource 名称。该数据源的数据必须已经存在于 Druid 中,并将作为输入数据。
dimensionsSpec指定数据的维度,必须是 base-dataSource 维度的子集。
metricsSpec聚合器列表,必须是 base-dataSource 指标的子集。详见 聚合(aggregations)。
tuningConfig必须是 HadoopTuningConfig。详见 Hadoop tuning config。
dataSource该派生数据源的名称。否(默认 =baseDataSource+ supervisor 的 hashCode)
hadoopDependencyCoordinatesHadoop 依赖坐标的 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 配置字段
字段描述是否必填
maxTaskCountsupervisor 同时提交的最大任务数。否(默认 = 1)

源码中除maxTaskCount外还支持另一个 context 参数minDataLagMs(MaterializedViewSupervisor.java):

  • maxTaskCount默认值为1DEFAULT_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)。

每次执行的完整链路为:

  1. 校验元数据:从元数据存储读取派生数据源的DataSourceMetadata,必须是DerivativeDataSourceMetadata且其 baseDataSource、dimensions、metrics 与 spec 完全一致(run())。首次启动时,supervisor 会调用insertDataSourceMetadata写入该元数据。
  2. 清理已结束任务:遍历runningTasks,将状态非 runnable 的任务从运行表中移除(checkSegmentsAndSubmitTasks())。
  3. 计算缺失区间checkSegments):对比基表与派生表的 segment 时间线:
    • 基表有数据而派生表没有的区间 → 需要构建(entriesOnlyOnLeft);
    • 派生表 segment 的 version 不是基表 segment 最大created_date的区间 → 需要重建(entriesDiffering且 base version 更大);
    • 派生表有数据而基表已无数据的区间 → 将对应 segment 标记为 unused(markSegmentAsUnused)。
    • 构建顺序按区间开始时间倒序排列,即最新区间的数据最先构建
  4. 提交 Hadoop 任务submitTasks):在runningTasks.size() < maxTaskCount的前提下,为每个待构建区间调用 createTask() 生成HadoopIndexTask并加入任务队列。从createTask的实现可以看到任务的内部构造:
    • 解析格式为timeAndDimsmap类型 parser;
    • 使用ArbitraryGranularitySpec+Granularities.NONE,即以单一 interval 为粒度构建 segment;
    • 输入为dataSource类型的HadoopIOConfigDatasourceIngestionSpec指向 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
queryview 查询内包裹的真实查询。真实查询必须是 groupBy、topN 或 timeseries 类型之一。

这一点在源码中得到强制校验:MaterializedViewQuery.java 的构造函数通过Preconditions.checkArgument仅接受TopNQueryTimeseriesQueryGroupByQuery三种类型,其他查询类型会直接抛出 "Only topN/timeseries/groupby query are supported" 异常。因此,view 查询内层不支持 Scan、Search 等其他查询类型

优化执行流程(源码级解析)

view 查询的执行分为两层:外层由MaterializedViewQueryQueryToolChest委托给真实查询的 ToolChest 处理(merge、排序、指标等均复用真实查询的实现);内层则由MaterializedViewQueryRunner调用DataSourceOptimizer.optimize()完成改写。

DataSourceOptimizer.optimize() 的核心逻辑如下:

  1. 类型与数据源检查:仅当查询是 TopN/Timeseries/GroupBy 且 dataSource 为TableDataSource时才可能被优化,否则原样返回。
  2. 获取候选派生表:从DerivativeDataSourceManager.getDerivatives(baseName)取该基表的所有派生数据源集合。该集合是有序的——DerivativeDataSource实现了Comparable,按"每个 segment 粒度上的平均数据大小(avgSizeBasedGranularity)"升序排列(DerivativeDataSource.java),平均数据量最小的派生表优先级最高,因为它的扫描成本最低。
  3. 提取必需字段:通过 MaterializedViewUtils.getRequiredFields() 从查询中提取所有依赖的列:包括 filter 所需列、聚合器所需字段(requiredFields(),若是FilteredAggregatorFactory还会计入其过滤列)、topN 的维度列、groupBy 的维度列。
  4. 筛选可用派生表:仅当派生表getColumns()(dimensions ∪ metrics,见 DerivativeDataSourceMetadata.java)包含全部必需字段时,该派生表才可作为候选。若没有候选,则记录missFields统计并原样返回查询。
  5. 按区间切分与改写:按优先级遍历候选派生表,对查询的每个 interval,通过 Broker 的TimelineServerView查找派生表在该区间的实际覆盖(serverView.getTimeline(...).lookup(interval));将覆盖到的子区间改写为指向该派生表的新子查询(query.withDataSource(...).withQuerySegmentSpec(...)),并从剩余区间集合中减去(MaterializedViewUtils.minus,实现为两个区间列表的差集运算,支持不连续区间的拆分与合并)。
  6. 兜底基表:所有派生表都覆盖完后若仍有剩余区间,则追加一个仍指向基表、仅包含剩余区间的子查询,保证结果完整。
  7. 结果合并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 环境中,启用物化视图的完整流程为:

  1. 加载扩展:在common.runtime.propertiesdruid.extensions.loadList中加入druid-materialized-view-maintenancedruid-materialized-view-selection,重启相关服务(Overlord 加载维护模块,Broker 加载选择模块)。
  2. 确保基表就绪:先完成 base-dataSource 的摄入,确保其 segment 已发布。
  3. 提交 supervisor:通过 Overlord 的 Supervisor API(POST /druid/indexer/v1/supervisor)提交derivativeDataSource类型的 spec,指定维度子集、指标子集与 Hadoop tuningConfig。supervisor 会周期性地对比基表与派生表时间线,自动补齐缺失区间并随基表更新而重建。
  4. 发起 view 查询:通过 Broker 的查询 API(POST /druid/v2)提交queryType: view的查询,内层为 groupBy/topN/timeseries 之一。Broker 自动完成派生表选择、区间切分与结果合并,应用层无需改动。
  5. 监控与调优:通过hitRateavgCostMSmissNum等指标评估物化视图收益,按需调整派生表的字段集合与数量。

注意事项与限制

官方文档在结尾给出明确提醒: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),仅供参考

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

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

立即咨询