Spark Thrift Server は Apache Spark が提供するサービスです。Java Database Connectivity (JDBC) またはオープンデータベースコネクティビティ (ODBC) を通じて接続および SQL クエリを実行できます。これにより、ご利用の Spark 環境を既存のビジネスインテリジェンス (BI) ツール、データビジュアライゼーションツール、その他のデータ分析ツールと簡単に統合できます。 Spark Thrift Server セッションを作成し、さまざまなクライアントから接続できます。
前提条件
ワークスペースを作成しました。詳細については、「ワークスペースの管理」をご参照ください。
Spark Thrift Server セッションの作成
Spark Thrift Server セッションを作成すると、Spark SQL ノードを作成する際に選択できます。
-
Sessions ページに移動します。
-
EMR コンソール にログインします。
-
左側のナビゲーションウィンドウで、EMR Serverless > Spark を選択します。
-
Spark ページで、対象のワークスペース名をクリックします。
-
EMR Serverless Spark ページで、左側のナビゲーションウィンドウの Sessions をクリックします。
-
-
Sessions ページで、Spark Thrift Server Session タブをクリックします。
-
Create Spark Thrift Server Session をクリックします。
-
Create Spark Thrift Server Session ページで、パラメーターを設定し、Create をクリックします。
パラメーター
説明
Name
新しい Spark Thrift Server の名前です。
名前は 1~64 文字で、英数字、ハイフン (-)、アンダースコア (_)、スペースを使用できます。
Resource Queue
セッションをデプロイする開発キュー、または開発環境と本番環境で共有されるキューを選択します。
キューの詳細については、「リソースキューの管理」をご参照ください。
Engine Version
このセッションのエンジンバージョンです。詳細については、「エンジンバージョン」をご参照ください。
Use Fusion Acceleration
Fusion を使用すると、Spark ワークロードを高速化し、ジョブの総コストを削減できます。請求情報については、「プロダクト課金」をご参照ください。Fusion エンジンの詳細については、「Fusion エンジン」をご参照ください。
Automatic Stop
デフォルトで有効です。Spark Thrift Server セッションは、45 分間アクティビティがない場合に自動的に停止します。
Normal Network Connection
VPC 内または外部サービスにあるデータソースにアクセスするための既存のネットワーク接続です。詳細については、「EMR Serverless Spark と他の VPC 間のネットワーク接続性」をご参照ください。
Spark Thrift Server Port
ポート番号は、パブリックエンドポイント経由でのアクセスでは 443、内部の同一リージョンエンドポイント経由でのアクセスでは 80 です。
Authentication Method
トークン方式のみサポートされています。
spark.driver.cores
ドライバープロセスの CPU コア数です。デフォルト値:1。
spark.driver.memory
ドライバープロセスに割り当てるメモリです。デフォルト値:3.5 GB。
spark.executor.cores
各エグゼキュータの CPU コア数です。デフォルト値:1。
spark.executor.memory
各エグゼキュータに割り当てるメモリです。デフォルト値:3.5 GB。
spark.executor.instances
割り当てるエグゼキュータ数です。デフォルト値:2。
Dynamic Resource Allocation
デフォルトで無効です。有効にした場合、以下のパラメーターを設定します。
-
Minimum Number of Executors:デフォルト値は 2 です。
-
Maximum Number of Executors:spark.executor.instances が設定されていない場合、デフォルト値は 10 です。
More Memory Configurations
-
spark.driver.memoryOverhead:ドライバーで利用可能な非ヒープメモリです。このパラメーターが設定されていない場合、Spark はデフォルトに基づいて自動的に値を割り当てます。
max(384 MB, 10% * spark.driver.memory)です。 -
spark.executor.memoryOverhead:各エグゼキュータで利用可能な非ヒープメモリです。このパラメーターが設定されていない場合、Spark はデフォルトに基づいて自動的に値を割り当てます。
max(384 MB, 10% * spark.executor.memory)です。 -
spark.memory.offHeap.size:Spark で利用可能なオフヒープメモリ量です。デフォルト値は 1 GB です。
このパラメーターは、
spark.memory.offHeap.enabledがtrueに設定されている場合にのみ有効です。Fusion エンジンを使用する場合、この機能はデフォルトで有効になり、1 GB のオフヒープメモリが割り当てられます。
Spark Configuration
Spark 設定情報を入力します。デフォルトでは、パラメーターはスペースで区切られます。例:
spark.sql.catalog.paimon.metastore dlf。 -
-
エンドポイント情報を取得します。
-
Spark Thrift Server Session タブで、作成した Spark Thrift Server の名前をクリックします。
-
Overview タブで、エンドポイント情報をコピーします。
ネットワーク環境に応じて、エンドポイントタイプを選択します。
-
パブリックエンドポイント:ローカルマシン、外部ネットワーク、またはクロスクラウド環境から EMR Serverless Spark にアクセスする場合に使用します。トラフィック料金が発生する場合があります。適切なセキュリティ対策を講じてください。
-
内部エンドポイント:同一リージョンの Alibaba Cloud ECS インスタンスから EMR Serverless Spark にアクセスする場合に使用します。内部アクセスは無料でより安全ですが、同一リージョンの Alibaba Cloud 内部ネットワーク内に限定されます。
-
-
トークンの作成
-
Spark Thrift Server Session タブで、作成した Spark Thrift Server セッションの名前をクリックします。
-
Token Management タブをクリックします。
-
Create Token をクリックします。
-
Create Token ダイアログボックスで、パラメーターを設定し、OK をクリックします。
パラメーター
説明
Name
新しいトークンの名前です。
Expired At
トークンの有効期限です。値は少なくとも 1 日以上である必要があります。デフォルトで有効になっており、有効期間は 365 日です。
-
トークン情報をコピーします。
重要トークン情報は作成直後に必ずコピーしてください。再度表示することはできません。トークンが有効期限切れになったり紛失したりした場合は、トークンを再作成またはリセットする必要があります。
Spark Thrift Server への接続
Spark Thrift Server に接続する際は、必要に応じて以下の情報を置き換えてください。
-
<endpoint>:Overview タブで取得した Endpoint(Public) または Endpoint(Private)。内部エンドポイントを使用する場合、Spark Thrift Server へのアクセスは同一 VPC 内のリソースに限定されます。
-
<port>:ポート番号。パブリックエンドポイント経由でのアクセスでは 443、内部の同一リージョンエンドポイント経由でのアクセスでは 80 です。 -
<username>:Token Management タブで作成したトークンの名前。 -
<token>:Token Management タブでコピーしたトークン情報。
Python を使用した Spark Thrift Server への接続
-
PyHive および Thrift パッケージをインストールするには、次のコマンドを実行します。
pip install pyhive thrift -
Spark Thrift Server に接続する Python スクリプトを作成します。
次のスクリプトは Hive に接続してすべてのデータベースを一覧表示します。ネットワーク環境に応じて接続方法を選択してください。
パブリックエンドポイントを使用した接続
from pyhive import hive if __name__ == '__main__': # <endpoint>、<username>、<token> を実際の情報に置き換えてください。 cursor = hive.connect('<endpoint>', port=443, scheme='https', username='<username>', password='<token>').cursor() cursor.execute('show databases') print(cursor.fetchall()) cursor.close()内部の同一リージョンエンドポイントを使用した接続
from pyhive import hive if __name__ == '__main__': # <endpoint>、<username>、<token> を実際の情報に置き換えてください。 cursor = hive.connect('<endpoint>', port=80, scheme='http', username='<username>', password='<token>').cursor() cursor.execute('show databases') print(cursor.fetchall()) cursor.close()
Java を使用した Spark Thrift Server への接続
-
pom.xmlファイルに次の Maven 依存関係を追加します。<dependencies> <dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-common</artifactId> <version>3.0.0</version> </dependency> <dependency> <groupId>org.apache.hive</groupId> <artifactId>hive-jdbc</artifactId> <version>2.1.0</version> </dependency> </dependencies>説明Serverless Spark の組み込み Hive バージョンは 2.x です。そのため、hive-jdbc 2.x のみがサポートされています。
-
Spark Thrift Server に接続する Java コードを作成します。
次のコードは Spark Thrift Server に接続してデータベースの一覧をクエリします。
パブリックエンドポイントを使用した接続
import java.sql.Connection; import java.sql.DriverManager; import java.sql.ResultSet; import java.sql.ResultSetMetaData; import org.apache.hive.jdbc.HiveStatement; public class Main { public static void main(String[] args) throws Exception { String url = "jdbc:hive2://<endpoint>:443/;transportMode=http;httpPath=cliservice/token/<token>"; Class.forName("org.apache.hive.jdbc.HiveDriver"); Connection conn = DriverManager.getConnection(url); HiveStatement stmt = (HiveStatement) conn.createStatement(); String sql = "show databases"; System.out.println("Running " + sql); ResultSet res = stmt.executeQuery(sql); ResultSetMetaData md = res.getMetaData(); String[] columns = new String[md.getColumnCount()]; for (int i = 0; i < columns.length; i++) { columns[i] = md.getColumnName(i + 1); } while (res.next()) { System.out.print("Row " + res.getRow() + "=["); for (int i = 0; i < columns.length; i++) { if (i != 0) { System.out.print(", "); } System.out.print(columns[i] + "='" + res.getObject(i + 1) + "'"); } System.out.println(")]"); } conn.close(); } }内部の同一リージョンエンドポイントを使用した接続
import java.sql.Connection; import java.sql.DriverManager; import java.sql.ResultSet; import java.sql.ResultSetMetaData; import org.apache.hive.jdbc.HiveStatement; public class Main { public static void main(String[] args) throws Exception { String url = "jdbc:hive2://<endpoint>:80/;transportMode=http;httpPath=cliservice/token/<token>"; Class.forName("org.apache.hive.jdbc.HiveDriver"); Connection conn = DriverManager.getConnection(url); HiveStatement stmt = (HiveStatement) conn.createStatement(); String sql = "show databases"; System.out.println("Running " + sql); ResultSet res = stmt.executeQuery(sql); ResultSetMetaData md = res.getMetaData(); String[] columns = new String[md.getColumnCount()]; for (int i = 0; i < columns.length; i++) { columns[i] = md.getColumnName(i + 1); } while (res.next()) { System.out.print("Row " + res.getRow() + "=["); for (int i = 0; i < columns.length; i++) { if (i != 0) { System.out.print(", "); } System.out.print(columns[i] + "='" + res.getObject(i + 1) + "'"); } System.out.println(")]"); } conn.close(); } }
Spark Beeline を使用した Spark Thrift Server への接続
-
セルフマネージドクラスターを使用している場合は、まず Spark の
binディレクトリに移動し、Beeline を使用して Spark Thrift Server に接続します。パブリックエンドポイントを使用した接続
cd /opt/apps/SPARK3/spark-3.4.2-hadoop3.2-1.0.3/bin/ ./beeline -u "jdbc:hive2://<endpoint>:443/;transportMode=http;httpPath=cliservice/token/<token>"内部の同一リージョンエンドポイントを使用した接続
cd /opt/apps/SPARK3/spark-3.4.2-hadoop3.2-1.0.3/bin/ ./beeline -u "jdbc:hive2://<endpoint>:80/;transportMode=http;httpPath=cliservice/token/<token>"説明コード内のパス
/opt/apps/SPARK3/spark-3.4.2-hadoop3.2-1.0.3は、EMR on ECS クラスター上の Spark インストールパスの例です。実際のクライアント上の Spark インストールパスに応じて調整してください。Spark インストールパスが不明な場合は、env | grep SPARK_HOMEコマンドを実行して確認できます。 -
EMR on ECS クラスターを使用している場合は、Spark Beeline クライアントを直接使用して Spark Thrift Server に接続できます。
パブリックエンドポイントを使用した接続
spark-beeline -u "jdbc:hive2://<endpoint>:443/;transportMode=http;httpPath=cliservice/token/<token>"内部の同一リージョンエンドポイントを使用した接続
spark-beeline -u "jdbc:hive2://<endpoint>:80/;transportMode=http;httpPath=cliservice/token/<token>"
Hive Beeline を使用して Serverless Spark Thrift Server に接続する際に次のエラーが発生する場合、原因は通常、Hive Beeline のバージョンと Spark Thrift Server の互換性の問題です。これを解決するには、Hive 2.x バージョンの Beeline を使用してください。
24/08/22 15:09:11 [main]: ERROR jdbc.HiveConnection: Error opening session
org.apache.thrift.transport.TTransportException: HTTP Response code: 404
Apache Superset を使用した Spark Thrift Server への接続の構成
Apache Superset は、幅広いチャートタイプをサポートするデータ探索および可視化プラットフォームです。詳細については、Superset ドキュメントをご参照ください。
-
依存関係をインストールします。
thriftパッケージのバージョン 0.20.0 がインストールされていることを確認してください。インストールされていない場合は、次のコマンドを実行してインストールします。pip install thrift==0.20.0 -
Superset を起動し、Superset インターフェイスにアクセスします。
Superset の起動方法の詳細については、Superset ドキュメントをご参照ください。
-
ページ右上隅の DATABASE をクリックして、Connect a database ページに移動します。
-
Connect a database ページで、Apache Spark SQL を選択します。

