EMR Serverless Spark は、Hologres Spark コネクタを介して Hologres に接続します。Serverless Spark 環境で Hologres のデータの読み取りと書き込みを行うために、必要な設定を追加できます。
制限事項
Spark コネクタには、Hologres バージョン 1.3 以降が必要です。インスタンスのバージョンは、Hologres コンソール の インスタンスの詳細 ページで確認できます。インスタンスがバージョン 1.3 より古い場合は、インスタンスをアップグレードするか、Hologres DingTalk グループ (ID: 32314975) に参加してアップグレードをリクエストしてください。
アクセス方法
Hologres には2つの方法でアクセスできます。ニーズに最も適した方法を選択してください:
|
アクセス方法 |
説明 |
シナリオ |
リファレンス |
|
方法1:タスク/セッションレベルの設定 |
各タスクまたはセッションで、Hologres 接続情報 (JDBC URL、ユーザー名、パスワード、その他のパラメーター) を個別に設定します。 |
|
このトピック |
|
方法2:データカタログによる統一設定 (推奨) |
EMR Serverless Spark の [Data Catalogs] 機能を使用して Hologres データカタログを追加します。カタログを追加すると、ワークスペース内のすべてのジョブとセッションが、デフォルトでアクセスが許可されたデータにアクセスできるようになります。 説明
エンジンバージョン esr-4.9.0 以降でのみサポートされます。 |
|
長期間にわたって Hologres データに頻繁にアクセスする必要があるワークスペースの場合、方法2 (データカタログ) を使用すると、設定の繰り返しが減り、開発効率が向上します。
操作手順
ステップ1:hologres-connector-spark JAR の取得とアップロード
EMR Serverless Spark esr-4.8.0 以降のバージョンには、Hologres コネクタが組み込まれています。この手順は、esr-4.8.0 より前のバージョンを使用している場合にのみ必要です。
-
Spark で Hologres のデータの読み取りと書き込みを行うには、コネクタの JAR ファイルが必要です。Maven Central Repository からダウンロードできます。このトピックでは、バージョン 1.5.6 の hologres-connector-spark-3.x-1.5.6-jar-with-dependencies.jar を使用します。
-
ダウンロードした hologres-connector-spark JAR ファイルを OSS にアップロードします。手順については、「単純アップロード」をご参照ください。
ステップ2:ネットワーク接続の追加
-
ネットワーク情報を取得します。
Hologres ページに移動し、対象の Hologres インスタンスのインスタンス詳細ページに移動して、VPC と vSwitch の情報を確認します。
-
ネットワーク接続を追加します。
Serverless Spark は、Hologres クラスターへのネットワーク接続を必要とします。詳細については、「EMR Serverless Sparkと他のVPC間のネットワーク接続」をご参照ください。
ステップ3:Hologresデータベースとテーブルの作成
-
Hologres インスタンスに接続します。詳細については、「インスタンスへの接続」をご参照ください。
-
[SQL Editor] タブで、新しい一時クエリに次の SQL 文を入力して実行します。
-- データベースを作成します。 CREATE DATABASE testdb; -- テーブルを作成します。 CREATE TABLE "public"."test" ( "id" text NULL, "name" text NULL); -- データを挿入します。 INSERT INTO public.test VALUES ('1001','jack'),('1002','tony'),('1003','mike'); -- データをクエリします。 SELECT * FROM public.test
ステップ4:Hologres のデータの読み取りと書き込み
例1:SQL セッション
SQL セッションを使用して Hologres のデータの読み取りと書き込みを行います。
-
SQL セッションを作成します。詳細については、「SQLセッションの管理」をご参照ください。
セッションを作成するときは、前の手順で作成したネットワーク接続をネットワーク接続リストから選択します。 Spark Configuration セクションで、hologres-connector-spark をロードするための次のパラメーターを追加します。
# hologres-connector JAR ファイルを追加します (esr-4.8.0 より前のバージョンでのみ必要です)。 spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar # Hologres カタログを設定します。 spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb次の表にパラメーターを示します。
パラメーター
例
説明
spark.emr.serverless.user.defined.jarsoss://<bucket>/hologres-connector-spark-3.x-<version>.jarユーザー定義 JAR ファイルへのパス。
spark.sql.catalog.hologres_external_test_dbcom.alibaba.hologres.spark3.HoloTableCatalogHologres データソースを Spark 3.x の外部カタログとして設定します。これは固定値です。
spark.sql.catalog.hologres_external_test_db.usernameLTAI******お使いの Alibaba Cloud アカウントの AccessKey ID。機密情報の処理にはシークレット管理を使用することを推奨します。詳細については、「シークレット管理を使用した機密情報の管理」をご参照ください。
spark.sql.catalog.hologres_external_test_db.passwordmXYV******お使いの Alibaba Cloud アcount の AccessKey Secret。機密情報の処理にはシークレット管理を使用することを推奨します。詳細については、「シークレット管理を使用した機密情報の管理」をご参照ください。
spark.sql.catalog.hologres_external_test_db.jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbHologres インスタンスの JDBC 接続 URL。
パラメーター名の中の
hologres_external_test_dbの部分はカスタマイズできます。 -
[Data Development] ページで [SparkSQL] ジョブを作成し、右上隅から作成した SQL セッションを選択します。
詳細については、「SparkSQL開発」をご参照ください。
-
次のコードを新しい SparkSQL タブにコピーし、Run をクリックします。
-- testdb データベースに切り替えます。 USE hologres_external_test_db; -- データを書き込みます。 INSERT INTO `public`.test VALUES ('1004','tom'); -- データをクエリします。 SELECT * FROM `public`.test;
例2:ストリーミングジョブ
この PySpark の例では、Kafka からデータを読み取り、ストリーミングジョブとして Hologres に書き込みます。
Kafka と Hologres 間のネットワーク接続が有効であることを確認してください。Kafka と Hologres を同じ VPC と vSwitch にデプロイすることを推奨します。
-
このコード例では、Kafka の情報と Hologres のテーブルを実際の値に置き換えます。
from pyspark.sql import SparkSession from pyspark.sql.functions import col # Kafka 情報を設定します。 servers = "alikafka-serverless-cn-xxxxx-vpc.alikafka.aliyuncs.com:9092" # ご使用の Kafka ブートストラップサーバーに置き換えます。 topic = "topic-name" # ご使用の Kafka トピックに置き換えます。 # SparkSession を作成します。 spark = SparkSession.builder \ .appName("test read kafka") \ .getOrCreate() # Kafka ストリームを読み取ります。 df = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", servers) \ .option("subscribe", topic) \ .load() # Hologres に書き込む関数を定義します (マイクロバッチごとに呼び出されます)。 def write_to_hologres(batch_df, batch_id): print(f"Writing batch {batch_id} to Hologres...") batch_df.write \ .format("hologres") \ .mode("append") \ .insertInto("hologres_external_test_db.public.test") # ご使用の Hologres テーブルに置き換えます。 # キーと値を文字列に変換してストリームを書き込みます。 query = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \ .writeStream \ .foreachBatch(write_to_hologres) \ .outputMode("append") \ .trigger(processingTime='30 seconds') \ .start() # ストリーミングクエリが終了するのを待ちます (ノートブックでは、これによりセルの実行がブロックされます)。 query.awaitTermination() -
ファイルをアップロードします。
-
Artifacts ページで、Upload File をクリックします。
-
Upload File ダイアログボックスで、アップロードエリアをクリックして前の手順の Python ファイルを選択するか、ファイルをアップロードエリアにドラッグします。
-
-
ストリーミングジョブを作成して実行します。
-
Development ページで、
(作成) アイコンをクリックします。 -
表示されるダイアログボックスで、Name を入力し、Application (Streaming) リストから PySpark を選択して、OK をクリックします。
-
新しい開発タブで、次のパラメーターを設定し、その他はデフォルト設定のままにします。次に、Publish をクリックします。
パラメーター
説明
[Main Python Resources]
前の手順でアップロードした Python ファイルを選択します。
Engine Version
互換性のある Spark バージョンを選択します。この例では
esr-4.6.0を使用します。Network Connection
ステップ2で作成したネットワーク接続を選択します。
[Spark Configuration]
# hologres-connector JAR ファイルを追加します (esr-4.8.0 より前のバージョンでのみ必要です)。 spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar # Hologres カタログを設定します。 spark.sql.catalog.hologres_external_test_db com.alibaba.hologres.spark3.HoloTableCatalog spark.sql.catalog.hologres_external_test_db.username *** spark.sql.catalog.hologres_external_test_db.password *** spark.sql.catalog.hologres_external_test_db.jdbcurl jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbパラメーターの詳細な説明については、「例1:SQL セッション」をご参照ください。
-
ジョブを公開したら、Go to O&M をクリックします。開いたページで Start をクリックします。
-
-
結果を検証します。
-
Kafka にメッセージを送信します。

