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.rootPathcomo 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 |
|
|
Job em lote do Spark |
|
|
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 |
|
|
O job foi concluído com sucesso |
|
|
O job falhou |
|
|
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);
}
}