Apache Paimon のプライマリキーテーブルは Aggregation マージエンジンをサポートしており、書き込み時に非プライマリキー列の集計関数を指定できます。正確な重複排除のため、Paimon は RoaringBitmap に基づく 2 つの集計関数 rbm32 と rbm64 を提供しています。本トピックでは、カスタム UDF と Paimon のビットマップ集計を組み合わせて、Spark SQL で正確な重複排除を実行する方法を説明します。
背景情報
Paimon の Aggregation マージエンジンを使用すると、書き込み時に非プライマリキー列に集計関数 (sum、max、min) などを指定できます。コンパクション時と読み取り時に、Paimon はこれらの集計関数を適用し、同じプライマリキーを共有する行を自動的にマージします。
正確な重複排除のために、Paimon は 2 つの集計関数、rbm32 と rbm64 を提供します。どちらも 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」をご参照ください。