-
Spark SQL を使用してデータをクエリします。

-
例3:Notebook セッション
-
Notebook セッションを作成します。詳細については、「Notebookセッションの管理」をご参照ください。
セッションを作成するときは、前の手順で作成したネットワーク接続をネットワーク接続リストから選択します。 Spark Configuration セクションで、hologres-connector-spark をロードするための次のパラメーターを追加します。
# hologres-connector JAR ファイルを追加します (esr-4.8.0 より前のバージョンでのみ必要です)。 spark.emr.serverless.user.defined.jars oss://<bucket>/hologres-connector-spark-3.x-<version>.jar -
[Data Development] ページで [Notebook] ジョブを作成し、右上隅から作成した Notebook セッションを選択します。
-
次のコードを新しい Notebook タブにコピーし、
をクリックします。import pandas as pd from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType # 1. Pandas DataFrame を準備します。 pdf = pd.DataFrame({ "id": ["1006"], "name": ["sl"] }) # 2. PySpark DataFrame に変換します。 # (オプション:正しいデータ型を保証するため、スキーマを明示的に定義します) schema = StructType([ StructField("id", StringType(), True), StructField("name", StringType(), True) ]) df = spark.createDataFrame(pdf, schema=schema) # Hologres に書き込みます。 df.write \ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb") \ .option("table", "test") \ .mode("append") \ .save() # データを読み取ります。 readDf = spark.read\ .format("hologres") \ .option("username", "LTAI******") \ .option("password", "mXYV******") \ .option("jdbcurl", "jdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdb") \ .option("table", "test") \ .load() readDf.select("id", "name").show(10)次の表にパラメーターを示します。
パラメーター
例
説明
usernameLTAI******お使いの Alibaba Cloud アカウントの AccessKey ID。機密情報の処理にはシークレット管理を使用することを推奨します。詳細については、「シークレット管理を使用した機密情報の管理」をご参照ください。
passwordmXYV******お使いの Alibaba Cloud アカウントの AccessKey Secret。機密情報の処理にはシークレット管理を使用することを推奨します。詳細については、「シークレット管理を使用した機密情報の管理」をご参照ください。
jdbcurljdbc:postgresql://hgpostcn-cn-***-vpc-st.hologres.aliyuncs.com:80/testdbHologres インスタンスの JDBC 接続 URL。
-
結果を検証します。

