Em um fluxo de trabalho de processamento de dados, use um nó for-each para executar a mesma subtarefa em cada item de uma lista, como nomes de arquivos ou partições. Esse nó itera automaticamente sobre o conjunto de resultados de um nó upstream (geralmente um nó de atribuição) e repete seu corpo de loop interno para cada elemento. Essa abordagem automatiza e otimiza seus fluxos de trabalho ao eliminar a criação manual e repetitiva de tarefas individuais.
Casos de uso
No desenvolvimento diário de dados, o nó for-each permite a execução parametrizada quando você precisa aplicar a mesma lógica de análise ou processamento a diferentes unidades de negócios, linhas de produtos ou itens de configuração. Por exemplo, se sua empresa possui várias linhas de produtos e você precisa gerar um relatório diário separado para cada uma, a lógica de processamento é idêntica; apenas os dados de destino diferem.
Assim como um loop for em linguagens de programação, o nó for-each itera automaticamente sobre uma lista, como nomes de tabelas, partições ou arquivos. Ele executa um subfluxo de trabalho predefinido para cada item da lista, o que aumenta significativamente a automação e a flexibilidade do seu fluxo de trabalho.
Pré-requisitos
Requisitos de versão: Este recurso está disponível apenas no DataWorks Standard Edition ou superior.
Permissões: Adicione sua conta RAM ao workspace de destino e atribua a ela a função Development ou Workspace Manager. Para mais informações, consulte Adicionar membros a um workspace.
Funcionamento
O nó for-each atua como um contêiner que encapsula um subfluxo de trabalho personalizável, conhecido como corpo do loop. O mecanismo funciona da seguinte maneira:
Entrada de dados: O nó for-each depende de um nó de atribuição upstream ou outro nó atribuível (como um nó EMR Hive). Ele recupera o conjunto de resultados formatado como array ao vincular-se ao parâmetro
loopDataArray.-
Execução do loop: Ao iniciar, o nó percorre sequencialmente cada elemento do conjunto de resultados. Para cada elemento, ele executa integralmente o corpo do loop interno uma vez, desde o nó
Startaté o nóEnd.NotaOs nós Start e End não são editáveis. Eles servem apenas para marcar o início e o fim do corpo do loop.
Passagem de dados: Durante cada iteração, o valor do elemento atual é transmitido aos nós internos do corpo do loop por meio de variáveis integradas. Os nós de negócio internos usam
${dag.foreach.current}para acessar o item de dados em processamento.
Parâmetros integrados
Variáveis no formato ${...} constituem uma sintaxe de modelo específica do DataWorks. O sistema analisa esses parâmetros diretamente e os substitui pelos respectivos valores antes da execução.
Os nós dentro do corpo do loop for-each podem usar as seguintes variáveis integradas para acessar o status do loop e os dados:
|
Parâmetro integrado |
Descrição |
Analogia com loop for |
|
|
O conjunto de resultados completo transmitido pelo nó de atribuição upstream. |
Considere o seguinte código de loop for:
|
|
|
O item de dados processado na iteração atual. |
|
|
|
O deslocamento atual do loop (indexado em 0). |
|
|
|
A contagem atual do loop (indexada em 1). |
Se a saída upstream for um array bidimensional, como o resultado de uma consulta SQL, use também a seguinte sintaxe para acessar valores específicos:
|
Outros parâmetros |
Descrição |
|
|
Retorna uma string separando os elementos da linha de dados atual (um array unidimensional) por vírgula |
|
|
O |
|
|
Os dados da Atualmente, o nó for-each não suporta loops aninhados. Este exemplo serve apenas para demonstrar a recuperação de valores. |
Limitações
Mecanismo de execução: O loop suporta execução serial e execução paralela. Escolha a execução paralela quando as iterações forem independentes entre si.
Limite de loops: O número máximo padrão de loops é 128, ajustável até 1024.
Restrições de depuração: Não execute um nó for-each diretamente no Data Studio. Implante a tarefa e teste-a no Operation Center usando o recurso de teste de fumaça.
Restrições de execução: Um nó for-each não pode ser executado isoladamente. Isso inclui testes de fumaça, backfill e execuções manuais.
Controle de fluxo no corpo do loop: Ao usar um nó de ramificação dentro do corpo de um loop for-each, garanta que todas as ramificações convirjam para um único nó de mesclagem antes de se conectarem ao nó
End. Isso assegura a integridade lógica do corpo do loop.Restrições de reexecução: Após a implantação de um nó, a reexecução automática em caso de falha retoma a partir do ponto de erro. No entanto, uma reexecução manual reinicia todo o nó for-each do zero.
Procedimento
Este procedimento usa um nó de atribuição como nó upstream e um nó Shell dentro do corpo do loop para imprimir os resultados. Esta seção orienta você na configuração de uma tarefa for-each completa:
-
Prepare os dados upstream (configure um nó de atribuição)
Crie e configure um nó de atribuição para fornecer um conjunto de resultados iterável para o nó for-each downstream.
No fluxo de trabalho, crie um nó de atribuição (por exemplo,
assign) e posicione-o upstream do nó for-each.-
Clique duas vezes no nó de atribuição e selecione um ambiente Python 2. Por exemplo, use
Python 2para gerar um array com quatro elementos:O nó envia [10,20,30,40] para os nós downstream dividindo automaticamente a última linha de saída em um array a cada vírgula.
print "10,20,30,40" O nó de atribuição gera automaticamente um parâmetro de saída chamado
outputs, que representa seu conjunto de resultados.Salve o nó de atribuição.
-
Configure o nó for-each para consumir dados
Configure o nó for-each para receber os dados upstream e usá-los dentro de seu corpo de loop.
Clique duas vezes no nó for-each para abrir sua tela interna.
-
No painel Scheduling à direita, localize o parâmetro
loopDataArrayem Scheduling Parameters e clique em Bind.Selecione o parâmetro outputs do nó assign para criar o vínculo. Após a conclusão, o valor do parâmetro loopDataArray refletirá seu status vinculado.
Na caixa de diálogo exibida, defina a Value Source como o nó de atribuição upstream (
assign) e selecione seu parâmetrooutputs. Essa ação cria automaticamente uma dependência entre os dois nós.-
No corpo do loop for-each, clique em Create Internal Node e crie um nó
Shell.Em um cenário real, você pode configurar qualquer tipo de nó.
-
Clique duas vezes no novo nó Shell e use variáveis integradas no código para recuperar e imprimir informações sobre o loop:
#!/bin/bash # Use ${dag.loopTimes} to get the current loop count echo "Current loop number is: ${dag.loopTimes}" # Use ${dag.foreach.current} to get the data item for the current iteration echo "Current item is: ${dag.foreach.current}" -
(Opcional) No painel Scheduling Settings à direita, configure as propriedades em Scheduling Policy.
-
Maximum Number of Loops: O padrão é 128 e o máximo é 1024.
ImportanteEste parâmetro determina o número máximo de iterações do corpo do loop. Se a quantidade de itens de dados upstream for grande, aumente esse valor para garantir o processamento de todos os itens.
-
Execute Policy: Selecione Serial para este exemplo.
Serial: Executa as iterações sequencialmente.
Parallel: Executa as iterações do loop simultaneamente para melhorar a eficiência da tarefa. No modo Parallel, a falha de uma iteração não afeta as demais. O agendador tenta executar todas as iterações até a conclusão. A concorrência padrão é 5, e o máximo é 20.
-
Salve o nó Shell.
-
Implante, execute e verifique
Implante o fluxo de trabalho no Operation Center para execução e verifique os resultados do nó for-each.
Retorne à tela principal do fluxo de trabalho e clique em Deploy na barra de ferramentas para publicar todo o fluxo.
-
Acesse e faça um teste de fumaça no fluxo de trabalho de destino.
ImportanteNão faça um teste de fumaça no nó for-each individualmente. Como o nó for-each depende da saída do nó de atribuição upstream, inicie o teste a partir do nó de atribuição para garantir que a linhagem de dados esteja completa.
Após a execução bem-sucedida da instância de teste, localize a instância do nó for-each na lista, abra-a e clique com o botão direito para selecionar View Internal Nodes.
-
Na visualização de nós internos, verifique as instâncias do nó Shell geradas por cada loop. Abra o log de execução de qualquer instância para visualizar a saída daquela iteração e confirmar se o resultado está correto.
O painel à esquerda mostra que todas as quatro iterações do loop foram concluídas. O log de execução da quarta iteração exibe
Current loop number is: 4eCurrent item is: 40, e o comando Shell encerra com o código 0, indicando execução bem-sucedida.
Caso de uso: Processar diferentes formatos de dados
Cenário 1: Processar um array unidimensional
Saída do nó de atribuição: 2025-11-01,2025-11-02,2025-11-03
Contagem de iterações: 3
-
Durante a segunda iteração:
O valor de
${dag.foreach.current}é2025-11-02.O valor de
${dag.loopTimes}é2.
Cenário 2: Processar um array bidimensional
-
Saída do nó de atribuição (MaxCompute SQL):
+-----+----------+ | id | city | +-----+----------+ | 101 | beijing | | 102 | shanghai | +-----+----------+ Contagem de iterações: 2
-
Durante a segunda iteração:
O valor de
${dag.foreach.current}é102,shanghai.O valor de
${dag.loopTimes}é2.O valor de
${dag.foreach.current[0]}é102.O valor de
${dag.foreach.current[1]}éshanghai.
Caso de uso: Processamento em lote de dados particionados
Este exemplo demonstra como usar um nó de atribuição e um nó for-each para processar em lote dados de comportamento do usuário provenientes de múltiplas linhas de negócio. Isso permite implementar o processamento automatizado de dados para várias linhas de produtos com um único conjunto de lógicas.
Contexto
Suponha que você seja um engenheiro de dados em uma grande empresa de internet, responsável por processar dados de três linhas de negócio principais: e-commerce (ecom), finanças (finance) e logística (logistics). Novas linhas podem ser adicionadas no futuro. Você precisa executar a mesma lógica diária de agregação nos logs de comportamento do usuário dessas linhas para calcular o PV (page view) de cada usuário e armazenar os resultados em uma tabela de resumo unificada.
-
Tabelas de origem upstream (camada DWD):
dwd_user_behavior_ecom_d: Tabela de comportamento do usuário de e-commerce.dwd_user_behavior_finance_d: Tabela de comportamento do usuário de finanças.dwd_user_behavior_logistics_d: Tabela de comportamento do usuário de logística.dwd_user_behavior_${biz_line}_d: Tabelas de comportamento do usuário para outras possíveis linhas de negócio.Essas tabelas possuem a mesma estrutura e são particionadas por dia (
dt).
-
Tabela de destino downstream (camada DWS):
dws_user_summary_d: Tabela de resumo do usuário.Esta tabela é particionada por linha de negócio (
biz_line) e dia (dt) para armazenar os resultados agregados de todas as linhas de negócio.
Criar uma tarefa separada para cada linha de negócio dificulta a manutenção e aumenta a propensão a erros. Usar um nó for-each simplifica a manutenção, pois basta manter um único conjunto de lógicas de processamento, enquanto o sistema itera automaticamente por todas as linhas de negócio para concluir os cálculos.
Preparação dos dados
Primeiramente, crie as tabelas de exemplo e insira dados de teste. Este exemplo usa a data de negócio 20251010.
Vincule um mecanismo de computação MaxCompute ao workspace.
Acesse o Data Studio e crie um nó MaxCompute SQL.
-
Crie as tabelas de origem (camada DWD): Adicione o seguinte código ao nó MaxCompute SQL, selecione-o e execute-o.
-- E-commerce user behavior table CREATE TABLE IF NOT EXISTS dwd_user_behavior_ecom_d ( user_id STRING COMMENT 'User ID', action_type STRING COMMENT 'Action type', event_time BIGINT COMMENT 'UNIX timestamp of the event in milliseconds' ) COMMENT 'Detailed e-commerce user behavior log table' PARTITIONED BY (dt STRING COMMENT 'Date partition, format yyyymmdd'); INSERT OVERWRITE TABLE dwd_user_behavior_ecom_d PARTITION (dt='20251010') VALUES ('user001', 'click', 1760004060000), -- 2025-10-10 10:01:00.000 ('user002', 'browse', 1760004150000), -- 2025-10-10 10:02:30.000 ('user001', 'add_to_cart', 1760004300000); -- 2025-10-10 10:05:00.000 -- Verify that the e-commerce user behavior table is created. SELECT * FROM dwd_user_behavior_ecom_d where dt='20251010'; -- Finance user behavior table CREATE TABLE IF NOT EXISTS dwd_user_behavior_finance_d ( user_id STRING COMMENT 'User ID', action_type STRING COMMENT 'Action type', event_time BIGINT COMMENT 'UNIX timestamp of the event in milliseconds' ) COMMENT 'Detailed finance user behavior log table' PARTITIONED BY (dt STRING COMMENT 'Date partition, format yyyymmdd'); INSERT OVERWRITE TABLE dwd_user_behavior_finance_d PARTITION (dt='20251010') VALUES ('user003', 'open_app', 1760020200000), -- 2025-10-10 14:30:00.000 ('user003', 'transfer', 1760020215000), -- 2025-10-10 14:30:15.000 ('user003', 'check_balance', 1760020245000), -- 2025-10-10 14:30:45.000 ('user004', 'open_app', 1760020300000); -- 2025-10-10 14:31:40.000 -- Verify that the finance user behavior table is created. SELECT * FROM dwd_user_behavior_finance_d where dt='20251010'; -- Logistics user behavior table CREATE TABLE IF NOT EXISTS dwd_user_behavior_logistics_d ( user_id STRING COMMENT 'User ID', action_type STRING COMMENT 'Action type', event_time BIGINT COMMENT 'UNIX timestamp of the event in milliseconds' ) COMMENT 'Detailed logistics user behavior log table' PARTITIONED BY (dt STRING COMMENT 'Date partition, format yyyymmdd'); INSERT OVERWRITE TABLE dwd_user_behavior_logistics_d PARTITION (dt='20251010') VALUES ('user001', 'check_status', 1760032800000), -- 2025-10-10 18:00:00.000 ('user005', 'schedule_pickup', 1760032920000); -- 2025-10-10 18:02:00.000 -- Verify that the logistics user behavior table is created. SELECT * FROM dwd_user_behavior_logistics_d where dt='20251010'; -
Crie a tabela de destino (camada DWS): Adicione o seguinte código ao nó MaxCompute SQL, selecione-o e execute-o.
CREATE TABLE IF NOT EXISTS dws_user_summary_d ( user_id STRING COMMENT 'User ID', pv BIGINT COMMENT 'PV' ) COMMENT 'User daily PV summary table' PARTITIONED BY ( dt STRING COMMENT 'Date partition, format yyyymmdd', biz_line STRING COMMENT 'Business line partition, such as ecom, finance, logistics' );ImportanteSe o workspace estiver na Standard Edition, implante este nó no ambiente de produção e faça um backfill de dados.
Implementação do fluxo de trabalho
Crie um fluxo de trabalho. No painel Scheduling Parameters à direita, defina o parâmetro de agendamento bizdate como o dia anterior:
$[yyyymmdd-1].-
No fluxo de trabalho, crie um nó de atribuição chamado get_biz_list e use a linguagem MaxCompute SQL para escrever o seguinte código. Este nó gera a lista de linhas de negócio a serem processadas:
-- Output all business lines to be processed SELECT 'ecom' AS biz_line UNION ALL SELECT 'finance' AS biz_line UNION ALL SELECT 'logistics' AS biz_line; -
Configure o nó for-each
Retorne à tela do fluxo de trabalho e crie um nó for-each downstream para o nó de atribuição get_biz_list.
Acesse a página de configurações do nó for-each. No painel de configuração de agendamento à direita, em , vincule o parâmetro loopDataArray às saídas (outputs) do nó get_biz_list.
-
No corpo do loop do nó for-each, clique em Create Internal Node para criar um nó MaxCompute SQL e escreva a lógica de processamento do corpo do loop.
NotaEste script é impulsionado pelo nó for-each e executa uma vez para cada linha de negócio.
Durante a execução, a variável integrada ${dag.foreach.current} é substituída dinamicamente pelo nome da linha de negócio atual. Os valores de iteração esperados são 'ecom', 'finance' e 'logistics'.
SET odps.sql.allow.dynamic.partition=true; INSERT OVERWRITE TABLE dws_user_summary_d PARTITION (dt='${bizdate}', biz_line) SELECT user_id, COUNT(*) AS pv, '${dag.foreach.current}' AS biz_line FROM dwd_user_behavior_${dag.foreach.current}_d WHERE dt = '${bizdate}' GROUP BY user_id;
-
Adicione um nó de verificação
Retorne ao fluxo de trabalho. Clique em Create Downstream Node para o nó for-each a fim de criar um nó MaxCompute SQL e adicione o seguinte código.
SELECT * FROM dws_user_summary_d WHERE dt='20251010' ORDER BY biz_line, user_id;
Implantação e resultados
Implante o fluxo de trabalho no ambiente de produção. Acesse , localize o fluxo de trabalho de destino e execute um teste de fumaça com a data de negócio definida como '20251010'.
Após a conclusão da execução, visualize o log de execução na instância de teste. O nó final deve produzir a seguinte saída:
|
user_id |
pv |
dt |
biz_line |
|
user001 |
2 |
20251010 |
ecom |
|
user002 |
1 |
20251010 |
ecom |
|
user003 |
3 |
20251010 |
finance |
|
user004 |
1 |
20251010 |
finance |
|
user001 |
1 |
20251010 |
logistics |
|
user005 |
1 |
20251010 |
logistics |
Vantagens
Alta escalabilidade: Para adicionar uma nova linha de negócio, basta inserir uma linha de SQL no nó de atribuição. A lógica de processamento não requer alterações.
Facilidade de manutenção: Todas as linhas de negócio compartilham a mesma lógica de processamento. Uma alteração em um único local aplica-se a todas.
Perguntas frequentes
-
P: Por que não consigo executar um nó for-each diretamente no Data Studio para testá-lo?
R: Isso ocorre por design. O nó requer um ambiente de agendamento completo para resolver o contexto do nó e suas dependências, portanto, não suporta execução direta no Data Studio. Implante a tarefa no Operation Center e teste-a usando backfill ou acionando uma execução agendada.
-
P: Por que um teste de fumaça em um nó for-each individual falha ou não faz nada?
R: Os dados do loop de um nó for-each provêm de seu parâmetro de entrada
loopDataArray, que deve estar vinculado ao parâmetrooutputsde um nó de atribuição upstream. Se você executar o nó for-each isoladamente, ele falhará ou será ignorado porque não conseguirá receber um conjunto de resultados de entrada. -
P: Por que meu loop executa apenas uma vez?
R: Geralmente, isso acontece porque a saída do nó de atribuição upstream é analisada como um único elemento. Verifique sua saída:
É uma única string sem delimitador?
Se você espera iterar sobre vários itens, garanta que estejam separados por vírgulas (
,). Por exemplo,'item1,item2,item3'resulta em três loops, enquanto'item1 item2 item3'resulta em apenas um.