Diferentemente de um workflow agendado, que executa conforme um cronograma predefinido (como diariamente à 1h), um workflow acionado é um modelo de processamento de dados sob demanda e orientado a eventos. Um sinal externo — como upload de arquivo, chegada de mensagem, chamada de API ou clique manual — dispara sua execução em tempo real, oferecendo alta capacidade de resposta e flexibilidade no processamento de dados.
|
Recurso |
Workflow agendado |
Workflow acionado |
|
Mecanismo de acionamento |
Cronograma fixo (expressão cron) |
Sinal externo (evento, API, manual) |
|
Modelo de execução |
Agendado e previsível |
Reativo e sob demanda |
|
Casos de uso |
Data warehousing em lote T+1, relatórios agendados |
Processar arquivos ao chegar, integrar com sistemas de negócios, reparo manual de dados |
|
Principais vantagens |
Confiabilidade e agendamento previsível |
Capacidade de resposta em tempo real e flexibilidade |
Métodos de acionamento compatíveis
Um workflow acionado oferece três métodos de acionamento. Escolha aquele adequado ao seu cenário.
|
Método de acionamento |
Iniciador |
Cenários principais |
Pontos-chave |
|
Acionamento por evento |
Fonte de evento externa (como OSS ou ApsaraMQ for Kafka) |
ETL orientado a eventos: processe arquivos assim que chegarem ou dispare computação em tempo real a partir de mensagens. |
Crie primeiro um trigger e associe-o ao workflow. Válido apenas no ambiente de produção. |
|
Acionamento manual |
Usuário (desenvolvedor/operador) |
Tarefas ad-hoc: processamento ou análise pontual de dados. |
Executa tanto nos ambientes de desenvolvimento quanto de produção. Recomendado como alternativa a um processo de negócios manual. |
|
Acionamento por API |
Sistema externo (via OpenAPI) |
Integração de sistemas: dispare o processamento de dados por meio de callback de um sistema de negócios, como CRM ou ERP. |
Requer chamada a uma OpenAPI com as permissões necessárias. |
Início rápido: crie um workflow acionado manualmente
Esta seção orienta você na criação de um workflow acionado simples e na sua execução manual para experimentar rapidamente todo o processo.
Etapa 1: crie um workflow acionado
Acesse a página Workspaces no console do DataWorks. Na barra de navegação superior, selecione a região desejada. Localize o workspace e escolha na coluna Actions.
Na barra de navegação à esquerda, clique em
. Em seguida, ao lado de Project Directory, clique em para abrir a página Create Workflow.Na caixa de diálogo Create Workflow, defina Scheduling Type como Triggered scheduling, insira um Name para o workflow e clique em Confirm.
Etapa 2: orquestre e desenvolva nós
Na barra de ferramentas, clique em + Add Node para abrir a lista de nós. Na lista de tipos de nó à esquerda, arraste um nó Shell para a tela e insira um nome para criá-lo.
-
Clique duas vezes no nó Shell para abrir o editor de código e insira o seguinte código:
echo "Hello, Trigger Workflow! Current time is ${bizdate}" Clique no botão Save na barra de ferramentas.
Etapa 3: depure e execute
Retorne à tela do workflow e clique no ícone
na barra de ferramentas superior.Na caixa de diálogo exibida, insira o Value Used in This Run para esta execução (por exemplo, se hoje for 20260310, substitua
bizdatepor20260309).Após alguns instantes, visualize o status de execução do nó e a saída do comando
echono log de execução abaixo.
Etapa 4: publique e execute
Acima da tela do workflow, clique no botão Publish
e siga as instruções para publicá-lo.Após publicar o workflow, acesse Operation Center > Manually Triggered Task O&M > Manually Triggered Task> Triggered Workflow.
Localize o workflow recém-publicado e clique em Run na coluna Operation.
Na janela exibida, clique em Run novamente para disparar uma instância do workflow no ambiente de produção. A página Triggered Workflow Instance exibe os detalhes desta execução.
Agora você conhece os conceitos básicos sobre workflows acionados. A seguir, exploraremos os recursos mais avançados de acionamento por eventos.
Casos de uso avançados: workflows acionados por eventos
Cenário 1: processe novos arquivos do OSS
Objetivo: quando um novo arquivo CSV for carregado em um diretório especificado no Object Storage Service (OSS), acione automaticamente um workflow que imprime o caminho do arquivo.
Etapa 1: crie um trigger do OSS
Acesse Operation Center > Scheduling Settings > Trigger Management.
-
Clique em Create Trigger e configure-o da seguinte forma:
NotaPara descrições detalhadas dos parâmetros, consulte OSS trigger.
Trigger Name: insira um nome personalizado, como
oss_new_file_trigger.Applicable Workspace: selecione o workspace de destino onde o workflow está localizado.
Trigger Event Type: selecione Object Storage OSS.
Trigger Event: selecione
oss:ObjectCreated:PutObject(ou outro evento de upload).Bucket Name: selecione seu bucket do OSS.
File Name: especifique o caminho e o formato do arquivo a monitorar. Há suporte para wildcards. Por exemplo, para monitorar todos os arquivos
.csvno diretórioinput/, insirainput/*.csv.-
Role Configuration: no primeiro uso, realize a autorização com um clique. Selecione a função gerada chamada DataWorks-EventBridge-OSS-MNS-Role-***.
O *** representa um ID de 13 dígitos gerado aleatoriamente para garantir unicidade.
Clique em Confirm para criar o trigger.
Etapa 2: crie e associe o workflow
Siga as etapas em Início rápido: crie um workflow acionado manualmente para criar um novo workflow acionado chamado
process_oss_file_workflow.No painel à direita da tela do workflow, selecione Scheduling Settings > Scheduling Policy.
-
Na lista suspensa Trigger, selecione o
oss_new_file_triggercriado anteriormente.Após selecionar um trigger, use a variável
${workflow.triggerMessage}nas tarefas internas para obter o corpo completo da mensagem ou utilize${workflow.triggerMessage.xxx}para obter um campo específico.
Etapa 3: analise os parâmetros do evento
Na barra de ferramentas, clique em + Add Node para abrir a lista de nós. Na lista de tipos de nó à esquerda, arraste um nó Shell para a tela e insira um nome para criá-lo.
-
Clique duas vezes no nó e escreva o código para obter e imprimir o caminho do arquivo a partir do evento de acionamento.
# When a trigger starts a workflow, the event information is passed through the built-in variable workflow.triggerMessage # We can get the full path of the uploaded file using ${workflow.triggerMessage.data.oss.object.key} echo "========= Start Processing OSS File =========" message="${workflow.triggerMessage}" echo "Raw Value: ${message}" # Extract the file name from the event message FILE_PATH="${workflow.triggerMessage.data.oss.object.key}" echo "A new file has arrived: ${FILE_PATH}" # You can write specific processing logic here echo "========= Finish Processing OSS File ========="Nota${workflow.triggerMessage}: obtém o corpo completo da mensagem de evento no formato JSON. Encontre o formato específico da mensagem para OSS no EventBridge navegando até Event Bus >
DATAWORKS_TRIGGER_FOR_BUCKET_<OSS_Bucket_Name>> Event Tracking> Event Detail.
Etapa 4: depure e publique
-
Depuração:
Retorne à tela do workflow e clique no botão Run
.-
Na caixa de entrada Trigger Message Body, cole um JSON simulando um evento do OSS. Copie o exemplo de formato de mensagem da página de configuração do trigger e modifique o valor de
key. Veja um exemplo mínimo abaixo.{ "data": { "oss":{ "object": { "key": "input/test_file_20260310.csv" } } } } Clique em Run e verifique os logs para confirmar se
input/test_file_20260310.csvfoi impresso com sucesso.
Publicação: após depurar com sucesso, clique no botão Publish para implantar o workflow no ambiente de produção. O acionamento por eventos só entra em vigor no ambiente de produção.
Etapa 5: verifique em produção
-
Use o console do OSS ou um cliente para carregar um arquivo CSV no bucket e no caminho configurados no trigger (por exemplo, o diretório
input/). Acesse o Operation Center do DataWorks > Manually Triggered Task O&M > Manually Triggered Task> Triggered Workflow. O workflow publicado
process_oss_file_workflowaparecerá na lista.-
Após uma breve espera, acesse o Operation Center do DataWorks > Manually Triggered Task O&M > Triggered Workflow Instance. Uma nova instância do workflow será acionada automaticamente. Clique para visualizar seu log e confirmar se o caminho do arquivo foi processado corretamente.
Melhor prática: design de idempotência
Eventos do OSS podem ser entregues mais de uma vez devido a fatores como flutuações de rede. Para evitar processamento duplicado, implemente idempotência na lógica de negócios. Uma abordagem comum é verificar uma tabela de registros (como uma tabela do MaxCompute) antes de processar um arquivo, usando o ETag ou o caminho exclusivo do arquivo como identificador. Se já tiver sido processado, ignore-o.
Cenário 2: processe novas mensagens do Kafka
Objetivo: monitore um tópico do Kafka em busca de logs de comportamento do usuário. Quando uma nova mensagem chegar, acione um workflow para analisá-la e executar lógicas diferentes com base no conteúdo.
Etapa 1: crie um trigger do Kafka
Acesse Operation Center > Scheduling Settings > Trigger Management e clique em Create Trigger.
-
Configure os seguintes itens:
Trigger Name:
kafka_user_action_trigger.Trigger Event Type: selecione ApsaraMQ for Kafka.
Kafka Instance, Topic: selecione a instância e o tópico que deseja monitorar.
ConsumerGroupId: recomendamos selecionar Quick Create. O sistema gera automaticamente um ID de grupo de consumidores para evitar conflitos com outras aplicações.
Key (opcional): especifique uma chave de mensagem. Apenas mensagens com correspondência exata de chave acionarão o workflow.
Clique em Confirm.
Etapa 2: crie e associe o workflow
Siga as etapas em Início rápido: crie um workflow acionado manualmente para criar um novo workflow acionado chamado
handle_user_action_workflow.No painel à direita da tela do workflow, selecione Scheduling Settings > Scheduling Policy.
-
Na lista suspensa Trigger, selecione o
kafka_user_action_triggercriado anteriormente.Após selecionar um trigger, use a variável
${workflow.triggerMessage}nas tarefas internas para obter o corpo completo da mensagem ou utilize${workflow.triggerMessage.xxx}para obter um campo específico. (Importante) Como as mensagens podem chegar com alta frequência, defina com um valor como
100para evitar que picos de mensagens sobrecarreguem os recursos de agendamento.
Etapa 3: analise JSON aninhado
Suponha que o campo value da mensagem do Kafka seja uma string JSON com o seguinte formato: {"user_id": "1001", "action_type": "login", "timestamp": 1688888888}.
Na barra de ferramentas, clique em + Add Node para abrir a lista de nós. Na lista de tipos de nó à esquerda, arraste um nó Python para a tela.
-
Escreva o código para analisar a mensagem. Como
valueé uma string, faça o parse da string JSON no código.import json # 1. Use the built-in variable to get the value field of the Kafka message. It is a JSON string. message_value_str = '${workflow.triggerMessage.value}' print(f'Received raw message value string: ${message_value_str}') try: # 2. Parse the string into a JSON object (dictionary) in Python. message_data = json.loads(message_value_str) user_id = message_data.get("user_id") action_type = message_data.get("action_type") print(f"Successfully parsed message. User ID: ${user_id}, Action: ${action_type}") # 3. Execute different business logic based on action_type. if action_type == 'login': # o.run_sql(f"INSERT OVERWRITE TABLE user_login_record PARTITION(ds='{bizdate}') VALUES ('{user_id}');") print("Processing login action...") elif action_type == 'purchase': print("Processing purchase action...") else: print("Unknown action type.") except json.JSONDecodeError as e: print(f"Error decoding JSON: {e}") # Error handling. For example, write the error message to a dedicated log table. raise e # Re-raise so the node fails and is easier to troubleshoot.
Etapa 4: depure e publique
-
Depuração:
Retorne à tela do workflow e clique no botão Run
.-
No Trigger Message Body, cole um evento simulado do Kafka. Observe que o campo
valueé uma string JSON escapada.{ "topic": "user-behavior-topic", "key": "some-key", "value": "{\"user_id\": \"1001\", \"action_type\": \"login\", \"timestamp\": 1688888888}" } Execute e verifique os logs para confirmar se o nó Python analisou corretamente
user_ideaction_type.
Publicação: após depurar com sucesso, publique o workflow no ambiente de produção.
Etapa 5: verifique em produção
-
Envie uma mensagem no formato correto para o tópico do Kafka configurado.
Na página de detalhes do tópico, clique em Quickly Experience Message Sending and Receiving e selecione o método de envio Console. Em Message Key, insira
some-key. Em Message Content, insira{"user_id": "1001", "action_type": "login", "timestamp": 1688888888}. Para Send to specified partition, selecione No e clique em Send. Um aviso de Message sent successfully indica que a verificação foi concluída. Acesse o Operation Center do DataWorks > Manually Triggered Task O&M > Manually Triggered Task> Triggered Workflow. O workflow publicado
handle_user_action_workflowaparecerá na lista.-
Em Operation Center > Manually Triggered Task O&M > Manual instance > Triggered Workflow Instance, observe se uma nova instância do workflow foi acionada e verifique seu log de execução.
2026-xxx 14:55:40 INFO ======================================================================== Received raw message value string: ${"user_id": "1001", "action_type": "login", "timestamp": 1688888888} Successfully parsed message. User ID: $1001, Action: $login Processing login action... 2026-xxx 14:55:40 INFO ========================================================================
Melhor prática: concorrência e ordenação
Controle de concorrência: sempre defina um número máximo razoável de instâncias paralelas para lidar com picos de mensagens.
Garantia de ordem: o agendamento do DataWorks não garante ordenação estrita de mensagens. Se for necessário assegurar que mensagens do mesmo usuário (ou partição) sejam executadas em ordem, implemente um bloqueio distribuído (por exemplo, baseado em Redis ou MaxCompute) no código de negócios ou delegue a lógica de processamento a um mecanismo de computação que garanta consumo ordenado por partição, como o Flink.
Design e configuração principais
Orquestração de workflow
O processo principal para orquestrar um workflow acionado é semelhante ao de um workflow agendado. Para mais informações, consulte Orchestrate nodes and workflows.
Parâmetros de agendamento
No painel Scheduling Settings no lado direito da tela do workflow, defina parâmetros globais para o workflow. Todos os nós dentro do workflow podem referenciar esses parâmetros.
Método de referência: no código do nó, referencie um parâmetro do workflow usando o formato
${workflow.parameter_name}.-
Prioridade de parâmetros: os parâmetros no DataWorks possuem uma relação hierárquica de substituição. A prioridade efetiva é: parâmetros de nó > parâmetros de workflow.
Para mais informações sobre parâmetros, consulte Parameter design and flow .
Políticas de agendamento
Quando múltiplos workflows ou tarefas são acionados simultaneamente e causam gargalo de recursos, utilize Priority e priority weighting policies para alcançar agendamento inteligente de recursos e garantir que as tarefas mais importantes executem primeiro.
Garantia de negócios essenciais: defina uma prioridade maior para workflows de negócios críticos para assegurar que sempre executem antes de outros workflows não essenciais.
-
Redução do tempo de processos críticos: dentro de uma única instância de workflow, influencie a ordem de execução dos nós usando a Priority weighting strategy. Por exemplo, usar a política de ponderação descendente atribui um peso dinâmico maior aos nós no caminho crítico que possuem mais dependências upstream. Isso faz com que executem primeiro, reduzindo o tempo total de execução do workflow.
Parâmetro
Descrição
Priority
Define o nível absoluto de prioridade de uma instância de workflow na fila de agendamento. Os níveis disponíveis são 1, 3, 5, 7 e 8 (quanto maior o número, maior a prioridade). Tarefas ou workflows de alta prioridade sempre obtêm recursos de agendamento antes dos de baixa prioridade.
Priority weighting strategy
Define como os pesos de nós individuais (tarefas) dentro de um workflow são calculados dinamicamente no mesmo nível de prioridade. Nós com pesos maiores recebem prioridade de execução.
-
Sem ponderação: todos os nós têm um peso base fixo.
-
Ponderação descendente: o peso de um nó é ajustado dinamicamente. Quanto mais nós upstream um nó depender, maior será seu peso. Esta política ajuda a priorizar nós no caminho crítico de um Grafo Acíclico Direcionado (DAG). O peso é calculado como:
Peso inicial + Soma das prioridades de todos os nós upstream.
Controla o número máximo de instâncias deste workflow que podem executar simultaneamente, usado para controle de concorrência e proteção de recursos. Quando o número de instâncias em execução atinge esse limite, novas instâncias acionadas entram em fila. Defina como Allowed ou um valor máximo personalizado (até 100.000).
NotaSe o limite definido exceder a capacidade do grupo de recursos, a concorrência real será limitada pelos limites físicos do grupo de recursos.
-
O sistema de prioridades do DataWorks segue uma regra hierárquica de substituição: especificação em tempo de execução > configuração no nível do nó > configuração no nível do workflow.
Configuração no nível do workflow (base): configurada na Scheduling Policy do workflow como configuração padrão para todos os nós.
Configuração no nível do nó (local): nas Scheduling Settings > Scheduling Policy de um único nó, defina uma Priority maior para um nó específico, substituindo a configuração no nível do workflow.
Especificação em tempo de execução (temporária): ao acionar manualmente uma execução no Operation and Maintenance Center, use a opção Runtime Priority Reset para especificar uma configuração temporária. Esta configuração tem a maior precedência e aplica-se apenas à execução atual, sem modificar configurações permanentes.
Operações e gerenciamento
Monitoramento de instâncias: na página Operation Center > Manually Triggered Node O&M > Triggered Workflow Instance, monitore todas as instâncias acionadas ou executadas manualmente. Isso inclui visualizar status, reexecutar ou terminar instâncias e verificar logs.
Clonar workflow: em Workspace Directories, clique com o botão direito em um workflow e selecione Clone para copiar rapidamente um novo workflow incluindo todos os nós e relacionamentos de dependência. Para mais informações, consulte Clone a workflow para workflows agendados.
Gerenciamento de versões: no painel Version no lado direito da tela do workflow, visualize, compare e reverta para versões históricas do workflow. Para mais informações, consulte Version management para workflows agendados.
Limitações e observações
Ambiente efetivo: o mecanismo de acionamento por evento só entra em vigor após a publicação do workflow no ambiente de produção (Operation Center).
Quantidade de nós: um único workflow suporta no máximo 400 nós. Para simplificar a manutenção, recomendamos manter a quantidade de nós abaixo de 100.
Limite de concorrência: o número máximo de instâncias paralelas é 100.000, mas a capacidade concorrente real é limitada pelas especificações do grupo de recursos de agendamento adquirido.
Agendamento no nível do nó: ao configurar o agendamento no nível do nó, há suporte apenas para Priority; a política de ponderação de prioridade não é compatível.
Tipos de nó não compatíveis: EMR Spark Streaming, Flink SQL Streaming, Flink JAR Streaming, Flink Python Streaming e nós de verificação de dependência não são compatíveis com workflows acionados. Eles só podem ser desenvolvidos e executados como nós independentes.