Descubra como gerir dependências com Cloud Workflows utilizando o Google Cloud APIs e sensores inteligentes
Introdução
Gerir as dependências entre vários workflows é um desafio comum em pipelines de dados complexos.
Esta publicação demonstra como utilizar a Cloud Workflows e o Google Cloud APIs para orquestrar processos complexos, destacando os benefícios de dividir workloads em workflows separados para melhorar a modularidade, a otimização dos custos e a resiliência.
Embora o preço da Cloud Workflows seja baseado no número de etapas executadas, as vantagens de um pipeline bem orquestrado ultrapassam de longe qualquer pequeno aumento de custos.
Antes de começar:
- Já deve ter o Cloud Workflows a funcionar corretamente no seu projeto de GCP.
Funções do utilizador:
- Workflows Editor: Esta função permite-lhe criar, editar e implementar workflows.
- Workflows Invoker: Esta função permite-lhe executar (run) workflows.
- Logging Viewer: Esta função permite-lhe ler e visualizar logs.
Funções da service account:
- Workflows Invoker
- Logging Log Writer
- Certifique-se de que o seu utilizador tem permissão para atuar como service account.
Para testar o código, certifique-se de que cumpre os requisitos acima referidos.
Amostra Workflow 1
Vamos começar por construir dois Cloud Workflows simples e dependentes. Crie um novo workflow, escolha uma região de baixo carbono e selecione a service account para o lançamento. (Não é necessário configurar variáveis de ambiente, etiquetas ou accionadores para este exemplo).
O workflow pré-construído inclui código funcional, que vamos modificar ligeiramente. Copie e cole o seguinte código:
main:
steps:
- init:
assign:
- searchTerm: "europe"
- readWikipedia:
call: http.get
args:
url: '<https://en.wikipedia.org/w/api.php>'
query:
action: opensearch
search: '${searchTerm}'
result: wikiResult
- returnOutput:
return: {"wiki":'${wikiResult.body[1]}'}
Este workflow pesquisa na Wikipedia por artigos relacionados com a “europa”. Implemente o workflow sem alterações e deverá obter resultados semelhantes.
Amostra Workflow 2
De seguida, vamos criar um segundo workflow que se baseie no output do primeiro. Este workflow irá converter uma lista de cadeias de caracteres para maiúsculas, demonstrando uma dependência simples. (Cenários mais complexos são possíveis, mas não são o foco aqui). Siga o mesmo processo de criação do primeiro workflow para começar.
Substitua o código predefinido por este:
main:
params: [wikiResults]
steps:
- initializeUppercaseResults:
assign:
- uppercaseResults: []
- convertToUppercase:
for:
value: wikiItem
in: ${wikiResults.wiki}
steps:
- appendToUppercaseResults:
assign:
- uppercaseResults: ${list.concat(uppercaseResults, text.to_upper(wikiItem))}
- returnOutput:
return: ${uppercaseResults}
Este código converte uma lista de cadeias de caracteres em maiúsculas. Implemente o workflow e tente executá-lo sem enviar um input. A execução vai falhar porque, no código, usamos o dicionário de input dentro do ciclo for e neste caso está vazio. No entanto, deverá funcionar com o seguinte input:
{
"wiki": [
"Europe",
"European Union",
"European Parliament",
"European Commission",
"European colonization of the Americas",
"European emission standards",
"European People's Party Group",
"European Cup and UEFA Champions League records and statistics",
"European Convention on Human Rights",
"European debt crisis"
]
}
O resultado deve ser o seguinte:
Orquestrador do Workflow
Com ambos os workflows implementados e testados, podemos agora criar um workflow para acionar automaticamente o workflow 2 após a conclusão bem sucedida do workflow 1. Para isso, vamos precisar do ID e da localização de ambos os workflows. Crie um novo workflow com o nome “orquestrador-workflow” (a descrição é opcional). Escolha uma região com baixo teor de carbono e a mesma service account utilizada para os workflows anteriores.
Neste exemplo, uma vez que precisamos de passar o output do workflow 1 para o workflow 2, o código é um pouco mais complexo do que um simples ciclo for. Caso a partilha de dados não seja necessária entre os diferentes workflows, no final, é fornecido um código do orquestrador de workflow simplificado.
Adicione o seguinte código:
main:
params: [args]
steps:
- init:
assign:
- project_id: ${sys.get_env("GOOGLE_CLOUD_PROJECT_ID")}
- firstWorkflowId: "sample-workflow-1"
- secondWorkflowId: "sample-workflow-2"
- location: "europe-west1"
- launch_first_workflow_execution:
call: http.post
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/projects/"+project_id+"/locations/"+location+"/workflows/"+firstWorkflowId+"/executions>"}'
result: firstWorkflowResult
- set_first_execution_id:
assign:
- execution_id: ${firstWorkflowResult.body.name}
- exponential_backoff_retries_for_first_workflow:
try:
steps:
- get_first_execution_state:
call: http.get
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/>"+execution_id}'
result: firstWorkflowState
- set_first_execution_state:
assign:
- execution_state: ${firstWorkflowState.body.state}
- log_state_first_workflow:
call: sys.log
args:
text: '${firstWorkflowId +" state: "+ execution_state}'
severity: INFO
- induce_backoff_retry_if_first_state_not_done:
switch:
- condition: ${execution_state != "SUCCEEDED"}
raise: ${execution_state}
retry:
predicate: ${execution_state_predicate}
max_retries: 2
backoff:
initial_delay: 4
max_delay: 45
multiplier: 2
- get_first_workflow_output:
call: http.get
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/>"+execution_id}'
result: firstWorkflowOutput
- log_first_output:
call: sys.log
args:
text: '${firstWorkflowOutput.body.result}'
severity: INFO
- launch_second_workflow_execution:
call: http.post
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/projects/"+project_id+"/locations/"+location+"/workflows/"+secondWorkflowId+"/executions>"}'
body:
argument: '${firstWorkflowOutput.body.result}'
result: secondWorkflowResult
- set_second_execution_id:
assign:
- execution_id: ${secondWorkflowResult.body.name}
- exponential_backoff_retries_for_second_workflow:
try:
steps:
- get_second_execution_state:
call: http.get
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/>"+execution_id}'
result: secondWorkflowState
- set_second_execution_state:
assign:
- execution_state: ${secondWorkflowState.body.state}
- log_state_second_workflow:
call: sys.log
args:
text: '${secondWorkflowId +" state: "+ execution_state}'
severity: INFO
- induce_backoff_retry_if_second_state_not_done:
switch:
- condition: ${execution_state != "SUCCEEDED"}
raise: ${execution_state}
retry:
predicate: ${execution_state_predicate}
max_retries: 2
backoff:
initial_delay: 4
max_delay: 45
multiplier: 2
- get_second_workflow_output:
call: http.get
args:
auth:
type: OAuth2
url: '${"<https://workflowexecutions.googleapis.com/v1/>"+execution_id}'
result: secondWorkflowOutput
- log_second_output:
call: sys.log
args:
text: '${secondWorkflowOutput.body.result}'
severity: INFO
- returnSuccess:
return: "All pipelines finished with success! Check logs for more details."
execution_state_predicate: # Subworkflow to check if execution is complete
params: [execution_state]
steps:
- init:
assign:
- failureStates: ["FAILED","CANCELLED"]
- condition_to_retry:
switch:
- condition: ${execution_state in failureStates}
return: False # Stop Workflow
- condition: ${execution_state != "SUCCEEDED"}
return: True # Continue waiting
- otherwise:
return: False # Stop Workflow
Deverá ver um resultado semelhante ao seguinte:
Abrindo o separador de logs podemos ver o conjunto de caracteres transformados.
O processo é simples: definir variáveis, lançar o primeiro workflow e monitorizar o seu estado com tentativas de retrocesso exponenciais. Isto é rápido para workflows simples (1 a 2 minutos), mas os mais complexos podem necessitar de tempos de repetição mais longos (4 a 5 minutos). Isto destaca o tratamento de dependências, onde o segundo workflow depende do output do primeiro.
Cloud Workflows simplificam a construção de pipelines inteligentes, destacando-se na gestão de dependências para orquestrações complexas.
Bónus: Orquestrador simplificado
Aqui está um orquestrador de workflow mais simples para casos em que a partilha de variáveis não é necessária, apenas o rastreio de dependências.
main:
params: [args]
steps:
- init:
assign:
- project_id: ${sys.get_env("GOOGLE_CLOUD_PROJECT_ID")}
- workflowIds: ["sample-workflow-1", "sample-workflow-2"]
- location: "europe-west1"
- failureStates: ["FAILED","CANCELLED"]
- executeWorkflows:
for:
value: workflowId
in: ${workflowIds}
steps:
- launchWorkflow:
call: http.post
args:
auth:
type: OAuth2
url: ${"<https://workflowexecutions.googleapis.com/v1/projects/>" + project_id + "/locations/" + location + "/workflows/" + workflowId + "/executions"}
result: executionResult
- trackExecution:
assign:
- execution_id: ${executionResult.body.name}
- waitForCompletion:
try:
steps:
- getExecutionState:
call: http.get
args:
auth:
type: OAuth2
url: ${"<https://workflowexecutions.googleapis.com/v1/>" + execution_id}
result: executionState
- set_second_execution_state:
assign:
- execution_state: ${executionState.body.state}
- logState:
call: sys.log
args:
text: ${workflowId + " state -> " + execution_state}
severity: INFO
- checkCompletion:
switch:
- condition: ${execution_state != "SUCCEEDED"}
raise: ${execution_state} # Retry if not in a failure state and not yet succeeded.
retry:
predicate: ${execution_state_predicate}
max_retries: 2
backoff:
initial_delay: 4
max_delay: 45
multiplier: 2
- returnSuccess:
return: "All pipelines finished with success! Check logs for more details."
execution_state_predicate: # Subworkflow to check if execution is complete
params: [execution_state]
steps:
- init:
assign:
- failureStates: ["FAILED","CANCELLED"]
- condition_to_retry:
switch:
- condition: ${execution_state in failureStates}
return: False # Stop Workflow
- condition: ${execution_state != "SUCCEEDED"}
return: True # Continue waiting
- otherwise:
return: False # Stop Workflow
Esta estrutura é semelhante ao exemplo anterior, mas utiliza um ciclo for para adicionar facilmente mais workflows à lista de workflows. Sinta-se à vontade para experimentar, mas repare que este orquestrador simplificado não é compatível com os workflows anteriores devido à sua dependência de partilha de dados, uma característica que não está incluída aqui.
Conclusão
Esta publicação descreve estratégias para uma orquestração simplificada da Cloud Workflows. Utilizando o Google Cloud APIs, é possível criar pipelines robustos para gerir workflows mais complexos. Os exemplos de código fornecidos servem como ponto de partida para os seus próprios projetos. Lembre-se, dividir os workloads em vários workflows optimiza os custos e aumenta a flexibilidade. Embora o orquestrador acrescente etapas, os benefícios de uma pipeline mais inteligente superam o pequeno aumento no tempo de execução.

