StarRocks は公式の Spark コネクタを提供しており、Spark と StarRocks 間のデータ転送が可能です。EMR Serverless Spark に必要な設定を追加することで、StarRocks インスタンスに接続できます。このトピックでは、EMR Serverless Spark を使用して StarRocks のデータを読み書きする方法について説明します。
アクセス方法
必要に応じて、2 つの方法で StarRocks にアクセスできます。
方法1:ジョブ/セッションレベルでの設定
この方法では、JDBC URL、ユーザー名、パスワードなどの StarRocks 接続情報を各ジョブまたはセッションで設定します。この方法は、次のシナリオに適しています。
-
一時的またはアドホックなデータアクセス
-
異なるジョブが異なる StarRocks クラスターにアクセスする必要がある場合
-
各ジョブに対してきめ細かなアクセス制御が必要な場合
方法2:データカタログによる一元化された設定
-
サポートされるエンジンバージョンは、esr-4.8.0 以降および esr-5.2.0 以降です。
EMR Serverless Spark の データカタログ 機能を使用して、StarRocks データカタログを追加します。一度設定すると、ワークスペース内のすべてのジョブとセッションは、デフォルトでデータカタログ内の承認済みデータにアクセスできるようになり、各ジョブで設定を繰り返す必要がなくなります。詳細については、「データカタログの管理」をご参照ください。
この方法は、次のシナリオに適しています。
-
StarRocks データに頻繁にアクセスする必要がある場合
-
複数のジョブが同じ StarRocks アクセス設定を共有する場合
-
ジョブ設定を簡素化し、開発効率を向上させたい場合
StarRocks への長期的かつ頻繁なアクセスを必要とするワークスペースでは、設定のオーバーヘッドを削減し、開発効率を向上させるために、データカタログ (方法 2) の使用を推奨します。
前提条件
-
Serverless Spark ワークスペースを作成します。詳細については、「ワークスペースの作成」をご参照ください。
-
EMR Serverless StarRocks インスタンスを作成します。詳細については、「インスタンスの作成」をご参照ください。
制限事項
Serverless Spark のエンジンバージョンは、esr-2.5.0、esr-3.1.0、esr-4.1.0、またはそれ以降のバージョンである必要があります。
手順
手順1:Spark コネクタ JAR の取得とアップロード
-
お使いのエンジンバージョンに対応する Spark コネクタの JAR ファイルをダウンロードします。詳細については、「」「Spark コネクタを使用した StarRocks からのデータ読み取り」をご参照ください。
説明エンジンバージョンが esr-4.8.0 以降および esr-5.2.0 以降の場合、Spark コネクタ JAR が組み込まれているため、この手順は省略できます。
この例では、Maven Central Repository からビルド済みの JAR をダウンロードします。
説明コネクタの JAR ファイル名は
starrocks-spark-connector-${spark_version}_${scala_version}-${connector_version}.jarという形式です。たとえば、エンジンバージョンが esr-4.1.0 (Spark 3.5.2, Scala 2.12) で、コネクタバージョンが 1.1.2 の場合は、starrocks-spark-connector-3.5_2.12-1.1.2.jarを選択します。 -
ダウンロードした Spark コネクタ JAR を OSS にアップロードします。詳細については、「簡易アップロード」をご参照ください。
手順2:ネットワーク接続の追加
-
ネットワーク情報を取得します。
EMR Serverless StarRocks のページで、対象の StarRocks インスタンスの Instance Details ページに移動し、Virtual Private Cloud (VPC) と vSwitch の情報を取得します。
-
ネットワーク接続を作成します。
-
EMR Serverless Spark のページで、お使いの Spark ワークスペースの Network Connection ページに移動し、Create Network Connection をクリックします。
-
Create Network Connection ダイアログボックスで、Name を入力し、StarRocks インスタンスの VPC と vSwitch を選択して、OK をクリックします。
説明ネットワーク接続は StarRocks インスタンスと一致している必要があります。StarRocks インスタンスと同じ VPC 内にある vSwitch を選択します。現在のアベイラビリティーゾーンに利用可能な vSwitch がない場合は、vSwitch をクリックして VPC コンソールに移動し、vSwitch を作成します。詳細については、「VPC と vSwitch」をご参照ください。
-
手順3:データベースとテーブルの作成
-
StarRocks インスタンスに接続します。詳細については、「EMR StarRocks Manager を使用した StarRocks インスタンスへの接続」をご参照ください。
-
[SQL Editor] の Queries ページで、File または右側の
アイコンをクリックし、OK をクリックして新しいファイルを作成します。 -
新しいファイルに次の SQL 文を入力し、Run をクリックします。
CREATE DATABASE `testdb`; CREATE TABLE `testdb`.`score_board` ( `id` int(11) NOT NULL COMMENT "", `name` varchar(65533) NULL DEFAULT "" COMMENT "", `score` int(11) NOT NULL DEFAULT "0" COMMENT "" ) ENGINE=OLAP PRIMARY KEY(`id`) COMMENT "OLAP" DISTRIBUTED BY HASH(`id`);
手順4:StarRocks の読み取りと書き込み
方法1:SQL セッションとノートブックセッション
セッションタイプの詳細については、「セッションマネージャー」をご参照ください。
SQL セッション
-
EMR Serverless Spark を使用して StarRocks にデータを書き込みます。
-
SQL セッションを作成します。詳細については、「SQL セッションの管理」をご参照ください。
セッションを作成する際、Spark コネクタのバージョンに対応するエンジンバージョンを選択し、Network Connection で前の手順で作成したネットワーク接続を選択し、[Spark Configuration]に以下のパラメーターを追加して Spark コネクタをロードします。
spark.emr.serverless.user.defined.jars oss://<bucketname>/path/connector.jaross://<bucketname>/path/connector.jarを、手順 1 でアップロードした Spark コネクタの OSS パスに置き換えます。例:oss://emr-oss/spark/starrocks-spark-connector-3.5_2.12-1.1.2.jar。 -
Development ページで、[SparkSQL] ジョブを作成します。次に、右上で作成した SQL セッションを選択します。
詳細については、「SparkSQL ジョブの開発」をご参照ください。
-
次のコードを新しい SparkSQL タブにコピーし、必要に応じてプレースホルダーの値を置き換えてから、Run をクリックします。
CREATE TEMPORARY VIEW score_board USING starrocks OPTIONS ( "starrocks.table.identifier" = "testdb.score_board", "starrocks.fe.http.url" = "<fe_host>:<fe_http_port>", "starrocks.fe.jdbc.url" = "jdbc:mysql://<fe_host>:<fe_query_port>", "starrocks.user" = "<user>", "starrocks.password" = "<password>" ); INSERT INTO `score_board` VALUES (1, "starrocks", 100), (2, "spark", 100);パラメーターの説明:
-
<fe_host>:EMR Serverless StarRocks インスタンスの FE ノードの内部エンドポイントまたはパブリックエンドポイントです。この値は、Instance Details ページの FE 詳細 セクションで確認できます。-
内部エンドポイントを使用する場合は、両方のサービスが同じ VPC 内にあることを確認してください。
-
パブリックエンドポイントを使用する場合は、セキュリティグループのルールで必要なポートでのトラフィックが許可されていることを確認してください。詳細については、「ネットワークアクセスとセキュリティ設定」をご参照ください。
-
-
<fe_http_port>:EMR Serverless StarRocks インスタンスの FE ノードの HTTP ポートです。デフォルトは 8030 です。この値は、Instance Details ページの FE 詳細 セクションで確認できます。 -
<fe_query_port>:EMR Serverless StarRocks インスタンスの FE ノードのクエリポートです。デフォルトは 9030 です。この値は、Instance Details ページの FE 詳細 セクションで確認できます。 -
<user>:Serverless StarRocks インスタンスのユーザー名です。デフォルトでは、管理者権限を持つ admin ユーザーが用意されています。User Management ページを使用して新しいユーザーを追加して接続することもできます。ユーザーの追加方法の詳細については、「ユーザーとデータ認可の管理」をご参照ください。 -
<password>:<user>に対応するパスワードです。
-
-
-
EMR Serverless Spark で書き込んだデータをクエリします。
この例では、Spark SQL ジョブで
test_viewという名前の一時ビューを作成し、score_boardからデータをクエリします。次のコードを新しい Spark SQL タブにコピーし、コードを選択してから、Run the selected をクリックします。CREATE TEMPORARY VIEW test_view USING starrocks OPTIONS ( "starrocks.table.identifier" = "testdb.score_board", "starrocks.fe.http.url" = "<fe_host>:<fe_http_port>", "starrocks.fe.jdbc.url" = "jdbc:mysql://<fe_host>:<fe_query_port>", "starrocks.user" = "<user>", "starrocks.password" = "<password>" ); SELECT * FROM test_view;出力
クエリにより、(1, "starrocks", 100) と (2, "spark", 100) の 2 つのレコードが返されます。
ノートブックセッション
-
EMR Serverless Spark を使用して StarRocks にデータを書き込みます。
-
ノートブックセッションを作成します。詳細については、「ノートブックセッションの管理」をご参照ください。
セッションを作成する際、Spark コネクタのバージョンに対応するエンジンバージョンを選択し、Network Connection で前の手順で作成したネットワーク接続を選択し、[Spark Configuration]に以下のパラメーターを追加して Spark コネクタをロードします。
spark.emr.serverless.user.defined.jars oss://<bucketname>/path/connector.jaross://<bucketname>/path/connector.jarを、手順 1 でアップロードした Spark コネクタの OSS パスに置き換えます。例:oss://emr-oss/spark/starrocks-spark-connector-3.5_2.12-1.1.2.jar。 -
Development ページで、[Interactive Development] > [Notebook] ジョブを作成し、右上で作成したノートブックセッションを選択します。
詳細については、「ノートブックセッションの管理」をご参照ください。
-
次のコードをノートブックの Python セルにコピーし、Run をクリックします。
# プレースホルダーをお使いのEMR Serverless StarRocksの設定に置き換えます。 fe_host = "<fe_host>" fe_http_port = "<fe_http_port>" fe_query_port = "<fe_query_port>" user = "<user>" password = "<password>" # ビューを作成します。 create_table_sql = f""" CREATE TEMPORARY VIEW score_board USING starrocks OPTIONS ( "starrocks.table.identifier" = "testdb.score_board", "starrocks.fe.http.url" = "{fe_host}:{fe_http_port}", "starrocks.fe.jdbc.url" = "jdbc:mysql://{fe_host}:{fe_query_port}", "starrocks.user" = "{user}", "starrocks.password" = "{password}" ) """ spark.sql(create_table_sql) # データを挿入します。 insert_data_sql = """ INSERT INTO `score_board` VALUES (1, "starrocks", 100), (2, "spark", 100) """ spark.sql(insert_data_sql)パラメーターの説明:
-
<fe_host>:EMR Serverless StarRocks インスタンスの FE ノードの内部エンドポイントまたはパブリックエンドポイントです。この値は、Instance Details ページの FE 詳細 セクションで確認できます。-
内部エンドポイントを使用する場合は、両方のサービスが同じ VPC 内にあることを確認してください。
-
パブリックエンドポイントを使用する場合は、セキュリティグループのルールで必要なポートでのトラフィックが許可されていることを確認してください。詳細については、「ネットワークアクセスとセキュリティ設定」をご参照ください。
-
-
<fe_http_port>:EMR Serverless StarRocks インスタンスの FE ノードの HTTP ポートです。デフォルトは 8030 です。この値は、Instance Details ページの FE 詳細 セクションで確認できます。 -
<fe_query_port>:EMR Serverless StarRocks インスタンスの FE ノードのクエリポートです。デフォルトは 9030 です。この値は、Instance Details ページの FE 詳細 セクションで確認できます。 -
<user>:Serverless StarRocks インスタンスのユーザー名です。デフォルトでは、管理者権限を持つ admin ユーザーが用意されています。User Management ページを使用して新しいユーザーを追加して接続することもできます。ユーザーの追加方法の詳細については、「ユーザーとデータ認可の管理」をご参照ください。 -
<password>:<user>に対応するパスワードです。
-
-
-
EMR Serverless Spark で書き込んだデータをクエリします。
新しい Python セルで、
test_viewという名前の一時ビューを作成してscore_boardテーブルをクエリします。次のコードをセルにコピーし、実行アイコン (
) をクリックします。# ビューを作成します。 create_view_sql=f""" CREATE TEMPORARY VIEW test_view USING starrocks OPTIONS ( "starrocks.table.identifier" = "testdb.score_board", "starrocks.fe.http.url" = "{fe_host}:{fe_http_port}", "starrocks.fe.jdbc.url" = "jdbc:mysql://{fe_host}:{fe_query_port}", "starrocks.user" = "{user}", "starrocks.password" = "{password}" ) """ spark.sql(create_view_sql) # データをクエリします。 query_sql="SELECT * FROM test_view" result_df = spark.sql(query_sql) result_df.show()出力
+---+---------+-----+ | id| name|score| +---+---------+-----+ | 2| spark| 100| | 1|starrocks| 100| +---+---------+-----+
方法2:Spark バッチジョブ
-
Spark バッチジョブを作成します。
-
EMR Serverless Spark のページで、左側メニューで Development をクリックします。
-
Development タブで、
アイコンをクリックします。 -
[Create] ダイアログボックスで、Name を入力し、タイプを に設定し、OK をクリックします。
必要に応じてタイプを調整できます。このトピックでは、SQL を例に説明します。ジョブタイプの詳細については、「アプリケーションの開発」をご参照ください。
-
-
Spark バッチジョブを使用して StarRocks の読み書きを行います。
-
新しいジョブ開発ページの右上で、キューを選択します。
キューの追加方法の詳細については、「リソースキューの管理」をご参照ください。
-
新しいジョブ開発ページで、以下の設定を行い、他のパラメーターはデフォルトのままにして、実行 をクリックします。
パラメーター
説明
SQL File
この例では、SQL セッションで使用した SQL 文を含む spark_sql_starrocks.sql ファイルを使用します。ファイルを使用する前に、ダウンロードして必要に応じて設定を変更し、Artifacts ページにアップロードしてください。
[Engine Version]
Spark コネクタのバージョンと一致するエンジンバージョンを選択します。
[Normal Network Connection]
前の手順で作成したネットワーク接続を選択します。
[Spark Configuration]
[Spark Configuration]セクションに以下のパラメーターを追加して Spark コネクタをロードします。
spark.emr.serverless.user.defined.jars oss://<bucketname>/path/connector.jaross://<bucketname>/path/connector.jarを、手順 1 でアップロードした Spark コネクタの OSS パスに置き換えます。例:oss://emr-oss/spark/starrocks-spark-connector-3.5_2.12-1.1.2.jar。
-
-
ログ情報を表示します。
-
画面下部の Execution Records セクションで、[Actions] 列の Details をクリックします。
-
Log Exploration タブをクリックして、ジョブのログ情報を表示します。
[Driver Log] > [Stdout] を選択して、stdout.log ファイルを表示します。ジョブが成功すると、ログ出力に
starrocks 100とspark 100が含まれます。これにより、データが StarRocks に正常に書き込まれ、その後読み取られたことが確認できます。
-
参考資料
詳細については、StarRocks の公式ドキュメントをご参照ください。