Este tópico descreve como usar a Open API para criar, atualizar e excluir tarefas do Data Integration que sincronizam dados de uma origem para um destino.
Pré-requisitos
Crie um workflow. Para mais informações, consulte Create a scheduled workflow.
Crie as fontes de dados necessárias para a tarefa de sincronização.
Limitações
Chame a operação CreateDISyncTask para criar uma tarefa do Data Integration. Configure o conteúdo da tarefa apenas no code editor.
A Open API não permite criar workflows no DataWorks. Utilize um workflow existente para a tarefa de sincronização de dados.
Obter um SDK
Obtenha o SDK mais recente no Alibaba Cloud Open API Portal.
Procedimento
Após configurar seu ambiente, siga estas etapas para criar e gerenciar uma tarefa de sincronização de dados com a API do DataWorks:
Etapas
-
Crie uma tarefa de integração de dados.
Chame a operação CreateDISyncTask para criar a tarefa. O código de exemplo a seguir mostra como configurar vários parâmetros principais. Para mais informações sobre os parâmetros, consulte CreateDISyncTask.
public void createFile() throws ClientException{ CreateDISyncTaskRequest request = new CreateDISyncTaskRequest(); request.setProjectId(181565L); request.setTaskType("DI_OFFLINE"); request.setTaskContent("{\"type\":\"job\",\"version\":\"2.0\",\"steps\":[{\"stepType\":\"mysql\",\"parameter\":{\"envType\":1,\"datasource\":\"dh_mysql\",\"column\":[\"id\",\"name\"],\"tableComment\":\"Comment for the table same\",\"connection\":[{\"datasource\":\"dh_mysql\",\"table\":[\"same\"]}],\"where\":\"\",\"splitPk\":\"id\",\"encoding\":\"UTF-8\"},\"name\":\"Reader\",\"category\":\"reader\"},{\"stepType\":\"odps\",\"parameter\":{\"partition\":\"pt=${bizdate}\",\"truncate\":true,\"datasource\":\"odps_source\",\"envType\":1,\"column\":[\"id\",\"name\"],\"emptyAsNull\":false,\"tableComment\":\"Comment for the table same\",\"table\":\"same\"},\"name\":\"Writer\",\"category\":\"writer\"}],\"setting\":{\"errorLimit\":{\"record\":\"\"},\"speed\":{\"throttle\":false,\"concurrent\":2}},\"order\":{\"hops\":[{\"from\":\"Reader\",\"to\":\"Writer\"}]}}"); request.setTaskParam("{\"FileFolderPath\":\"Workflow/new_biz/Data Integration\",\"ResourceGroup\":\"S_res_group_280749521950784_1602767279794\"}"); request.setTaskName("new_di_task_0607_1416"); String regionId = "cn-hangzhou"; // Please ensure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are set.https://www.alibabacloud.com/help/zh/alibaba-cloud-sdk-262060/latest/configure-credentials-378659 IClientProfile profile = DefaultProfile.getProfile("cn-shanghai", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")); DefaultProfile.addEndpoint("cn-hangzhou","dataworks-public","dataworks.cn-hangzhou.aliyuncs.com"); IAcsClient client; client = new DefaultAcsClient(profile); CreateDISyncTaskResponse response1 = client.getAcsResponse(request); Gson gson1 = new Gson(); System.out.println(gson1.toJson(response1)); } -
Chame a operação UpdateFile para atualizar os parâmetros de agendamento.
A tabela a seguir descreve os parâmetros da solicitação.
Nome
Tipo
Obrigatório
Exemplo
Descrição
Action
String
Sim
UpdateFile
Operação a executar.
FileFolderPath
String
Não
Workflows/MyFirstWorkflow/DataIntegration/Folder1/Folder2
Caminho do arquivo.
ProjectId
Long
Não
10000
ID do workspace do DataWorks. Acesse a página de gerenciamento de workspaces no DataWorks Console para encontrar o ID do workspace.
FileName
String
Não
ods_user_info_d
Nome do arquivo. Para renomear um arquivo, defina um novo valor para este parâmetro.
Encontre o
FileIdnecessário usando a operaçãoListFiles.FileDescription
String
Não
Esta é uma descrição de arquivo
Descrição do arquivo.
Content
String
Não
SELECT "1";
Código do arquivo. O formato varia conforme o tipo de arquivo (
fileType). Para ver um exemplo, clique em com o botão direito em uma tarefa no Operation Center e selecione View Code.AutoRerunTimes
Integer
Sim
3
Número de reexecuções automáticas após erro.
AutoRerunIntervalMillis
Integer
Não
120000
Intervalo entre reexecuções automáticas, em milissegundos. O valor máximo é 1.800.000 (30 minutos).
Este parâmetro corresponde ao intervalo de reexecução na aba Scheduling no DataWorks Console.
O console exibe esse valor em minutos. Converta a unidade ao usar a API.
RerunMode
String
Não
ALL_ALLOWED
Política de reexecução. Valores válidos:
-
ALL_ALLOWED: Permite reexecução após sucesso ou falha.
-
FAILURE_ALLOWED: Permite reexecução apenas após falha.
-
ALL_DENIED: Não permite reexecução.
Este parâmetro corresponde à propriedade de reexecução na aba Scheduling no DataWorks console.
Stop
Boolean
Não
false
Define se deve pausar o agendamento da tarefa. Valores válidos:
-
true: Pausa o agendamento.
-
false: Não pausa o agendamento.
Este parâmetro corresponde à configuração de pausa de agendamento na aba Scheduling no DataWorks Console.
ParaValue
String
Não
x=a y=b z=c
Parâmetros de agendamento.
Este parâmetro corresponde aos parâmetros na aba Scheduling no DataWorks Console. Para obter informações sobre como configurar parâmetros de agendamento, consulte Configurar parâmetros de agendamento.
StartEffectDate
Long
Não
936923400000
Horário de início do agendamento automático, especificado como timestamp Unix em milissegundos.
Este parâmetro corresponde à hora de início especificada para a data efetiva na aba Scheduling no DataWorks console.
EndEffectDate
Long
Não
4155787800000
Horário de término do agendamento automático, especificado como timestamp Unix em milissegundos.
Este parâmetro corresponde à hora de término especificada para a data efetiva na aba Scheduling no DataWorks console.
CronExpress
String
Não
00 00-59/5 1-23 ?
Expressão cron para agendamento periódico. Este parâmetro corresponde à expressão cron na aba Scheduling no DataWorks Console. Após configurar o agendamento e o horário programado, o DataWorks gera automaticamente uma expressão cron.
Exemplos:
-
Executar a tarefa às 05:30 todos os dias:
00 30 05 ?. -
Executar a tarefa aos 15 minutos de cada hora:
00 15 * ?. -
Executar a tarefa a cada 10 minutos:
00 00/10 * ?. -
Executar a tarefa a cada 10 minutos das 08:00 às 17:00 todos os dias:
00 00-59/10 8-17 * ?. -
Executar a tarefa às 00:20 no primeiro dia de cada mês:
00 20 00 1 * ?. -
Executar a tarefa a cada 3 meses, começando às 00:10 de 1º de janeiro:
00 10 00 1 1-12/3 ?. -
Executar a tarefa às 00:05 em todas as terças e sextas-feiras:
00 05 00 2,5.
O sistema de agendamento do DataWorks possui os seguintes limites para expressões cron:
-
O intervalo mínimo de agendamento é de 5 minutos.
-
O horário mais cedo em que uma tarefa pode ser agendada para execução é às 00:05 de cada dia.
CycleType
String
Não
NOT_DAY
Tipo do ciclo de agendamento. Os valores válidos são
NOT_DAY(minuto ou hora) eDAY(dia, semana ou mês).Este parâmetro corresponde ao agendamento na aba Scheduling no DataWorks Console.
DependentType
String
Não
USER_DEFINE
Modo de dependência entre ciclos. Valores válidos:
-
SELF: Selecione o nó atual como nó dependente.
-
CHILD: Selecione os nós filhos de primeiro nível como nós dependentes.
-
USER_DEFINE: Selecione outros nós como nós dependentes.
-
NONE: Nenhum nó dependente selecionado. A tarefa não depende do ciclo anterior.
DependentNodeIdList
String
Não
5,10,15,20
Se
DependentTypeforUSER_DEFINE, este parâmetro é obrigatório. Especifique os IDs dos nós dependentes, separados por vírgulas (,).
Este parâmetro corresponde às configurações de USER_DEFINE exibidas ao configurar a dependência entre ciclos na aba Scheduling no DataWorks console.InputList
String
Não
project_root,project.file1,project.001_out
Nomes de saída dos nós upstream dos quais este nó depende. Separe vários nomes com vírgulas (,).
Este parâmetro corresponde ao nome de saída do nó upstream na aba Scheduling no DataWorks console.NotaEste parâmetro é obrigatório ao criar uma tarefa de sincronização em lote com a operação
CreateDISyncTaskouUpdateFile.ProjectIdentifier
String
Não
dw_project
Nome do workspace do DataWorks. Faça login no DataWorks Console e acesse a página de configuração do workspace para obter o nome do workspace.
Defina este parâmetro ou o parâmetroProjectIdpara especificar o workspace do DataWorks.FileId
Long
Sim
100000001
ID do arquivo. Chame a operação
ListFilespara obter o ID do arquivo.OutputList
String
Não
dw_project.ods_user_info_d
Nome de saída do arquivo.
Este parâmetro corresponde ao nome de saída do nó atual na aba Scheduling no DataWorks console.ResourceGroupIdentifier
String
Não
default_group
Grupo de recursos onde a tarefa executa após a implantação. Chame a operação
ListResourceGroupspara obter os grupos de recursos disponíveis no workspace.ConnectionName
String
Não
odps_source
Identificador da fonte de dados usada pela tarefa. Chame a operação
ListDataSourcespara obter a lista de fontes de dados disponíveis.Owner
String
Não
18023848927592
ID de usuário do proprietário do arquivo.
AutoParsing
Boolean
Não
true
Define se deve ativar a análise automática para o arquivo. Valores válidos:
-
true: Analisa o código automaticamente. -
false: Não analisa o código automaticamente.
Este parâmetro corresponde à opção de analisar entrada/saída do código na aba Scheduling no DataWorks Console.
SchedulerType
String
Não
NORMAL
Tipo de agendamento. Valores válidos:
-
NORMAL: Tarefa agendada normalmente.
-
MANUAL: Tarefa manual. Executa apenas quando acionada manualmente e não possui agendamento automático.
-
PAUSE: Tarefa pausada.
-
SKIP: Tarefa de execução simulada (dry-run). Tarefas dry-run são agendadas regularmente, mas assumem status de sucesso imediatamente ao iniciar o agendamento.
AdvancedSettings
String
Não
{"queue":"default","SPARK_CONF":"--conf spark.driver.memory=2g"}
Configurações avançadas da tarefa.
Este parâmetro corresponde às configurações avançadas na página de edição de uma tarefa EMR Spark Streaming ou EMR Streaming SQL no DataWorks Console.
Atualmente, este parâmetro é suportado apenas para tarefas EMR Spark Streaming e EMR Streaming SQL. O valor do parâmetro deve estar no formato JSON.
StartImmediately
Boolean
Não
true
Define se deve iniciar a tarefa imediatamente após a implantação. Valores válidos:
-
true: Inicia a tarefa imediatamente após a implantação.
-
false: Não inicia a tarefa imediatamente após a implantação.
Este parâmetro corresponde ao modo de inicialização na aba Scheduling na página de edição de uma tarefa EMR Spark Streaming ou EMR Streaming SQL no DataWorks Console.
InputParameters
String
Não
[{"ValueSource": "project_001.first_node:bizdate_param","ParameterName": "bizdate_input"}]
Parâmetros de entrada de contexto para o nó. O valor é uma string JSON. Para obter informações sobre os campos, consulte o parâmetro
InputContextParameterListretornado pela operaçãoGetFile.Este parâmetro corresponde aos parâmetros de entrada do nó atual na aba Scheduling no DataWorks Console.
OutputParameters
String
Não
[{"Type": 1,"Value": "${bizdate}","ParameterName": "bizdate_param"}]
Parâmetros de saída de contexto para o nó. O valor é uma string JSON. Para obter informações sobre os campos, consulte o parâmetro
OutputContextParameterListretornado pela operaçãoGetFile.
Este parâmetro corresponde aos parâmetros de saída do nó atual na aba Scheduling no DataWorks Console. -
-
Envie a tarefa do Data Integration.
Chame a operação
SubmitFilepara enviar a tarefa do Data Integration ao ambiente de desenvolvimento do sistema de agendamento. A resposta retorna umdeploymentId. Use este ID com a operaçãoGetDeploymentpara recuperar os detalhes da implantação.public void submitFile() throws ClientException{ SubmitFileRequest request = new SubmitFileRequest(); request.setProjectId(78837L); request.setProjectIdentifier("zxy_8221431"); // This node ID is the ID that is returned when you create the node. It corresponds to the file_id in the File table of the database. request.setFileId(501576542L); request.setComment("Comment"); SubmitFileResponse acsResponse = client.getAcsResponse(request); // The DeploymentId is the return value of the commit or publish operation. Long deploymentId = acsResponse.getData(); log.info(acsResponse.toString()); }O código anterior exemplifica a configuração de alguns parâmetros. Para mais informações sobre os parâmetros, consulte SubmitFile e GetDeployment.
-
Implante a tarefa de sincronização no ambiente de produção.
Chame a operação
DeployFilepara publicar a tarefa de sincronização do Data Integration no ambiente de produção.NotaEsta operação aplica-se apenas a workspaces no modo Standard.
public void deploy() throws ClientException{ DeployFileRequest request = new DeployFileRequest(); request.setProjectIdentifier("zxy_8221431"); request.setFileId(501576542L); request.setComment("Comment"); // Specify either NodeId or file_id. The value of NodeId is the node ID in the basic properties of the scheduling configuration. request.setNodeId(700004537241L); DeployFileResponse acsResponse = client.getAcsResponse(request); // The DeploymentId is the return value of the commit or publish operation. Long deploymentId = acsResponse.getData(); log.info(acsResponse.getData().toString()); }O código anterior exemplifica a configuração de alguns parâmetros. Para mais informações sobre os parâmetros, consulte DeployFile.
-
Obtenha os detalhes do pacote de implantação.
A operação retorna um
deploymentId. Use este ID com a operaçãoGetDeploymentpara verificar o status da implantação. UmStatusigual a1indica uma implantação bem-sucedida.public void getDeployment() throws ClientException{ GetDeploymentRequest request = new GetDeploymentRequest(); request.setProjectId(78837L); request.setProjectIdentifier("zxy_8221431"); // The DeploymentId is the return value of the commit or publish operation. Call the GetDeployment operation to get the details of this deployment. request.setDeploymentId(2776067L); GetDeploymentResponse acsResponse = client.getAcsResponse(request); log.info(acsResponse.getData().toString()); }O código anterior exemplifica a configuração de alguns parâmetros. Para mais informações sobre os parâmetros, consulte GetDeployment.
Modificar a configuração de uma tarefa de sincronização
Para modificar uma tarefa, chame a operação UpdateDISyncTask para atualizar o Content da tarefa ou atualize o grupo de recursos exclusivo usando o parâmetro TaskParam. Após atualizar a tarefa, envie-a e implante-a novamente. Para mais informações, consulte Procedure.
Excluir uma tarefa de sincronização
Para excluir uma tarefa de sincronização do Data Integration, chame a operação DeleteFile.
Código de exemplo
import com.alibaba.fastjson.JSONObject;
import com.aliyun.dataworks_public20200518.Client;
import com.aliyun.dataworks_public20200518.models.*;
import com.aliyun.teaopenapi.models.Config;
public class createofflineTask {
static Long createTask(String fileName) throws Exception {
Long projectId = 2043L;
String taskType = "DI_OFFLINE";
String taskContent = "{\n" +
" \"type\": \"job\",\n" +
" \"version\": \"2.0\",\n" +
" \"steps\": [\n" +
" {\n" +
" \"stepType\": \"mysql\",\n" +
" \"parameter\": {\n" +
" \"envType\": 0,\n" +
" \"datasource\": \"mysql_autotest_dev\",\n" +
" \"column\": [\n" +
" \"id\",\n" +
" \"name\"\n" +
" ],\n" +
" \"connection\": [\n" +
" {\n" +
" \"datasource\": \"mysql_autotest_dev\",\n" +
" \"table\": [\n" +
" \"user\"\n" +
" ]\n" +
" }\n" +
" ],\n" +
" \"where\": \"\",\n" +
" \"splitPk\": \"\",\n" +
" \"encoding\": \"UTF-8\"\n" +
" },\n" +
" \"name\": \"Reader\",\n" +
" \"category\": \"reader\"\n" +
" },\n" +
" {\n" +
" \"stepType\": \"odps\",\n" +
" \"parameter\": {\n" +
" \"partition\": \"pt=${bizdate}\",\n" +
" \"truncate\": true,\n" +
" \"datasource\": \"odps_source\",\n" +
" \"envType\": 0,\n" +
" \"column\": [\n" +
" \"id\",\n" +
" \"name\"\n" +
" ],\n" +
" \"emptyAsNull\": false,\n" +
" \"tableComment\": \"null\",\n" +
" \"table\": \"user\"\n" +
" },\n" +
" \"name\": \"Writer\",\n" +
" \"category\": \"writer\"\n" +
" }\n" +
" ],\n" +
" \"setting\": {\n" +
" \"executeMode\": null,\n" +
" \"errorLimit\": {\n" +
" \"record\": \"\"\n" +
" },\n" +
" \"speed\": {\n" +
" \"concurrent\": 2,\n" +
" \"throttle\": false\n" +
" }\n" +
" },\n" +
" \"order\": {\n" +
" \"hops\": [\n" +
" {\n" +
" \"from\": \"Reader\",\n" +
" \"to\": \"Writer\"\n" +
" }\n" +
" ]\n" +
" }\n" +
"}";
CreateDISyncTaskRequest request = new CreateDISyncTaskRequest();
request.setProjectId(projectId);
request.setTaskType(taskType);
request.setTaskContent(taskContent);
request.setTaskName(fileName);
request.setTaskParam("{\"FileFolderPath\":\"Workflow/Automated_Testing_Workspace_Do_Not_Modify/Data_Integration\",\"ResourceGroup\":\"S_res_group_XXX\"}");
// Use an exclusive resource group for Data Integration.
CreateDISyncTaskResponse response1 = client.getAcsResponse(request);
return response1.getData().getFileId();
}
public static void updateFile(Long fileId) throws Exception {
UpdateFileRequest request = new UpdateFileRequest();
request.setProjectId(2043L);
request.setFileId(fileId);
request.setAutoRerunTimes(3);
request.setRerunMode("FAILURE_ALLOWED");
request.setCronExpress("00 30 05 * * ?");
request.setCycleType("DAY");
request.setResourceGroupIdentifier("S_res_group_XXX");
// Use an exclusive resource group for scheduling.
request.setInputList("dataworks_di_autotest_root");
request.setAutoParsing(true);
request.setDependentNodeIdList("5,10,15,20");
request.setDependentType("SELF");
request.setStartEffectDate(0L);
request.setEndEffectDate(4155787800000L);
request.setFileDescription("description");
request.setStop(false);
request.setParaValue("x=a y=b z=c");
request.setSchedulerType("NORMAL");
request.setAutoRerunIntervalMillis(120000);
UpdateFileResponse response1 = client.getAcsResponse(request);
}
static void deleteTask(Long fileId) throws Exception {
DeleteFileRequest request = new DeleteFileRequest();
Long projectId = 63845L;
request.setProjectId(projectId);
request.setFileId(fileId);
String akId = "XXX";
String akSecret = "XXXX";
String regionId = "cn-hangzhou";
IClientProfile profile = DefaultProfile.getProfile(regionId, akId, akSecret);
DefaultProfile.addEndpoint("cn-hangzhou","dataworks-public","dataworks.cn-hangzhou.aliyuncs.com");
IAcsClient client;
client = new DefaultAcsClient(profile);
DeleteFileResponse response1 = client.getAcsResponse(request);
System.out.println(JSONObject.toJSONString(response1));
}
static IAcsClient client;
public static void main(String[] args) throws Exception {
String akId = "XX";
// Please ensure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are set.https://www.alibabacloud.com/help/zh/alibaba-cloud-sdk-262060/latest/configure-credentials-378659
IClientProfile profile = DefaultProfile.getProfile("cn-shanghai", System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
DefaultProfile.addEndpoint(regionId, "dataworks-public", "dataworks." + regionId + ".aliyuncs.com");
client = new DefaultAcsClient(profile);
String taskName = "offline_job_0930_1648";
Long fileId = createTask(taskName); // Create a data integration task.
updateFile(fileId); // Modify the scheduling properties of the data integration task.
}
}