Todos os produtos
Search
Central de documentação

AnalyticDB:Desenvolva aplicações Spark com o SDK para Java

Última atualização: Jun 27, 2026

O AnalyticDB for MySQL Data Lakehouse Edition (V3.0) permite gerenciar jobs do Spark programaticamente por meio do SDK para Java. Este guia aborda como enviar um job do Spark, consultar seu status, recuperar logs e detalhes, listar jobs históricos e encerrar um job em execução.

Pré-requisitos

Antes de começar, verifique se você possui:

  • JDK 1.8 ou posterior instalado

  • Um cluster do AnalyticDB for MySQL Data Lakehouse Edition (V3.0). Consulte Criar um cluster Data Lakehouse Edition

  • Um grupo de recursos de job para o cluster. Consulte Criar um grupo de recursos

  • Um caminho de armazenamento de log configurado, utilizando um dos métodos a seguir:

    • No console do AnalyticDB for MySQL, acesse a página Spark JAR Development e clique em Log Settings no canto superior direito

    • Defina o parâmetro spark.app.log.rootPath como um caminho do Object Storage Service (OSS)

Adicione dependências do Maven

Adicione as seguintes dependências ao seu arquivo pom.xml. Utilize a versão 1.0.16 do adb20211201 para garantir estabilidade.

<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>

Configure a autenticação

Armazene suas credenciais AccessKey em variáveis de ambiente. Codificar credenciais diretamente no código-fonte traz riscos de exposição de informações sensíveis.

export ALIBABA_CLOUD_ACCESS_KEY_ID=<your-access-key-id>
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<your-access-key-secret>

Para obter instruções sobre como definir variáveis de ambiente no Linux, macOS e Windows, consulte Configurar variáveis de ambiente.

Inicialize o cliente lendo as credenciais das variáveis de ambiente:

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"));
// Replace with the region ID where your cluster resides
config.setRegionId("cn-hangzhou");
// Replace with the endpoint for that region
config.setEndpoint("adb.cn-hangzhou.aliyuncs.com");

Client client = new Client(config);

Operações do SDK

As seções a seguir apresentam cada operação como um exemplo independente. Todas seguem o mesmo padrão: construir um objeto de requisição, chamar o método do cliente e extrair o resultado do corpo da resposta.

Envie um job do Spark

Chame submitSparkApp() informando o id do cluster, o nome do grupo de recursos, a configuração do job e o tipo de job. O método retorna o id do job (appId), necessário para todas as operações subsequentes.

O parâmetro appType aceita dois valores:

Valor

Descrição

Batch

Job em lote do Spark

SQL

Job Spark SQL

Para jobs do tipo Batch, passe uma string JSON com a configuração do job. Para jobs SQL, passe uma instrução 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);
    // Save the returned appId — you need it for all subsequent operations
    return response.getBody().getData().getAppId();
}

O exemplo abaixo envia um job em lote SparkPi:

String clusterId = "amv-bp1mhnosdb38****";
String resourceGroupName = "test";

String data = "{\n" +
        "    \"comments\": [\"-- Here is just an example of SparkPi. Modify the content and run your spark program.\"],\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("Submitted job ID: " + appId);

Consulte o status do job

Utilize getSparkAppState() com o id do job para obter o status atual.

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();
}

Estados finais do job

Ao terminar, o job atinge um dos seguintes estados finais:

Estado

Descrição

COMPLETED

O job foi concluído com sucesso

FAILED

O job falhou

FATAL

Ocorreu um erro fatal durante a execução

Para verificar a conclusão via polling, crie um loop até que o job atinja um estado final (COMPLETED, FAILED ou FATAL):

long maxRunningTimeMs = 60000;  // 60 seconds
long pollIntervalMs = 2000;     // 2 seconds

