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

E-MapReduce:Paimon と Spark の連携

最終更新日:Jun 22, 2026

Paimon を使用して HDFS または OSS 上にデータレイクを構築し、Spark コンピュートエンジンを使用してデータを分析できます。このトピックでは、EMR の Spark SQL を使用して Paimon テーブルのデータの読み書きを行う方法について説明します。

前提条件

Spark と Paimon がインストールされた DataLake またはカスタム EMR クラスターを作成済みであること。詳細については、「クラスターの作成」をご参照ください。

制限事項

  • EMR-3.46.0 以降、または EMR-5.12.0 以降を実行している EMR クラスター上の Paimon に対して、Spark SQL を使用して読み書きができます。

  • Spark 3 の Spark SQL のみ、カタログを介して Paimon のデータの読み書きができます。

操作手順

ステップ 1:カタログの設定

Spark はカタログを介して Paimon テーブルの読み書きを行います。利用可能なカタログには、Paimon カタログと spark_catalog の 2 種類があります。ユースケースに応じていずれかを選択できます。

  • Paimon カタログ:Paimon フォーマットでメタデータを管理します。Paimon テーブルのクエリと書き込みにのみ使用できます。

  • spark_catalog:Spark のデフォルトの組み込みカタログです。通常、Spark SQL の内部テーブルのメタデータを管理し、Paimon テーブルと非 Paimon テーブルの両方のクエリと書き込みが可能です。

Paimon カタログ

メタデータは、HDFS などのファイルシステムや、OSS などのオブジェクトストレージサービスに保存できます。また、メタデータを DLF や Hive に同期して、他のサービスから Paimon にアクセスできるようにすることも可能です。

spark.sql.catalog.paimon.warehouse パラメーターは、ウェアハウスのルートパスを指定します。ルートパスが存在しない場合、システムによって自動的に作成されます。ルートパスが既に存在する場合、このカタログを使用してパス内の既存のテーブルにアクセスできます。

  1. SSH 経由でクラスターのマスターノードにログインします。詳細については、「クラスターへのログイン」をご参照ください。

  2. メタデータの種類に基づいて設定するカタログを選択します。対応するコマンドを実行して Spark SQL を起動します。

    ファイルシステムカタログ

    ファイルシステムカタログは、ファイルシステムまたはオブジェクトストレージにメタデータを格納します。

    spark-sql --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
    --conf spark.sql.catalog.paimon.metastore=filesystem \
    --conf spark.sql.catalog.paimon.warehouse=oss://<yourBucketName>/warehouse \
    --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
    説明
    • spark.sql.catalog.paimon:paimon という名前のカタログを定義します。

    • spark.sql.catalog.paimon.metastore:カタログのメタストアタイプを指定します。値を filesystem に設定すると、メタデータがファイルシステムに格納されることを意味します。

    • spark.sql.catalog.paimon.warehouse:データウェアハウスの場所を設定します。<yourBucketName> をご利用の OSS バケットの名前に置き換えてください。バケットの作成方法の詳細については、「バケットの作成」をご参照ください。

    DLF カタログ

    DLF カタログはメタデータを DLF に同期します。

    重要

    クラスターを作成する際、MetadataDLF 統合メタデータ に設定する必要があります。

    spark-sql --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
    --conf spark.sql.catalog.paimon.metastore=dlf \
    --conf spark.sql.catalog.paimon.warehouse=oss://<yourBucketName>/warehouse \
    --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
    説明
    • spark.sql.catalog.paimon:paimon という名前のカタログを定義します。

    • spark.sql.catalog.paimon.metastore:カタログのメタストアタイプを指定します。値を dlf に設定すると、メタデータが Data Lake Formation (DLF) に同期されることを意味します。

    • spark.sql.catalog.paimon.warehouse:データウェアハウスの場所を設定します。<yourBucketName> をご利用の OSS バケットの名前に置き換えてください。バケットの作成方法の詳細については、「バケットの作成」をご参照ください。

    Hive カタログ

    Hive カタログはメタデータを Hive メタストアに同期します。Hive カタログで作成されたテーブルは、Hive で直接クエリできます。Hive から Paimon をクエリする方法の詳細については、「Paimon と Hive の連携」をご参照ください。

    spark-sql --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
    --conf spark.sql.catalog.paimon.metastore=hive \
    --conf spark.sql.catalog.paimon.uri=thrift://master-1-1:9083 \
    --conf spark.sql.catalog.paimon.warehouse=oss://<yourBucketName>/warehouse \
    --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
    説明
    • spark.sql.catalog.paimon:paimon という名前のカタログを定義します。

    • spark.sql.catalog.paimon.metastore:カタログのメタストアタイプを指定します。値を hive に設定すると、メタデータが Hive メタストアに同期されることを意味します。

    • spark.sql.catalog.paimon.uri:Hive メタストアサービスのアドレスとポートです。値 thrift://master-1-1:9083 は、Spark SQL が master-1-1 ホストで実行され、ポート 9083 でリッスンしている Hive メタストアサービスに接続してメタデータを取得することを意味します。

    • spark.sql.catalog.paimon.warehouse:データウェアハウスの場所を設定します。<yourBucketName> をご利用の OSS バケットの名前に置き換えてください。バケットの作成方法の詳細については、「バケットの作成」をご参照ください。