-
接続文字列を入力し、データソースパラメーターを構成します。
パブリックエンドポイントを使用した接続
hive+https://<username>:<token>@<endpoint>:443/<db_name>内部の同一リージョンエンドポイントを使用した接続
hive+http://<username>:<token>@<endpoint>:80/<db_name> -
FINISH をクリックして、接続と検証を確定します。
Hue を使用した Spark Thrift Server への接続の構成
Hue は Hadoop エコシステムと対話するためのオープンソース Web インターフェイスです。詳細については、Hue ドキュメントをご参照ください。
-
依存関係をインストールします。
thriftパッケージのバージョン 0.20.0 がインストールされていることを確認してください。インストールされていない場合は、次のコマンドを実行してインストールします。pip install thrift==0.20.0 -
Hue 設定ファイルに Spark SQL 接続文字列を追加します。
Hue 設定ファイル(通常は
/etc/hue/hue.conf)を見つけ、ファイルに次の内容を追加します。パブリックエンドポイントを使用した接続
[[[sparksql]]] name = Spark Sql interface=sqlalchemy options='{"url": "hive+https://<username>:<token>@<endpoint>:443/"}'内部の同一リージョンエンドポイントを使用した接続
[[[sparksql]]] name = Spark Sql interface=sqlalchemy options='{"url": "hive+http://<username>:<token>@<endpoint>:80/"}' -
Hue を再起動します。
設定を変更した後は、変更を有効にするために次のコマンドを実行して Hue サービスを再起動する必要があります。
sudo service hue restart -
接続を確認します。
正常に再起動した後、Hue インターフェイスにアクセスし、Spark SQL オプションを探します。設定が正しい場合、Spark Thrift Server に接続して SQL クエリを実行できます。

