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

AnalyticDB:Java SDK を使用した Spark アプリケーションの開発

最終更新日:Aug 27, 2026

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 ジョブで致命的なエラーが発生しました

ポーリングで完了を確認するには、ジョブが最終状態 (COMPLETEDFAILED、または 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));
    }
}

次のステップ