String state;
long startTime = System.currentTimeMillis();
do {
    state = getAppState(appId, client);
    if (System.currentTimeMillis() - startTime > maxRunningTimeMs) {
        System.out.println("Timed out waiting for job to complete.");
        break;
    }
    System.out.println("Current state: " + state);
    Thread.sleep(pollIntervalMs);
} while (!"COMPLETED".equalsIgnoreCase(state)
        && !"FATAL".equalsIgnoreCase(state)
        && !"FAILED".equalsIgnoreCase(state));

Consulte os detalhes do job

Invoque getSparkAppInfo() para recuperar metadados do job, como o endereço da UI do Spark e os timestamps de início e término.

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();
}

Acesse campos específicos do objeto SparkAppInfo retornado:

SparkAppInfo appInfo = getAppInfo(appId, client);
System.out.println("State:          " + appInfo.getState());
System.out.println("Spark UI:       " + appInfo.getDetail().webUiAddress);
System.out.println("Submitted at:   " + appInfo.getDetail().submittedTimeInMillis);
System.out.println("Terminated at:  " + appInfo.getDetail().terminatedTimeInMillis);

Recupere os logs do job

Use getSparkAppLog() para obter o log do driver de um job.

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);

Liste jobs históricos

Execute listSparkApps() para recuperar jobs históricos do Spark de um cluster. Os resultados são paginados e a numeração das páginas começa em 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();
}
// Retrieve the first 50 jobs for the cluster
List<SparkAppInfo> jobs = listSparkApps(clusterId, 1, 50, client);
for (SparkAppInfo job : jobs) {
    System.out.printf("AppId: %s | State: %s | Spark UI: %s%n",
            job.getAppId(),
            job.getState(),
            job.getDetail().webUiAddress);
}

Encerre um job

Chame killSparkApp() passando o id do job para interromper uma execução em andamento.

import com.aliyun.adb20211201.models.KillSparkAppRequest;

KillSparkAppRequest request = new KillSparkAppRequest();
request.setAppId(appId);
client.killSparkApp(request);
System.out.println("Job terminated: " + appId);

Exemplo completo

O exemplo de ponta a ponta a seguir integra todas as operações: envio de um job SparkPi, aguardo pela conclusão, recuperação de detalhes e logs, listagem de jobs históricos e encerramento do job.

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("Submitting job: " + gson.toJson(request));
        SubmitSparkAppResponse response = client.submitSparkApp(request);
        System.out.println("Submit response: " + 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 {
        // Initialize the client using credentials from environment variables
        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\": [\"-- Here is just an example of SparkPi. Modify the content and run your spark program.\"],\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";

        // Step 1: Submit the job
        String appId = submitSparkApp(clusterId, resourceGroupName, data, "Batch", client);

        // Step 2: Poll until the job reaches a terminal state
        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("Timed out.");
                break;
            }
            System.out.println("Current state: " + state);
            Thread.sleep(pollIntervalMs);
        } while (!"COMPLETED".equalsIgnoreCase(state)
                && !"FATAL".equalsIgnoreCase(state)
                && !"FAILED".equalsIgnoreCase(state));

        // Step 3: Retrieve job details
        SparkAppInfo appInfo = getAppInfo(appId, client);
        System.out.printf("State: %s | Spark UI: %s | Submitted: %s | Terminated: %s%n",
                state,
                appInfo.getDetail().webUiAddress,
                appInfo.getDetail().submittedTimeInMillis,
                appInfo.getDetail().terminatedTimeInMillis);

        // Step 4: Retrieve driver logs
        String log = getAppDriverLog(appId, client);
        System.out.println(log);

        // Step 5: List historical jobs
        List<SparkAppInfo> jobs = listSparkApps(clusterId, 1, 50, client);
        jobs.forEach(job -> System.out.printf("AppId: %s | State: %s | Spark UI: %s%n",
                job.getAppId(),
                job.getState(),
                job.getDetail().webUiAddress));

        // Step 6: Terminate the job
        KillSparkAppRequest killRequest = new KillSparkAppRequest();
        killRequest.setAppId(appId);
        client.killSparkApp(killRequest);
    }
}

Próximos passos