DataGrip を使用した Spark Thrift Server への接続
DataGrip は、ローカル、リモート、またはクラウド環境のデータベースをクエリ、作成、管理するためのデータベース管理ツールです。詳細については、DataGrip ウェブサイトをご参照ください。
-
DataGrip をインストールします。詳細については、「DataGrip のインストール」をご参照ください。
この例で使用する DataGrip のバージョンは 2025.1.2 です。
-
DataGrip クライアントを開きます。DataGrip インターフェイスが表示されます。
-
プロジェクトを作成します。
-
をクリックし、 を選択します。
-
New Project ダイアログボックスで、プロジェクト名(例:
Spark)を入力し、OK をクリックします。
-
-
[データベースエクスプローラー] メニューバーの
アイコンをクリックします。[データソース] > [その他] > [Apache Spark] を選択します。
-
Data Sources and Drivers ダイアログボックスで、次のパラメーターを構成します。

タブ
パラメーター
説明
General
Name
カスタム接続名。例:spark_thrift_server。
Authentication
認証方式を選択します。このトピックでは、No auth を選択します。
本番環境では、User & Password を選択して、許可されるユーザーのみが SQL ジョブを送信できるようにし、システムのセキュリティを向上させてください。
Driver
Apache Spark をクリックし、次に Go to Driver をクリックして、ドライバーバージョンが
ver. 1.2.2であることを確認します。説明現在の Serverless Spark エンジンバージョンは 3.x であるため、システムの安定性と機能の互換性を確保するために、ドライバーバージョン 1.2.2 を選択する必要があります。

