Apache Paimon的主键表支持Aggregation Merge Engine,可在写入时为非主键字段指定聚合函数。对于精确去重场景,Paimon提供了rbm32和rbm64两种聚合函数,底层基于RoaringBitmap实现。本文介绍如何通过自定义UDF配合Paimon的bitmap聚合功能,在Spark SQL中实现精确去重。
背景信息
Paimon的Aggregation Merge Engine可以在写入时为非主键字段指定聚合函数(如sum、max、min等),Paimon会在compaction和读取时自动对相同主键的数据进行聚合合并。
Paimon提供了rbm32和rbm64两种聚合函数,底层基于RoaringBitmap实现。写入时只需将整数值序列化为bitmap字节写入BINARY列,Paimon引擎会在compaction时自动对相同主键的bitmap执行OR合并。由于Paimon未内置将整数转为bitmap字节的UDF,需要自行实现。
本文提供了一组示例UDF,演示如何在Spark SQL中配合Paimon的bitmap聚合功能实现精确去重。
前提条件
已创建EMR Serverless Spark工作空间。
已下载Paimon Bitmap UDF的JAR文件(paimon-spark-bitmap-1.0-SNAPSHOT.jar)并上传至OSS。
已创建SQL会话并启动。
已创建SparkSQL类型的任务。
UDF说明
本文提供4个统一的UDF,自动兼容32位和64位bitmap。
UDF | 功能 | 32/64位行为 |
to_bitmap | 将整数转为bitmap序列化字节。 | 输入INT时产出32位格式(配合rbm32),输入BIGINT时产出64位格式(配合rbm64)。 |
bitmap_count | 获取bitmap基数(去重后的元素数量)。 | 自动识别格式。 |
bitmap_contains | 判断指定值是否存在于bitmap中。 | 自动识别格式。 |
bitmap_to_string | 将bitmap转为逗号分隔的字符串(调试用)。 | 自动识别格式。 |
rbm32示例:用户页面访问去重
以下示例演示如何使用rbm32聚合函数,实现用户页面访问的精确去重统计。rbm32适用于ID值在INT(32位整数)范围内的场景。
步骤一:注册UDF
注册bitmap相关的UDF。本文使用临时UDF进行说明,您可以根据需要注册临时或永久UDF。
请将示例中的oss://bucket/path/to/paimon-spark-bitmap-1.0-SNAPSHOT.jar替换为JAR文件的实际OSS路径。
CREATE TEMPORARY FUNCTION to_bitmap AS 'org.apache.paimon.spark.bitmap.ToBitmapUDF' USING JAR 'oss://bucket/path/to/paimon-spark-bitmap-1.0-SNAPSHOT.jar';
CREATE TEMPORARY FUNCTION bitmap_count AS 'org.apache.paimon.spark.bitmap.BitmapCountUDF' USING JAR 'oss://bucket/path/to/paimon-spark-bitmap-1.0-SNAPSHOT.jar';
CREATE TEMPORARY FUNCTION bitmap_contains AS 'org.apache.paimon.spark.bitmap.BitmapContainsUDF' USING JAR 'oss://bucket/path/to/paimon-spark-bitmap-1.0-SNAPSHOT.jar';
CREATE TEMPORARY FUNCTION bitmap_to_string AS 'org.apache.paimon.spark.bitmap.BitmapToStringUDF' USING JAR 'oss://bucket/path/to/paimon-spark-bitmap-1.0-SNAPSHOT.jar';步骤二:创建Paimon聚合表
创建使用rbm32聚合函数的Paimon主键表。
CREATE TABLE user_page_visits (
user_id INT,
page_visits BINARY
) USING paimon
TBLPROPERTIES (
'primary-key' = 'user_id',
'merge-engine' = 'aggregation',
'fields.page_visits.aggregate-function' = 'rbm32'
);关键参数说明:
'primary-key' = 'user_id':指定user_id为主键。'merge-engine' = 'aggregation':使用聚合合并引擎。'fields.page_visits.aggregate-function' = 'rbm32':对page_visits字段使用rbm32聚合函数,适用于INT类型的值。
步骤三:写入数据
使用to_bitmap将页面ID转为bitmap字节后写入。
-- 初始写入
INSERT INTO user_page_visits VALUES
(1, to_bitmap(100, 101, 102)), -- 用户1访问了页面100, 101, 102
(2, to_bitmap(101, 103)), -- 用户2访问了页面101, 103
(3, to_bitmap(102, 104)); -- 用户3访问了页面102, 104
-- 追加写入(Paimon自动对相同user_id的bitmap做OR合并)
INSERT INTO user_page_visits VALUES
(1, to_bitmap(103, 105)), -- 用户1又访问了103, 105
(2, to_bitmap(104, 106)); -- 用户2又访问了104, 106步骤四:查询去重结果
查询每个用户访问了多少个不同页面。
SELECT user_id, bitmap_count(page_visits) AS unique_pages
FROM user_page_visits;预期结果如下。
user_id | unique_pages | 说明 |
1 | 5 | 100, 101, 102, 103, 105自动合并。 |
2 | 4 | 101, 103, 104, 106自动合并。 |
3 | 2 | 102, 104。 |
您还可以使用以下查询进行更多操作。
-- 查询用户是否访问过某个页面
SELECT user_id, bitmap_contains(page_visits, 101) AS visited_page_101
FROM user_page_visits;
-- 查看bitmap中的具体值(调试用)
SELECT user_id, bitmap_to_string(page_visits) AS pages
FROM user_page_visits;rbm64示例:大ID场景(如设备指纹)
当ID值超过INT范围时(如设备指纹ID),需要使用rbm64聚合函数,底层基于64位RoaringBitmap实现。rbm64与rbm32使用同一套UDF,无需额外注册。
如果在同一个会话中已经注册过上述4个UDF,此处无需重复注册。
步骤一:创建表
创建使用rbm64聚合函数的Paimon主键表。设备ID超过INT范围,因此使用rbm64。
CREATE TABLE app_device_visits (
app_id INT,
device_bitmap BINARY
) USING paimon
TBLPROPERTIES (
'primary-key' = 'app_id',
'merge-engine' = 'aggregation',
'fields.device_bitmap.aggregate-function' = 'rbm64'
);步骤二:写入数据
to_bitmap输入BIGINT时自动产出64位格式。
-- 初始写入
INSERT INTO app_device_visits VALUES
(1001, to_bitmap(CAST(8000000001 AS BIGINT), CAST(8000000002 AS BIGINT), CAST(8000000003 AS BIGINT))),
(1002, to_bitmap(CAST(8000000002 AS BIGINT), CAST(8000000004 AS BIGINT)));
-- 追加写入(自动合并)
INSERT INTO app_device_visits VALUES
(1001, to_bitmap(CAST(8000000004 AS BIGINT), CAST(8000000006 AS BIGINT)));步骤三:查询去重结果
SELECT app_id,
bitmap_count(device_bitmap) AS unique_devices,
bitmap_contains(device_bitmap, CAST(8000000002 AS BIGINT)) AS has_device
FROM app_device_visits;预期结果如下。
app_id | unique_devices | has_device | 说明 |
1001 | 5 | true | 合并后包含8000000001, 8000000002, 8000000003, 8000000004, 8000000006。 |
1002 | 2 | true | 包含8000000002, 8000000004。 |
您还可以查看bitmap中的具体值。
SELECT app_id, bitmap_to_string(device_bitmap) AS devices
FROM app_device_visits;相关文档
Paimon Aggregation Merge Engine的更多用法,请参见Aggregation。
RoaringBitmap的更多信息,请参见RoaringBitmap。