spark_catalog

  1. SSH 経由でクラスターのマスターノードにログインします。詳細については、「クラスターへのログイン」をご参照ください。

  2. 次のコマンドを実行してカタログを設定し、Spark SQL を起動します。

    spark-sql --conf spark.sql.catalog.spark_catalog=org.apache.paimon.spark.SparkGenericCatalog \
    --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions
    説明
    • spark.sql.catalog.spark_catalog:spark_catalog という名前のカタログを定義します。

    • spark_catalog のウェアハウスのルートパスは、spark.sql.warehouse.dir パラメーターで指定されます。ほとんどの場合、このパラメーターを変更する必要はありません。

ステップ 2:Paimon テーブルの読み書き

次の Spark SQL ステートメントを実行してカタログにテーブルを作成し、そのテーブルのデータの読み書きを行います。

Paimon カタログ

Paimon テーブルにアクセスするには、paimon.<db_name>.<tbl_name> というフォーマットを使用します。ここで、<db_name> はデータベース名、<tbl_name> はテーブル名です。

-- データベースを作成します。
CREATE DATABASE IF NOT EXISTS paimon.ss_paimon_db;

-- Paimon テーブルを作成します。
CREATE TABLE paimon.ss_paimon_db.paimon_tbl (id INT, name STRING) USING paimon;

-- Paimon テーブルにデータを書き込みます。
INSERT INTO paimon.ss_paimon_db.paimon_tbl VALUES (1, "apple"), (2, "banana"), (3, "cherry");

-- 書き込み結果をクエリします。
SELECT * FROM paimon.ss_paimon_db.paimon_tbl ORDER BY id;

-- データベースを削除します。
DROP DATABASE paimon.ss_paimon_db CASCADE;
説明

Hive カタログを設定した後にデータベースを作成する際に metastore: Failed to connect to the Metastore Server エラーが発生した場合、Hive メタストアサービスが実行されていないことを意味します。次のコマンドを実行してサービスを起動する必要があります。サービスが起動したら、再度 Hive カタログを設定するコマンドを実行してください。

hive --service metastore &

クラスターの作成時に DLF 統合メタデータ を選択した場合は、DLF カタログを設定してメタデータを DLF に同期することを推奨します。

spark_catalog

spark_catalog.<db_name>.<tbl_name> を使用して、Paimon テーブルと非 Paimon テーブルの両方にアクセスできます。spark_catalog は Spark のデフォルトの組み込みカタログであるため、カタログ名を省略して <db_name>.<tbl_name> を使用して直接テーブルにアクセスできます。これらのフォーマットでは、<db_name> はデータベース名、<tbl_name> はテーブル名です。

-- データベースを作成します。
CREATE DATABASE IF NOT EXISTS ss_paimon_db;
CREATE DATABASE IF NOT EXISTS ss_parquet_db;

-- Paimon テーブルと Parquet テーブルを作成します。
CREATE TABLE ss_paimon_db.paimon_tbl (id INT, name STRING) USING paimon;
CREATE TABLE ss_parquet_db.parquet_tbl USING parquet AS SELECT 3, "cherry";

-- Paimon テーブルにデータを書き込みます。
INSERT INTO ss_paimon_db.paimon_tbl VALUES (1, "apple"), (2, "banana");
INSERT INTO ss_paimon_db.paimon_tbl SELECT * FROM ss_parquet_db.parquet_tbl;

-- 書き込み結果をクエリします。
SELECT * FROM ss_paimon_db.paimon_tbl ORDER BY id;

-- データベースを削除します。
DROP DATABASE ss_paimon_db CASCADE;
DROP DATABASE ss_parquet_db CASCADE;

クエリは次の結果を返します:

1       apple   
2       banana
3       cherry 

よくある質問

Paimon サービスをクラスターに追加すると、パラメーター spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions は自動的に設定されますか?

はい。Paimon サービスをクラスターに追加した後、次の手順に従って設定を表示します:

  1. 対象クラスターの Services タブに移動します。

  2. Spark サービスの設定を表示します。

    1. Spark サービスの右側にある Configure をクリックします。

    2. By Name 検索ボックスで spark.sql.extensions を検索して、設定を表示します。

      image

Spark Shell を使用して Paimon のデータの読み書きはできますか?

はい。Spark Shell を使用して Paimon のデータの読み書きを行うには、次の手順に従います:

  1. 次のコマンドを実行して Spark Shell を起動します。

    spark-shell
  2. Spark Shell で、次の Scala コードを実行して、指定されたディレクトリに格納されている Paimon テーブルへの書き込みとクエリを行います。

    val dataset = spark.read.format("paimon").load("oss://<yourBucketName>/warehouse/test_db.db/test_tbl")
    dataset.createOrReplaceTempView("test_tbl")
    spark.sql("INSERT INTO test_tbl VALUES (4, 'apple1', 3.5), (5, 'banana1', 4.0), (6, 'cherry1', 20.5)")
    spark.sql("SELECT * FROM test_tbl").show()
    説明
    • paimon:Paimon をデータ形式として指定する固定値です。

    • oss://<yourBucketName>/warehouse/test_db.db/test_tbl:Paimon テーブルへのパスです。<yourBucketName> はご利用の OSS バケットの名前です。

関連ドキュメント

Paimon の使用方法と設定の詳細については、「Apache Paimon ドキュメント」をご参照ください。