全部产品
Search
文档中心

开源大数据平台E-MapReduce:使用RoaringBitmap实现精确去重

更新时间:Jun 24, 2026

Apache Paimon的主键表支持Aggregation Merge Engine,可在写入时为非主键字段指定聚合函数。对于精确去重场景,Paimon提供了rbm32和rbm64两种聚合函数,底层基于RoaringBitmap实现。本文介绍如何通过自定义UDF配合Paimon的bitmap聚合功能,在Spark SQL中实现精确去重。

背景信息

Paimon的Aggregation Merge Engine可以在写入时为非主键字段指定聚合函数(如sum、max、min等),Paimon会在compaction和读取时自动对相同主键的数据进行聚合合并。

Paimon提供了rbm32rbm64两种聚合函数,底层基于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