- 大数据
- 数据库
- 后端
【免费下载链接】presto
The official home of the Presto distributed SQL query engine for big data
导读
KHyperLogLog(KHLL)是 Presto 内置的一种用于**估算大规模数据中两列关联关系(reidentifiability / joinability)**的紧凑型数据草图(data sketch)。本文以仓库内 KHyperLogLog 格式说明文档 为骨架,完整讲解 KHLL 的内存结构与二进制序列化布局(含各字段的字节序、含义与内存估算),并结合 KHyperLogLog.java 的实现与 khyperloglog.rst 的函数说明,带你掌握khyperloglog_agg、cardinality、intersection_cardinality、jaccard_index、uniqueness_distribution、reidentification_potential等 SQL 函数的使用方法,以及它们在分布式聚合与序列化存储场景下的底层原理。读完本文,你将能够在 Presto 中构建、合并、存储并分析两列关联关系的 KHLL 草图。
KHyperLogLog 是什么
KHyperLogLog 源自论文"KHyperLogLog: Estimating Reidentifiability and Joinability of Large Data at Scale"(Chia et al., 2019)。从源码注释与实现看,它在 Presto 中是一种两级数据结构:
- 外层是一个k 大小的 MinHash 结构,其每个条目(entry)以一个
long类型的哈希值作为键(key); - 每个键映射到一个HyperLogLog(HLL)草图,用于统计与该键相关联的另一列取值(uii,即 unique identifier input)的去重基数。
简单来说,假设有两列x与y,KHLL 用 MinHash 摘要x的取值(键),用每个键对应的 HLL 表示"与某个x值关联的所有y值"。因此一个 KHLL 草图可以回答诸如"有多少个x只关联了少量y(高度唯一,存在再识别风险)"之类的问题。
该类型在 Presto 中的类型名为KHyperLogLog,定义见 KHyperLogLogType.java,它是一个可变长度类型(继承AbstractVariableWidthType),底层以Slice承载序列化后的字节。其内部实现依赖 airlift 的com.facebook.airlift.stats.cardinality.HyperLogLog(KHyperLogLog.java)。
序列化格式:完整二进制布局
原文档 khll.md 明确指出:除非另行说明,所有字段均为小端序(little-endian)。
整体布局自上而下依次为:
| 字段 | 类型 | 说明 |
|---|---|---|
| format(版本字节) | byte | 序列化格式版本,当前实现中固定为1(VERSION_BYTE,见 KHyperLogLog.java) |
| K(maxSize) | int | MinHash 结构中最多可容纳的条目数,默认4096(DEFAULT_MAX_SIZE) |
| HLL buckets | int | 每个 HLL 草图的桶(bucket)数,默认256(DEFAULT_HLL_BUCKETS) |
| # of entries(minhashSize) | int | 当前 MinHash 结构中实际条目数,最多为 K |
| total HLL size | int | 所有序列化 HLL 草图的字节数总和 |
| HLL sizes | int[] | 每个序列化 HLL 草图的大小(字节数),顺序与 keys 一致 |
| keys | long[] | MinHash 结构中按升序排列的键(哈希值)序列 |
| HLLs | 字节序列 | 与 keys 顺序一一对应的序列化 HLL 草图字节 |
该布局与序列化实现完全吻合。查看 KHyperLogLog.java 的serialize()方法可见其写入顺序:先写 1 个版本字节(appendByte(VERSION_BYTE)),再依次写maxSize、hllBuckets、minhash.size()、totalHllSize四个int,随后写出每个 HLL 的 sizes 数组与 keys 数组,最后逐个追加序列化后的 HLL 字节。反序列化newInstance(Slice)(KHyperLogLog.java)则按同样顺序读取,并校验版本字节("Unexpected version")。
两个值得注意的细节:
- 原文档提到"currently just one format exists,
0",而当前仓库实现中VERSION_BYTE = 1。这说明该字段的语义是"格式/版本标识",实际取值以当前仓库源码为准,读取时会做严格校验,版本不匹配即抛出异常。 - HLL 草图自身的序列化遵循 airlift 的 HyperLogLog 文档格式(
HyperLogLog.newInstance(serializedHll)/hll.serialize()),本文不再展开,具体可见 KHyperLogLog.java 中对HyperLogLog.newInstance的调用。
内存与序列化体积估算
KHLL 内部维护了两个计数器hllsTotalEstimatedInMemorySize与hllsTotalEstimatedSerializedSize,在每次add、mergeWith以及溢出淘汰时通过increaseTotalHllSize/decreaseTotalHllSize保持同步(KHyperLogLog.java)。
estimatedInMemorySize()约为:对象本身 + 红黑树结构 +minhash.size() * 8(键的long) + 各 HLL 内存估算之和;estimatedSerializedSize()约为:1 + 4 * 4(版本字节 + 4 个int)+minhash.size() * (8 + 4)(每个键的long+ 每个 HLL 大小的int)+ 各 HLL 序列化字节数之和(KHyperLogLog.java)。
这两个估算值同时被聚合状态的内存记账所使用,见下文"聚合状态与内存管理"。
核心算法与源码级实现
数据插入:update 与溢出淘汰
add(long value, long uii)与add(Slice value, long uii)先将value通过Murmur3Hash128哈希为 64 位键,再调用私有update(long hash, long uii)(KHyperLogLog.java):
if (!(minhash.containsKey(hash) || isExact() || hash < minhash.lastLongKey())) { return; }即:只有当该哈希已存在、MinHash 尚未装满(isExact(),即条目数< maxSize)、或新哈希小于当前最大键时,才真正插入。随后computeIfAbsent为该键创建/复用 HLL 并hll.add(uii)。最后removeOverflowEntries()会循环淘汰最大的键(minhash.lastLongKey()),保证条目数不超过 K。
这个"只保留最小的 K 个哈希"的策略正是 MinHash 的精髓:它使得两个数据集的 MinHash 集合可以近似其集合交集,而每个键内用 HLL 压缩了与该键关联的y值集合。
基数估计:cardinality()
- 当
isExact()为真(未满 K 个条目)时,直接返回minhash.size(),即精确基数; - 否则按"哈希密度外推"估算:
long hashesRange = minhash.lastLongKey() - Long.MIN_VALUE,结合Long.divideUnsigned计算密度,再外推到哈希输出范围的一半(Long.MAX_VALUE),并引用 Beyer 等人的论文进行偏差修正(KHyperLogLog.java)。
合并:merge 与 mergeWith
merge(KHyperLogLog a, KHyperLogLog b)有一个关键设计(KHyperLogLog.java):总是保留 K 值较小(分辨率更高)的一方作为合并基座,因为若把小 K 的草图并入大 K 的草图,前者的 MinHash 空间无法覆盖后者的全部 MinHash 空间,会损失分辨率。mergeWith逐个键合并:键相同时把两个 HLL 合并,键不同则直接插入,最后同样执行removeOverflowEntries()(KHyperLogLog.java)。
交集与 Jaccard 指数
exactIntersectionCardinality(a, b):仅当两个草图都处于 exact 状态时可用,直接取Sets.intersection(a.minhash.keySet(), b.minhash.keySet()).size();jaccardIndex(a, b):取两集合键的并集,在较小的集合大小范围内统计共同键的比例(KHyperLogLog.java)。
在 SQL 函数层,intersection_cardinality会优先走精确路径;否则用jaccard * union.cardinality()估算,并修正为不超过较小集合的基数(KHyperLogLogFunctions.java)。
再识别潜力与唯一性分布
reidentificationPotential(long threshold):统计基数(该键关联的y值去重数)不超过阈值的键所占比例,即"有多少x值只关联了少量y值";uniquenessDistribution(long histogramSize):默认直方图大小 256(DEFAULT_HISTOGRAM_SIZE),对每个 HLL 的基数取min(cardinality, histogramSize)落入对应桶,桶内值为相对频率1 / minhash.size()(KHyperLogLog.java)。
Presto 中的 SQL 函数与使用方式
依据官方函数文档 khyperloglog.rst,KHLL 可通过khyperloglog_agg创建,并可 cast 为varbinary以便存储复用。具体函数如下:
| 函数 | 返回类型 | 说明 |
|---|---|---|
khyperloglog_agg(x, y) | KHyperLogLog | 返回表示x与y两列关联关系的草图;MinHash 摘要x,HLL 表示与各x关联的y |
cardinality(khll) | bigint | MinHash 草图基数,即x的基数估计 |
intersection_cardinality(khll1, khll2) | bigint | 两个草图 MinHash 结构所代表数据的集合交集基数 |
jaccard_index(khll1, khll2) | double | 两个草图数据的 Jaccard 指数 |
uniqueness_distribution(khll) | map<bigint,double> | 唯一性分布直方图,默认桶数为当前 MinHash 条目数 |
uniqueness_distribution(khll, histogramSize) | map<bigint,double> | 指定桶数的唯一性直方图,超过histogramSize的唯一性全部累计到最后一个桶 |
reidentification_potential(khll, threshold) | double | 唯一性低于threshold的x值占比(再识别潜力) |
merge(khll) | KHyperLogLog | 多个草图聚合后的并集(聚合函数) |
merge_khll(array[khll]) | KHyperLogLog | 数组形式 KHLL 的并集 |
SQL 使用示例:
-- 构建草图:x 为 bigint 列,y 为 bigint 列 SELECT khyperloglog_agg(x, y) AS khll FROM source_table; -- 估算 x 的基数 SELECT cardinality(khyperloglog_agg(x, y)) FROM source_table; -- 估算两组数据的 Jaccard 指数与交集基数 SELECT jaccard_index(khll_a, khll_b), intersection_cardinality(khll_a, khll_b) FROM (SELECT khyperloglog_agg(x, y) AS khll_a FROM table_a) a CROSS JOIN (SELECT khyperloglog_agg(x, y) AS khll_b FROM table_b) b; -- 评估再识别风险:唯一性不超过 5 的 x 值占比 SELECT reidentification_potential(khyperloglog_agg(x, y), 5) FROM source_table; -- 草图与 varbinary 互转,便于落盘存储 SELECT CAST(khyperloglog_agg(x, y) AS varbinary) AS stored FROM source_table;输入类型支持
khyperloglog_agg的第一参数(x)与第二参数(uii,即y)支持多种组合:bigint/varchar/double均可作为x;bigint/varchar可作为uii。当uii为varchar时,会先经XxHash64哈希为long再写入 HLL(见 KHyperLogLogAggregationFunction.java 与 KHyperLogLogWithLimitAggregationFunction.java)。
序列化存储与 cast
KHyperLogLogOperators.java 提供了KHyperLogLog <-> varbinary的双向 cast(直接透传底层Slice)。因此你可以把草图 cast 成varbinary存入外部表,下次读取后再 cast 回来继续做合并与分析。
聚合状态与内存管理
在 Presto 聚合框架中,KHLL 的中间状态由KHyperLogLogState接口描述,其序列化器与工厂分别为 KHyperLogLogStateSerializer.java(序列化类型即KHyperLogLog,空状态写 NULL)与 KHyperLogLogStateFactory.java(提供单值与分组两种状态)。
- 分组状态
GroupedKHyperLogLogState使用ObjectBigArray<KHyperLogLog>按 group 存放草图,并实时维护getEstimatedSize()(对象大小 + 各草图内存估算 + 数组开销),供查询引擎做内存控制; - 工厂支持
groupLimit参数:当 group 数超过限制时抛出NOT_SUPPORTED异常(错误信息中提示由khyperloglog-agg-group-limit配置控制),用于防止分组过多导致内存爆炸(KHyperLogLogStateFactory.java)。
KHyperLogLogWithLimitAggregationFunction是khyperloglog_agg的一个带分组上限的变体实现,其getDescription()明确描述了语义:"MinHash structure summarizes x and the HyperLogLog sketches represent y values linked to x values"。
合并函数的实现
merge聚合函数(@AggregationFunction("merge"))直接以KHyperLogLog作为输入,逐个mergeWith后输出序列化结果(MergeKHyperLogLogAggregationFunction.java);merge_khll(array(khyperloglog))则遍历数组(跳过 NULL 元素),对首个非空元素依次合并,空数组返回 NULL(KHyperLogLogFunctions.java)。
使用建议与注意事项
- 版本字节校验:序列化首字节当前固定为
1(与 khll.md 中"当前仅一种格式"的描述略有出入,以源码为准),跨版本读取不兼容的字节流会直接抛 "Unexpected version"; - K 与桶数的默认值:
K = 4096、hllBuckets = 256。K 决定 MinHash 的精度与内存上限,桶数决定单个 HLL 的精度,两者都可通过构造函数指定(KHyperLogLog.java); - 合并方向:
merge始终以 K 较小者为基础,避免分辨率损失;自建合并流程时也应遵循这一约定; - 精确与近似:未装满(条目数
< K)时,cardinality与intersection_cardinality走精确路径;装满后为近似估计,误差特性与 MinHash/HLL 的参数直接相关; - 存储:若需持久化草图,请通过
CAST(... AS varbinary)落库,读取后转回KHyperLogLog再参与merge聚合。
相关源码与文档索引
- 格式说明:presto-main-base/src/main/java/com/facebook/presto/type/khyperloglog/docs/khll.md(本文骨架,含布局图 khll_layout.png)
- 核心实现:KHyperLogLog.java
- 类型定义:KHyperLogLogType.java
- 标量函数:KHyperLogLogFunctions.java
- 聚合函数:KHyperLogLogAggregationFunction.java、KHyperLogLogWithLimitAggregationFunction.java、MergeKHyperLogLogAggregationFunction.java
- 状态与序列化:KHyperLogLogState.java、KHyperLogLogStateFactory.java、KHyperLogLogStateSerializer.java
- 官方函数文档:presto-docs/src/main/sphinx/functions/khyperloglog.rst
- 大数据
- 数据库
- 后端
【免费下载链接】presto
The official home of the Presto distributed SQL query engine for big data
相关推荐
Presto KHyperLogLog 函数完全指南:基于 MinHash 与 HyperLogLog 的双列关联数据草图
Presto KHyperLogLog 函数完全指南:基于 MinHash 与 HyperLogLog 的双列关联数据草图 导读 KHyperLogLog(KH
大数据数据库后端Presto Set Digest 函数完全指南:基于 MinHash 与 HyperLogLog 的集合相似度估算
Presto Set Digest 函数完全指南:基于 MinHash 与 HyperLogLog 的集合相似度估算 导读 本文深入讲解 Presto 分布式
大数据数据库后端Presto HyperLogLog 函数完全指南:approx_distinct 背后的数据草图与增量去重实战
Presto HyperLogLog 函数完全指南:approx_distinct 背后的数据草图与增量去重实战 HyperLogLog 是一种以固定内存估算海
大数据数据库后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考