AnalyticDB for MySQL Data Lakehouse Edition (V3.0) では、Java SDK を使用して Spark ジョブをプログラムで管理できます。このガイドでは、Spark ジョブの送信、状態のポーリング、ログと詳細の取得、履歴ジョブの一覧表示、実行中のジョブの終了方法について説明します。
前提条件
開始する前に、以下を完了していることを確認してください。
-
JDK 1.8 以降がインストールされていること
-
AnalyticDB for MySQL Data Lakehouse Edition (V3.0) クラスターが作成されていること。詳細については、「Data Lakehouse Edition クラスターの作成」をご参照ください。
-
クラスター用のジョブリソースグループが作成されていること。詳細については、「リソースグループの作成」をご参照ください。
-
ログストレージパスが、次のいずれかの方法で設定されていること:
-
AnalyticDB for MySQL コンソールで、[Spark JAR 開発] ページに移動し、右上隅の [ログ設定] をクリックします
-
spark.app.log.rootPathパラメータを Object Storage Service (OSS) パスに設定します。
-
Maven 依存関係の追加
pom.xml ファイルに次の依存関係を追加します。安定性のため、adb20211201 のバージョン 1.0.16 を使用してください。
<dependencies>
<dependency>
<groupId>com.aliyun</groupId>
<artifactId>adb20211201</artifactId>
<version>1.0.16</version>
</dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>
<version>2.10.1</version>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.30</version>
</dependency>
</dependencies>
認証の設定
AccessKey 認証情報を環境変数に保存します。認証情報をソースコードにハードコーディングすると、機密情報が漏洩するリスクがあります。
export ALIBABA_CLOUD_ACCESS_KEY_ID=<your-access-key-id>
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<your-access-key-secret>
Linux、macOS、Windows で環境変数を設定する手順については、「環境変数の設定」をご参照ください。
環境変数から認証情報を読み込んでクライアントを初期化します。
import com.aliyun.adb20211201.Client;
import com.aliyun.teaopenapi.models.Config;
Config config = new Config();
config.setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
config.setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
// クラスターが存在するリージョン ID に置き換えます
config.setRegionId("cn-hangzhou");
// そのリージョンのエンドポイントに置き換えます
config.setEndpoint("adb.cn-hangzhou.aliyuncs.com");
Client client = new Client(config);
SDK の操作
次のセクションでは、各操作を個別の例として示します。すべての操作は同じパターンに従います。リクエストオブジェクトを作成し、クライアントメソッドを呼び出し、レスポンスボディから結果を取得します。
Spark ジョブの送信
クラスター ID、リソースグループ名、ジョブ設定、ジョブタイプを指定して submitSparkApp() を呼び出します。このメソッドはジョブ ID (appId) を返します。このジョブ ID は、以降のすべての操作で使用します。
appType には次の 2 つの値を指定できます:
| 値 | 説明 |
|---|---|
Batch |
Spark バッチジョブ |
SQL |
Spark SQL ジョブ |
Batch ジョブの場合は、ジョブ設定を含む JSON 文字列を渡します。SQL ジョブの場合は、SQL ステートメントを渡します。
import com.aliyun.adb20211201.models.SubmitSparkAppRequest;
import com.aliyun.adb20211201.models.SubmitSparkAppResponse;
@SneakyThrows
public static String submitSparkApp(String clusterId, String rgName, String data, String type, Client client) {
SubmitSparkAppRequest request = new SubmitSparkAppRequest();
request.setDBClusterId(clusterId);
request.setResourceGroupName(rgName);
request.setData(data);
request.setAppType(type);
SubmitSparkAppResponse response = client.submitSparkApp(request);
// 返された appId を保存します。以降のすべての操作で必要になります
return response.getBody().getData().getAppId();
}
次の例では、SparkPi バッチジョブを送信します。
String clusterId = "amv-bp1mhnosdb38****";
String resourceGroupName = "test";
String data = "{\n" +
" \"comments\": [\"-- SparkPi のサンプルです。内容はご自身のプログラムに合わせて変更してください。\"],\n" +
" \"args\": [\"1000\"],\n" +
" \"file\": \"local:///tmp/spark-examples.jar\",\n" +
" \"name\": \"SparkPi\",\n" +
" \"className\": \"org.apache.spark.examples.SparkPi\",\n" +
" \"conf\": {\n" +
" \"spark.driver.resourceSpec\": \"medium\",\n" +
" \"spark.executor.instances\": 2,\n" +
" \"spark.executor.resourceSpec\": \"medium\"}\n" +
"}\n";
String appId = submitSparkApp(clusterId, resourceGroupName, data, "Batch", client);
System.out.println("送信したジョブ ID: " + appId);
ジョブ状態の照会
ジョブ ID を指定して getSparkAppState() を呼び出し、現在の状態を取得します。
import com.aliyun.adb20211201.models.GetSparkAppStateRequest;
import com.aliyun.adb20211201.models.GetSparkAppStateResponse;
@SneakyThrows
public static String getAppState(String appId, Client client) {
GetSparkAppStateRequest request = new GetSparkAppStateRequest();
request.setAppId(appId);
GetSparkAppStateResponse response = client.getSparkAppState(request);
return response.getBody().getData().getState();
}
ジョブの最終状態
ジョブが終了すると、次のいずれかの最終状態になります:
| 状態 | 説明 |
|---|---|
COMPLETED |
ジョブが正常に終了しました |
FAILED |
ジョブが失敗しました |
FATAL |
ジョブで致命的なエラーが発生しました |
ポーリングで完了を確認するには、ジョブが最終状態 (COMPLETED、FAILED、または FATAL) になるまでループします。
long maxRunningTimeMs = 60000; // 60 秒
long pollIntervalMs = 2000; // 2 秒
String state;
long startTime = System.currentTimeMillis();
do {
state = getAppState(appId, client);
if (System.currentTimeMillis() - startTime > maxRunningTimeMs) {
System.out.println("ジョブの完了を待機中にタイムアウトしました。");
break;
}
System.out.println("現在の状態: " + state);
Thread.sleep(pollIntervalMs);
} while (!"COMPLETED".equalsIgnoreCase(state)
&& !"FATAL".equalsIgnoreCase(state)
&& !"FAILED".equalsIgnoreCase(state));
ジョブの詳細の照会
getSparkAppInfo() を呼び出して、Spark UI アドレスや開始/終了タイムスタンプなどのジョブメタデータを取得します。
import com.aliyun.adb20211201.models.GetSparkAppInfoRequest;
import com.aliyun.adb20211201.models.GetSparkAppInfoResponse;
import com.aliyun.adb20211201.models.SparkAppInfo;
@SneakyThrows
public static SparkAppInfo getAppInfo(String appId, Client client) {
GetSparkAppInfoRequest request = new GetSparkAppInfoRequest();
request.setAppId(appId);
GetSparkAppInfoResponse response = client.getSparkAppInfo(request);
return response.getBody().getData();
}
返された SparkAppInfo オブジェクトから特定のフィールドにアクセスします。
SparkAppInfo appInfo = getAppInfo(appId, client);
System.out.println("状態: " + appInfo.getState());
System.out.println("Spark UI: " + appInfo.getDetail().webUiAddress);
System.out.println("送信日時: " + appInfo.getDetail().submittedTimeInMillis);
System.out.println("終了日時: " + appInfo.getDetail().terminatedTimeInMillis);
ジョブログの取得
getSparkAppLog() を呼び出して、ジョブのドライバーログを取得します。
import com.aliyun.adb20211201.models.GetSparkAppLogRequest;
import com.aliyun.adb20211201.models.GetSparkAppLogResponse;
@SneakyThrows
public static String getAppDriverLog(String appId, Client client) {
GetSparkAppLogRequest request = new GetSparkAppLogRequest();
request.setAppId(appId);
GetSparkAppLogResponse response = client.getSparkAppLog(request);
return response.getBody().getData().getLogContent();
}String log = getAppDriverLog(appId, client);
System.out.println(log);
履歴ジョブの一覧表示
listSparkApps() を呼び出して、クラスターの履歴 Spark ジョブを取得します。結果はページ分割されます。ページ番号は 1 から始まります。
import com.aliyun.adb20211201.models.ListSparkAppsRequest;
import com.aliyun.adb20211201.models.ListSparkAppsResponse;
import java.util.List;
@SneakyThrows
public static List<SparkAppInfo> listSparkApps(String clusterId, long pageNumber, long pageSize, Client client) {
ListSparkAppsRequest request = new ListSparkAppsRequest();
request.setDBClusterId(clusterId);
request.setPageNumber(pageNumber);
request.setPageSize(pageSize);
ListSparkAppsResponse response = client.listSparkApps(request);
return response.getBody().getData().getAppInfoList();
}// クラスターの最初の 50 件のジョブを取得します
List<SparkAppInfo> jobs = listSparkApps(clusterId, 1, 50, client);
for (SparkAppInfo job : jobs) {
System.out.printf("AppId: %s | 状態: %s | Spark UI: %s%n",
job.getAppId(),
job.getState(),
job.getDetail().webUiAddress);
}
ジョブの終了
ジョブ ID を指定して killSparkApp() を呼び出し、実行中のジョブを停止します。
import com.aliyun.adb20211201.models.KillSparkAppRequest;
KillSparkAppRequest request = new KillSparkAppRequest();
request.setAppId(appId);
client.killSparkApp(request);
System.out.println("ジョブを終了しました: " + appId);
完全な例
次のエンドツーエンドの例では、すべての操作をまとめています。SparkPi ジョブを送信し、完了するまで待機し、詳細とログを取得し、履歴ジョブを一覧表示します。
import com.aliyun.adb20211201.Client;
import com.aliyun.adb20211201.models.*;
import com.aliyun.teaopenapi.models.Config;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
import lombok.SneakyThrows;
import java.util.List;
public class SparkExample {
private static Gson gson = new GsonBuilder().disableHtmlEscaping().setPrettyPrinting().create();
@SneakyThrows
public static String submitSparkApp(String clusterId, String rgName, String data, String type, Client client) {
SubmitSparkAppRequest request = new SubmitSparkAppRequest();
request.setDBClusterId(clusterId);
request.setResourceGroupName(rgName);
request.setData(data);
request.setAppType(type);
System.out.println("ジョブを送信中: " + gson.toJson(request));
SubmitSparkAppResponse response = client.submitSparkApp(request);
System.out.println("送信レスポンス: " + gson.toJson(response));
return response.getBody().getData().getAppId();
}
@SneakyThrows
public static String getAppState(String appId, Client client) {
GetSparkAppStateRequest request = new GetSparkAppStateRequest();
request.setAppId(appId);
GetSparkAppStateResponse response = client.getSparkAppState(request);
return response.getBody().getData().getState();
}
@SneakyThrows
public static SparkAppInfo getAppInfo(String appId, Client client) {
GetSparkAppInfoRequest request = new GetSparkAppInfoRequest();
request.setAppId(appId);
GetSparkAppInfoResponse response = client.getSparkAppInfo(request);
return response.getBody().getData();
}
@SneakyThrows
public static String getAppDriverLog(String appId, Client client) {
GetSparkAppLogRequest request = new GetSparkAppLogRequest();
request.setAppId(appId);
GetSparkAppLogResponse response = client.getSparkAppLog(request);
return response.getBody().getData().getLogContent();
}
@SneakyThrows
public static List<SparkAppInfo> listSparkApps(String clusterId, long pageNumber, long pageSize, Client client) {
ListSparkAppsRequest request = new ListSparkAppsRequest();
request.setDBClusterId(clusterId);
request.setPageNumber(pageNumber);
request.setPageSize(pageSize);
ListSparkAppsResponse response = client.listSparkApps(request);
return response.getBody().getData().getAppInfoList();
}
public static void main(String[] args) throws Exception {
// 環境変数から認証情報を使用してクライアントを初期化します
Config config = new Config();
config.setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
config.setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
config.setRegionId("cn-hangzhou");
config.setEndpoint("adb.cn-hangzhou.aliyuncs.com");
Client client = new Client(config);
String clusterId = "amv-bp1mhnosdb38****";
String resourceGroupName = "test";
String data = "{\n" +
" \"comments\": [\"-- SparkPi のサンプルです。内容はご自身のプログラムに合わせて変更してください。\"],\n" +
" \"args\": [\"1000\"],\n" +
" \"file\": \"local:///tmp/spark-examples.jar\",\n" +
" \"name\": \"SparkPi\",\n" +
" \"className\": \"org.apache.spark.examples.SparkPi\",\n" +
" \"conf\": {\n" +
" \"spark.driver.resourceSpec\": \"medium\",\n" +
" \"spark.executor.instances\": 2,\n" +
" \"spark.executor.resourceSpec\": \"medium\"}\n" +
"}\n";
// ステップ 1: ジョブを送信します
String appId = submitSparkApp(clusterId, resourceGroupName, data, "Batch", client);
// ステップ 2: ジョブが最終状態になるまでポーリングします
long maxRunningTimeMs = 60000;
long pollIntervalMs = 2000;
String state;
long startTime = System.currentTimeMillis();
do {
state = getAppState(appId, client);
if (System.currentTimeMillis() - startTime > maxRunningTimeMs) {
System.out.println("タイムアウトしました。");
break;
}
System.out.println("現在の状態: " + state);
Thread.sleep(pollIntervalMs);
} while (!"COMPLETED".equalsIgnoreCase(state)
&& !"FATAL".equalsIgnoreCase(state)
&& !"FAILED".equalsIgnoreCase(state));
// ステップ 3: ジョブの詳細を取得します
SparkAppInfo appInfo = getAppInfo(appId, client);
System.out.printf("状態: %s | Spark UI: %s | 送信日時: %s | 終了日時: %s%n",
state,
appInfo.getDetail().webUiAddress,
appInfo.getDetail().submittedTimeInMillis,
appInfo.getDetail().terminatedTimeInMillis);
// ステップ 4: ドライバーログを取得します
String log = getAppDriverLog(appId, client);
System.out.println(log);
// ステップ 5: 履歴ジョブを一覧表示します
List<SparkAppInfo> jobs = listSparkApps(clusterId, 1, 50, client);
jobs.forEach(job -> System.out.printf("AppId: %s | 状態: %s | Spark UI: %s%n",
job.getAppId(),
job.getState(),
job.getDetail().webUiAddress));
}
}