一般的なHologresカタログコマンド
Hologres カタログは、Hologres データベースを 外部カタログとして Spark SQL にマッピングします。各カタログは単一の Hologres データベースにバインドされ、データベース間のアクセスはサポートされていません。カタログ内の論理構造は Hologres と一致しています:
|
Sparkの概念 |
Hologresの概念 |
説明 |
|
カタログ |
データベース |
たとえば、 |
|
名前空間 |
スキーマ |
たとえば、 |
|
テーブル |
テーブル |
|
Hologresカタログの使用
Spark の Hologres カタログは、厳密に 1 つの Hologres データベースに対応し、作成後に変更することはできません。
USE hologres_external_test_db;
すべての名前空間の一覧表示
Spark の名前空間は Hologres のスキーマに対応します。デフォルトのスキーマは public です。USE コマンドを使用して、セッションのデフォルトスキーマを変更します。
-- Hologres カタログ内のすべての名前空間を表示します。これは Hologres 内のすべてのスキーマに対応します。
SHOW NAMESPACES;
名前空間内のテーブルの一覧表示
-
すべてのテーブルを一覧表示します。
SHOW TABLES; -
特定の名前空間のテーブルを一覧表示します。
USE test_schema; SHOW TABLES; -- または、次を使用します。 SHOW TABLES IN test_schema;
関連ドキュメント
Spark を使用した Hologres のデータの読み取りと書き込みの詳細については、「Sparkを使用したHologresのデータの読み取りと書き込み」をご参照ください。