すべてのプロダクト
Search
ドキュメントセンター

E-MapReduce:RoaringBitmap を使用した正確な重複排除

最終更新日:Jun 25, 2026

Apache Paimon のプライマリキーテーブルは Aggregation マージエンジンをサポートしており、書き込み時に非プライマリキー列の集計関数を指定できます。正確な重複排除のため、Paimon は RoaringBitmap に基づく 2 つの集計関数 rbm32 と rbm64 を提供しています。本トピックでは、カスタム UDF と Paimon のビットマップ集計を組み合わせて、Spark SQL で正確な重複排除を実行する方法を説明します。

背景情報

Paimon の Aggregation マージエンジンを使用すると、書き込み時に非プライマリキー列に集計関数 (sum、max、min) などを指定できます。コンパクション時と読み取り時に、Paimon はこれらの集計関数を適用し、同じプライマリキーを共有する行を自動的にマージします。

正確な重複排除のために、Paimon は 2 つの集計関数、rbm32rbm64 を提供します。どちらも RoaringBitmap に基づいて構築されています。書き込み時に、整数値をビットマップバイトにシリアライズし、BINARY カラムに書き込みます。Paimon エンジンは、コンパクション中に、同じプライマリーキーを持つビットマップに対して自動的に OR マージを実行します。Paimon には整数をビットマップバイトに変換するための組み込み UDF が含まれていないため、自身で実装する必要があります。

本トピックでは、Paimon のビットマップ集計機能を使用して Spark SQL で正確な重複排除を実行する方法を示す、サンプル UDF のセットを提供します。

前提条件

  • EMR Serverless Spark ワークスペースが作成済みであること。

  • Paimon Bitmap UDF の JAR ファイル (paimon-spark-bitmap-1.0-SNAPSHOT.jar) がダウンロードされ、OSS にアップロード済みであること。

  • SQL Compute セッションが作成され、起動済みであること。

  • Spark SQL セッションタイプのタスクが作成済みであること。

UDF リファレンス

本トピックでは、32 ビットと 64 ビットの両方のビットマップに自動的に対応する、4 つの統合 UDF を提供します。

UDF

説明

32 ビットと 64 ビットの動作

to_bitmap

整数をシリアル化されたビットマップバイトに変換します。

入力タイプが INT の場合は 32 ビット形式 (rbm32 で使用) を、BIGINT の場合は 64 ビット形式 (rbm64 で使用) を生成します。

bitmap_count

ビットマップのカーディナリティ (一意な要素の数) を返します。

形式を自動的に検出します。

bitmap_contains

指定された値がビットマップに存在するかどうかを確認します。

形式を自動的に検出します。

bitmap_to_string

ビットマップをカンマ区切りの文字列に変換します。この UDF はデバッグ用です。

形式を自動的に検出します。

rbm32 の例:ユーザーのページ訪問の重複排除

次の例では、rbm32 集計関数を使用してユーザーのページ訪問に対する正確な重複排除を実行する方法を示します。rbm32 は、ID が INT (32 ビット整数) の範囲内であるシナリオに適しています。

ステップ 1:UDF の登録

ビットマップ関連の UDF を登録します。この例では一時 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';

ステップ 2: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': Aggregation マージエンジンを有効にします。

  • 'fields.page_visits.aggregate-function' = 'rbm32'page_visits 列に rbm32 集計関数を適用し、INT 値を処理します。

ステップ 3:データの書き込み

to_bitmap を使用してページ ID をビットマップバイトに変換し、テーブルに書き込みます。

-- 初回書き込み
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 のビットマップを自動的に OR マージします)
INSERT INTO user_page_visits VALUES
    (1, to_bitmap(103, 105)),        -- ユーザー 1 がページ 103、105 も訪問
    (2, to_bitmap(104, 106));        -- ユーザー 2 がページ 104、106 も訪問

ステップ 4:重複排除結果のクエリ

各ユーザーが訪問した一意のページ数をクエリします。

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;

-- デバッグ用にビットマップ内の値を表示
SELECT user_id, bitmap_to_string(page_visits) AS pages
FROM user_page_visits;

rbm64 の例:デバイスフィンガープリントなどの大規模 ID シナリオ

デバイスフィンガープリント ID など、ID が INT の範囲を超える場合は、64 ビット RoaringBitmap に基づく rbm64 集計関数を使用します。rbm64 は rbm32 と同じ UDF セットを共有するため、追加の登録は不要です。

説明

前述の 4 つの UDF が現在のセッションに登録済みの場合、再登録は不要です。

ステップ 1:テーブルの作成

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'
);

ステップ 2:データの書き込み

BIGINT 値を to_bitmap に渡すと、自動的に 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)));

ステップ 3:重複排除結果のクエリ

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 を含みます。

ビットマップ内の値を表示することもできます。

SELECT app_id, bitmap_to_string(device_bitmap) AS devices
FROM app_device_visits;

関連ドキュメント

  • Paimon の Aggregation マージエンジンの詳細については、「Aggregation」をご参照ください。

  • RoaringBitmap の詳細については、「RoaringBitmap」をご参照ください。