- 数据库
- OLAP
- 大数据
- 后端
【免费下载链接】druid
Apache Druid: a high performance real-time analytics database.
导读
本文围绕 Apache Druid 的druid-distinctcount扩展,系统讲解如何加载该扩展、在 Timeseries / TopN / GroupBy 三类查询中使用distinctCount聚合器计算去重计数(典型场景如 UV 统计),并深入剖析其基于 Bitmap 的实现原理、两个必须满足的前置分区条件以及使用限制。读完本文,你将能够正确配置并安全使用该聚合器,避免因分区与粒度设置不当导致的计数错误。
该扩展位于仓库的 extensions-contrib/distinctcount 模块,对应官方文档为 docs/development/extensions-contrib/distinctcount.md。
一、加载 druid-distinctcount 扩展
druid-distinctcount是 Apache Druid 的 contrib 扩展,与其它扩展一样,需要通过扩展加载列表启用。官方文档指出,使用前需在扩展加载列表中包含druid-distinctcount,具体加载方式参见 加载扩展说明。
在实际部署中,通常在common.runtime.properties中通过druid.extensions.loadList属性声明:
druid.extensions.loadList=["druid-distinctcount"]也可以与其它扩展一起列出,例如:
druid.extensions.loadList=["postgresql-metadata-storage", "druid-hdfs-storage", "druid-distinctcount"]从仓库的 extensions-contrib/distinctcount/pom.xml 可以看到,该模块的artifactId为druid-distinctcount,依赖druid-processing、guava、jackson、fastutil等组件(均为 provided 或 test 作用域),并通过 DistinctCountDruidModule.java 注册了"distinctCount"这一聚合器 JSON 类型名:
public static final String DISTINCT_COUNT = "distinctCount"; ... new SimpleModule("DistinctCountModule").registerSubtypes( new NamedType(DistinctCountAggregatorFactory.class, DISTINCT_COUNT) )这意味着在查询 JSON 的aggregations数组中,只要将type设为"distinctCount",Jackson 反序列化即可定位到 DistinctCountAggregatorFactory.java。
二、使用前的两个关键前置条件(务必遵守)
官方文档明确强调,使用distinctCount聚合器前必须完成以下两步,否则结果可能错误:
按单一维度做 Hash 分区:使用基于单维度的 hash 分区规格,按该维度(例如
visitor_id)对数据进行分区,确保该维度上具有同一值的所有行都落入同一个 segment。因为该聚合器只在单个 segment 内做去重,若相同键分散到多个 segment,会导致重复计数(over count)。保证 queryGranularity 能被 segmentGranularity 整除:即查询粒度必须能整除分段粒度(例如查询粒度为
day而分段粒度为month时,day能整除month;反之若查询粒度比分段粒度更粗且无法整除,结果会出错)。官方原文的表述是:make sure queryGranularity is divided exactly by segmentGranularity or else the result will be wrong。
在 Druid 的摄取配置中,Hash 分区可通过partitionsSpec与partitionDimensions实现,例如使用type: "hashed"的分区规格并将partitionDimensions设为["visitor_id"]。Hashed 分区的基本思路是:先选定 segment 数量,再按分区维度哈希把行均匀分布到各 segment,相关背景可参考 Hadoop 摄取文档中的 Hashed 分区说明 与 本地批量任务文档。
三、三种查询类型中的 distinctCount 用法
以下三个示例完整继承自官方文档,分别演示 Timeseries、TopN、GroupBy 查询中的distinctCount聚合器配置。三者均对visitor_id字段做去重计数,输出命名为uv。
1. Timeseries 查询
{ "queryType": "timeseries", "dataSource": "sample_datasource", "granularity": "day", "aggregations": [ { "type": "distinctCount", "name": "uv", "fieldName": "visitor_id" } ], "intervals": [ "2016-03-01T00:00:00.000/2013-03-20T00:00:00.000" ] }2. TopN 查询
{ "queryType": "topN", "dataSource": "sample_datasource", "dimension": "sample_dim", "threshold": 5, "metric": "uv", "granularity": "all", "aggregations": [ { "type": "distinctCount", "name": "uv", "fieldName": "visitor_id" } ], "intervals": [ "2016-03-06T00:00:00/2016-03-06T23:59:59" ] }3. GroupBy 查询
{ "queryType": "groupBy", "dataSource": "sample_datasource", "dimensions": ["sample_dim"], "granularity": "all", "aggregations": [ { "type": "distinctCount", "name": "uv", "fieldName": "visitor_id" } ], "intervals": [ "2016-03-06T00:00:00/2016-03-06T23:59:59" ] }三个示例的公共参数含义如下:
type:聚合器类型,固定为"distinctCount",对应源码中注册的 JSON 类型名;name:输出结果中的字段名(如uv);fieldName:要去重计数的维度字段(如visitor_id),该字段必须是维度字段,因为聚合器内部通过DimensionSelector读取它(见下文源码分析)。
从 DistinctCountAggregatorFactory.java 的构造器可以看出,它接受三个 JSON 属性:name、fieldName和可选的bitmapFactory(name与fieldName均通过Preconditions.checkNotNull强制非空):
@JsonCreator public DistinctCountAggregatorFactory( @JsonProperty("name") String name, @JsonProperty("fieldName") String fieldName, @JsonProperty("bitmapFactory") BitMapFactory bitMapFactory ) { Preconditions.checkNotNull(name); Preconditions.checkNotNull(fieldName); ... this.bitMapFactory = bitMapFactory == null ? DEFAULT_BITMAP_FACTORY : bitMapFactory; }其中bitmapFactory未指定时默认使用RoaringBitMapFactory。如果需要对底层位图实现做精细调优,可以在聚合器配置中显式指定,例如:
{ "type": "distinctCount", "name": "uv", "fieldName": "visitor_id", "bitmapFactory": { "type": "roaring" } }bitmapFactory支持的类型由 BitMapFactory.java 中的 Jackson 子类型注册定义,共有三种:
| type 值 | 实现类 | 底层位图 |
|---|---|---|
roaring(默认) | RoaringBitMapFactory | RoaringBitmapFactory(见 RoaringBitMapFactory.java) |
java | JavaBitMapFactory | BitSetBitmapFactory(见 JavaBitMapFactory.java) |
concise | ConciseBitMapFactory | ConciseBitmapFactory(见 ConciseBitMapFactory.java) |
四、源码实现原理:Bitmap 驱动的单段去重
理解distinctCount的实现,才能明白为何它有上述前置条件与限制。整个扩展只有 10 个 Java 源文件,逻辑非常聚焦。
1. 聚合核心:DistinctCountAggregator
DistinctCountAggregator.java 是核心实现:它在factorize时通过DimensionSelector拿到维度值,每次aggregate()把当前行对应的维度值在 segment 字典中的整数索引写入MutableBitmap位图:
public void aggregate() { IndexedInts row = selector.getRow(); for (int i = 0, rowSize = row.size(); i < rowSize; i++) { int index = row.get(i); mutableBitmap.add(index); } }去重结果就是位图中置位比特的数量:
public Object get() { return mutableBitmap.size(); }注意这里的关键点:位图去重的是字典索引(整数),而不是原始字符串本身。由于同一 segment 内同一个维度值必然映射到同一个字典索引,位图的“集合去重”语义天然成立——这也正解释了为什么必须用单一维度 hash 分区把相同键的所有行收敛到同一个 segment:一旦相同键落入不同 segment,每个 segment 各自维护一个独立位图,最终合并阶段无法再对原始值去重。
2. 合并逻辑为何是求和
从 DistinctCountAggregatorFactory.java 可以看到,该聚合器的combine()与getCombiningFactory()均按数值求和处理:
public Object combine(Object lhs, Object rhs) { ... return ((Number) lhs).longValue() + ((Number) rhs).longValue(); } @Override public AggregatorFactory getCombiningFactory() { return new LongSumAggregatorFactory(name, name); }也就是说:单个 segment 内是精确的位图去重,而跨 segment 的结果合并是“各 segment 去重计数之和”,无法消除跨 segment 重复。这也是官方文档提醒“否则可能重复计数”的根源。getIntermediateType()与getResultType()均为LONG,getMaxIntermediateSize()为 8 字节,位图只在内存计算期间存活,最终以 long 形式输出。
3. Buffer 聚合与空值兜底
在需要缓冲式聚合(如 groupBy 的批处理)时,使用 DistinctCountBufferAggregator.java:它通过Int2ObjectMap<MutableBitmap>按 buffer 位置维护一组WrappedRoaringBitmap,aggregate()后将位图大小写入ByteBuffer中的 long 位置,init()时置 0,同样只返回 long 值。
当字段不是维度或无法构造DimensionSelector时,工厂会回退到 NoopDistinctCountAggregator.java(及对应的NoopDistinctCountBufferAggregator),直接返回 0,保证查询不会因缺失维度而报错。
4. 缓存键与模块注册
该聚合器的查询结果缓存键由 getCacheKey() 生成,包含类型标识、fieldName与bitmapFactory的字符串表示;其中类型字节0x10定义于 processing 模块的 AggregatorUtil.java 的DISTINCT_COUNT_CACHE_KEY,与其它聚合器共享同一套缓存键协议。
五、使用限制与注意事项
官方文档明确列出了distinctCount在使用中的两类限制:
与 groupBy 一起使用:每个 segment 内 groupBy 键的数量不应超过
maxIntermediateRows(这是 groupBy 查询的中间结果行数上限,可通过查询 context 调整)。一旦超过,结果将不正确。这是因为缓冲聚合需要为每个 groupBy 键维护独立的位图集合,超出上限后会导致部分键的位图被截断或丢弃。与 topN 一起使用:
numValuesPerPass(topN 单遍扫描处理的候选值数量)不应设置过大。该值过大时,distinctCount会为大量候选维度值各维护一份位图,占用大量内存,可能导致 JVM 内存溢出(OOM)。从DistinctCountBufferAggregator的Int2ObjectMap<MutableBitmap>结构可以直观看到,内存占用与“候选键数量 × 位图大小”成正比。
此外,结合前文源码结论,还需重申两点使用纪律:
- 必须按去重维度做单一维度 Hash 分区,否则合并阶段的求和会重复计数;
- 必须保证 queryGranularity 能被 segmentGranularity 整除,否则跨粒度聚合的数值会失真。
六、测试验证与延伸阅读
仓库为扩展提供了三类查询的完整测试,可直接作为行为参考:
- DistinctCountTimeseriesQueryTest.java:构造三条包含不同
visitor_id的记录,断言 Timeseries 查询得到UV = 3; - DistinctCountGroupByQueryTest.java:验证 groupBy 场景下的去重计数;
- DistinctCountTopNQueryTest.java:验证 topN 场景下的去重计数。
如果你希望了解 Druid 提供的其它去重相关聚合方案(例如基于基数估计的近似去重),可以继续阅读 聚合器文档;而distinctCount本身的特点是单 segment 内精确去重,代价是需要严格的分区与粒度配合,并受限于 groupBy / topN 的中间行数与内存参数。在实践中,请务必先完成第一节的分区与粒度设计,再在生产环境大规模使用。
- 数据库
- OLAP
- 大数据
- 后端
【免费下载链接】druid
Apache Druid: a high performance real-time analytics database.
相关推荐
Apache Druid DistinctCount 聚合器扩展详解:精确去重计数(UV)的原理、安装与三种查询实战
Apache Druid DistinctCount 聚合器扩展详解:精确去重计数(UV)的原理、安装与三种查询实战 本篇技术指南围绕 Apache Druid
数据库数据分析OLAP大数据实时分析数据仓库后端DrawerKit动画原理深度剖析:如何实现流畅的抽屉过渡效果
DrawerKit动画原理深度剖析:如何实现流畅的抽屉过渡效果 DrawerKit是一个优秀的iOS自定义视图控制器转场库,它能够让你的应用实现类似Apple
Apache Druid RabbitMQ Firehose 扩展实战指南:从配置、原理到源码解析
Apache Druid RabbitMQ Firehose 扩展实战指南:从配置、原理到源码解析 导读 本文围绕 Apache Druid(当前仓库版本 0.
数据库数据分析OLAP大数据实时分析数据仓库后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考