URL
Spark Thrift Server に接続するための URL です。ネットワーク環境に応じて接続方法を選択します。
-
パブリックエンドポイントを使用した接続
jdbc:hive2://<endpoint>:443/;transportMode=http;httpPath=cliservice/token/<token> -
内部の同一リージョンエンドポイントを使用した接続
jdbc:hive2://<endpoint>:80/;transportMode=http;httpPath=cliservice/token/<token>
Options
Run keep-alive query
オプションです。これを有効にすると、アイドルタイムアウトによる自動切断を防ぐことができます。
-
-
Test Connection をクリックして、接続を検証します。

-
OK をクリックして、構成を完了します。
-
DataGrip を使用して Spark Thrift Server を管理できるようになりました。
DataGrip が Spark Thrift Server に接続されたら、データ開発を実行できます。詳細については、DataGrip ヘルプドキュメントをご参照ください。
たとえば、作成した接続で対象のテーブルを右クリックし、 を選択し、開いた SQL エディターで SQL スクリプトを作成・実行して、テーブルデータを表示します。

Redash を使用した Spark Thrift Server への接続
Redash は、Web ベースのデータベースクエリおよびデータ可視化のためのオープンソース BI ツールです。詳細については、Redash ドキュメントをご参照ください。
-
Redash をインストールします。詳細については、公式 Redash ドキュメントをご参照ください。
-
依存関係をインストールします。
thriftパッケージのバージョン 0.20.0 がインストールされていることを確認してください。インストールされていない場合は、次のコマンドを実行してインストールします。pip install thrift==0.20.0 -
Redash にログインします。
-
左側のナビゲーションウィンドウで Settings をクリックし、Data Sources タブで +New Data Source をクリックします。
-
表示されるダイアログボックスでパラメーターを構成し、Create をクリックします。

パラメーター
説明
Type Selection
データソースタイプ。検索ボックスで Hive (HTTP) を見つけて選択します。
Configuration
Name
データソース名。任意の名前を指定できます。
Host
Spark Thrift Server のエンドポイント。
Overview タブで Endpoint(Public) または Endpoint(Private) を取得できます。
Port
-
パブリックエンドポイント経由でアクセスする場合、ポート番号は 443 です。
-
内部の同一リージョンエンドポイント経由でアクセスする場合、ポート番号は 80 です。
HTTP Path
/cliserviceに設定します。Username
ユーザー名。任意の名前(例:
root)を入力できます。Password
作成したトークン情報を入力します。
HTTP Scheme
-
パブリックエンドポイント経由でアクセスする場合、
httpsに設定します。 -
内部の同一リージョンエンドポイントの場合、
httpを入力します。
-
-
ページ上部で を選択します。エディターで SQL ステートメントを記述できます。

dbt を使用した Spark Thrift Server への接続
dbt (data build tool) を使用すると、データアナリストおよびエンジニアが SQL ベースの変換ロジックを記述し、バージョン管理やテストなどのソフトウェアエンジニアリングのベストプラクティスを使用してデプロイメントを管理できます。詳細については、dbt ドキュメントをご参照ください。
-
dbt をインストールします。
pip install dbt-sparkHive コネクタを使用するには、次のようにインストールすることもできます。
pip install dbt-spark[HIVE] -
dbt プロジェクトを作成します。
dbt init my_spark_project cd my_spark_project -
dbt プロファイルを構成します。
~/.dbt/profiles.ymlファイルに Spark Thrift Server 接続情報を構成します。パブリックエンドポイントを使用した接続
my_spark_project: target: dev outputs: dev: type: spark method: thrift host: <endpoint> port: 443 user: <username> password: <token> schema: default connect_retries: 5 connect_timeout: 60 retry_all: true use_ssl: true server_side_parameters: "hive.exec.dynamic.partition": "true" "hive.exec.dynamic.partition.mode": "nonstrict"内部の同一リージョンエンドポイントを使用した接続
my_spark_project: target: dev outputs: dev: type: spark method: thrift host: <endpoint> port: 80 user: <username> password: <token> schema: default connect_retries: 5 connect_timeout: 60 retry_all: true use_ssl: false server_side_parameters: "hive.exec.dynamic.partition": "true" "hive.exec.dynamic.partition.mode": "nonstrict" -
接続をテストします。
dbt debug接続が成功すると、次のような出力が表示されます。
Connection test: [OK connection ok] -
dbt モデルを作成します。
models/ディレクトリに SQL ファイル(例:models/example_model.sql)を作成します。{{ config(materialized='table') }} select col1, col2, current_timestamp() as created_at from {{ ref('source_table') }} where col1 is not null -
dbt プロジェクトを実行します。
# Run all models dbt run # Run a specific model dbt run --models example_model # Run tests dbt test # Generate documentation dbt docs generate dbt docs serve -
ソーステーブルとテストを構成します。
models/schema.ymlファイルでソーステーブルとテストを定義します。version: 2 sources: - name: raw_data description: "生データ" tables: - name: source_table description: "生データを含むソーステーブル" columns: - name: col1 description: "主識別子" tests: - not_null - unique - name: col2 description: "データカラム" models: - name: example_model description: "変換済みデータモデル" columns: - name: col1 description: "主識別子" tests: - not_null - unique - name: col2 description: "処理済みデータカラム" - name: created_at description: "レコード作成タイムスタンプ"
注意事項:
-
dbt のバージョンが Spark Thrift Server と互換性があることを確認してください。dbt-spark 1.3.0 以降を使用してください。
-
本番環境では、トークンなどの機密情報を設定ファイルに直接記述せず、環境変数を使用して構成してください。
-
接続タイムアウトが発生する場合は、
connect_timeoutおよびconnect_retriesパラメーターを調整してください。 -
大規模なデータセットの場合は、パフォーマンスを向上させるために増分データモデルを使用してください。
{{ config(
materialized='incremental',
unique_key='id',
incremental_strategy='merge'
) }} select * from source_table {% if is_incremental() %} where updated_at > (select max(updated_at) from {{ this }}) {% endif %}
これらの構成を完了すると、dbt を使用して Spark Thrift Server に接続し、強力なデータ変換および管理機能を利用できます。
実行履歴の表示
ジョブが完了したら、その実行履歴を表示できます。
-
SQL Sessions ページで、目的のセッション名をクリックします。
-
Execution Records タブをクリックします。
このタブでは、実行ごとの詳細(実行 ID、開始時刻、Spark UI へのリンクなど)を表示できます。
