Un nœud for-each parcourt un jeu de résultats en amont, tel qu'une liste de noms de fichiers ou de partitions, et exécute la même sous-tâche pour chaque élément. Cette approche évite de créer manuellement des tâches individuelles et permet de mettre en place des workflows dynamiques et automatisés.
Cas d'usage
Le nœud for-each permet une exécution paramétrée lorsque vous devez appliquer la même logique d'analyse ou de traitement à différentes unités commerciales, lignes de produits ou éléments de configuration. Par exemple, si votre entreprise gère plusieurs lignes de produits et que vous devez générer un rapport quotidien distinct pour chacune d'elles, la logique de traitement reste identique ; seules les données cibles diffèrent.
À l'instar d'une boucle for dans un langage de programmation, le nœud for-each itère sur une liste (noms de tables, noms de partitions ou noms de fichiers) et exécute un sous-workflow prédéfini pour chaque élément.
Remarques d'utilisation
Exigences de version : Disponible uniquement à partir de DataWorks Standard Edition.
Permissions : Votre compte RAM doit être ajouté au workspace cible et se voir attribuer le rôle de developer ou de workspace administrator. Pour plus d'informations, consultez Ajouter des membres à un workspace.
Fonctionnement
Le nœud for-each agit comme un conteneur encapsulant un sous-workflow personnalisable, appelé corps de boucle. Son fonctionnement est le suivant :
Entrée des données : Le nœud for-each dépend d'un nœud d'affectation en amont ou d'un autre nœud assignable (tel qu'un nœud EMR Hive). Il récupère le jeu de résultats au format tableau en se liant au paramètre
loopDataArray.-
Exécution de la boucle : Au démarrage du nœud, celui-ci parcourt séquentiellement chaque élément du jeu de résultats. Pour chaque élément, il exécute intégralement le corps de boucle interne une fois, du nœud
Startau nœudEnd.RemarqueLes nœuds Start et End ne sont pas modifiables. Ils servent uniquement à délimiter le début et la fin du corps de boucle.
Transmission des données : Lors de chaque itération, la valeur de l'élément courant est transmise aux nœuds du corps de boucle via des variables intégrées. Les nœuds métier internes utilisent
${dag.foreach.current}pour accéder à l'élément de données en cours de traitement.
Paramètres intégrés
Les variables au format ${...} constituent une syntaxe de modèle propre à DataWorks. DataWorks analyse directement ces paramètres et les remplace par leurs valeurs avant l'exécution.
Les nœuds situés dans le corps de boucle for-each peuvent utiliser les variables intégrées suivantes pour accéder à l'état de la boucle et aux données :
|
Paramètre intégré |
Description |
Analogie avec une boucle for |
|
|
Le jeu de résultats complet transmis par le nœud d'affectation en amont. |
Prenons le code de boucle for suivant :
|
|
|
L'élément de données traité lors de l'itération courante. |
|
|
|
Le décalage actuel de la boucle (indexé à partir de 0). |
|
|
|
Le compteur actuel de la boucle (indexé à partir de 1). |
Si la sortie en amont est un tableau à deux dimensions, comme le résultat d'une requête SQL, vous pouvez également utiliser la syntaxe suivante pour accéder à des valeurs spécifiques :
|
Autres paramètres |
Description |
|
|
Renvoie une chaîne obtenue en séparant les éléments de la ligne de données courante (un tableau unidimensionnel) par une virgule |
|
|
Le |
|
|
Les données situées à la Le nœud for-each ne prend actuellement pas en charge les boucles imbriquées. Cet exemple sert uniquement à illustrer la récupération de valeurs. |
Limites
Mécanisme d'exécution : La boucle prend en charge l'exécution sérielle et l'exécution parallèle. Privilégiez le mode parallèle lorsque les itérations sont indépendantes les unes des autres.
Limite d'itérations : Le nombre maximal d'itérations par défaut est de 128, mais cette valeur peut être augmentée jusqu'à 1024.
Contraintes de débogage : L'exécution directe d'un nœud for-each dans Data Studio n'est pas possible. Vous devez déployer la tâche, puis la tester dans Operation Center à l'aide de la fonctionnalité de test fumée.
Contraintes d'exécution : Un nœud for-each ne peut pas être exécuté de manière isolée. Cela s'applique aux tests fumée, au backfill et aux exécutions manuelles.
Contrôle de flux dans le corps de boucle : Si vous utilisez un nœud de branchement dans le corps d'une boucle for-each, assurez-vous que toutes les branches convergent vers un seul nœud de fusion avant de se connecter au nœud
End. Cette convergence garantit l'intégrité logique du corps de boucle.Contraintes de réexécution : Après le déploiement d'un nœud, une réexécution automatique en cas d'échec reprend au point de défaillance. En revanche, une réexécution manuelle déclenche une reprise complète de l'ensemble du nœud for-each.
Procédure
Cette procédure utilise un nœud d'affectation comme nœud en amont et un nœud Shell dans le corps de boucle pour afficher les résultats. Elle détaille la configuration d'une tâche for-each complète :
-
Préparer les données en amont (configurer un nœud d'affectation)
Créez et configurez un nœud d'affectation afin de fournir un jeu de résultats itérable au nœud for-each en aval.
Dans le workflow, créez un nœud d'affectation (par exemple,
assign) et placez-le en amont du nœud for-each.-
Double-cliquez sur le nœud d'affectation et sélectionnez un environnement Python 2. Utilisez par exemple
Python 2pour produire un tableau de quatre éléments :Le nœud transmet [10,20,30,40] aux nœuds en aval en divisant automatiquement la dernière ligne de sortie en un tableau à chaque virgule.
print "10,20,30,40" Le nœud d'affectation génère automatiquement un paramètre de sortie nommé
outputs, qui représente son jeu de résultats.Enregistrez le nœud d'affectation.
-
Configurer le nœud for-each pour consommer les données
Configurez le nœud for-each pour recevoir les données en amont et les utiliser dans son corps de boucle.
Double-cliquez sur le nœud for-each pour ouvrir son canevas interne.
-
Dans le panneau Scheduling de droite, localisez le paramètre
loopDataArraysous Scheduling Parameters et cliquez sur Bind.Sélectionnez le paramètre outputs du nœud assign pour créer la liaison. Une fois la liaison établie, la valeur du paramètre loopDataArray reflète son état lié.
Dans la boîte de dialogue qui s'affiche, définissez la Value Source sur le nœud d'affectation en amont (
assign) et sélectionnez son paramètreoutputs. Cette action crée automatiquement une dépendance entre les deux nœuds.-
Dans le corps de boucle for-each, cliquez sur Create Internal Node et créez un nœud
Shell.Dans un scénario réel, vous pouvez configurer n'importe quel type de nœud.
-
Double-cliquez sur le nouveau nœud Shell et utilisez les variables intégrées dans le code pour récupérer et afficher les informations relatives à la boucle :
#!/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}" -
(Facultatif) Dans le panneau Scheduling Settings de droite, configurez les propriétés sous Scheduling Policy.
-
Maximum Number of Loops : La valeur par défaut est 128 et la valeur maximale est 1024.
ImportantCe paramètre détermine le nombre maximal d'itérations du corps de boucle. Si le volume de données en amont est important, augmentez cette valeur pour garantir le traitement de tous les éléments.
-
Execute Policy : Sélectionnez Serial pour cet exemple.
Serial : Exécute les itérations de manière séquentielle.
Parallel : Exécute les itérations de boucle simultanément pour améliorer l'efficacité de la tâche. En mode Parallel, l'échec d'une itération n'affecte pas les autres. Le planificateur tente d'exécuter toutes les itérations jusqu'à leur terme. La concurrence par défaut est de 5, avec un maximum de 20.
-
Enregistrez le nœud Shell.
-
Déployer, exécuter et vérifier
Déployez le workflow dans Operation Center pour l'exécuter et vérifiez les résultats du nœud for-each.
Retournez au canevas principal du workflow et cliquez sur le bouton Deploy de la barre d'outils pour publier l'ensemble du workflow.
-
Accédez à et effectuez un test fumée sur le workflow cible.
ImportantN'effectuez pas de test fumée sur le nœud for-each individuellement. Étant donné que le nœud for-each dépend de la sortie du nœud d'affectation en amont, vous devez lancer le test à partir du nœud d'affectation afin de garantir l'intégrité de la lignée des données.
Une fois l'instance de test exécutée avec succès, localisez l'instance du nœud for-each dans la liste, ouvrez-la, puis faites un clic droit et sélectionnez View Internal Nodes.
-
Dans la vue des nœuds internes, examinez les instances de nœud Shell générées par chaque boucle. Ouvrez le journal d'exécution de n'importe quelle instance pour consulter la sortie de cette itération et vérifier son exactitude.
Le panneau de gauche indique que les quatre itérations de la boucle sont terminées. Le journal d'exécution de la quatrième itération affiche
Current loop number is: 4etCurrent item is: 40, et la commande Shell se termine avec le code 0, ce qui confirme la réussite de l'exécution.
Outre l'utilisation d'un nœud d'affectation classique comme nœud en amont, un nœud for-each permet également d'obtenir le même effet d'itération grâce à la fonctionnalité de paramètre d'affectation d'un nœud SQL en amont. Pour les types de nœuds prenant en charge les paramètres d'affectation, tels que EMR Hive, Hologres SQL, EMR Spark SQL, AnalyticDB for PostgreSQL, ClickHouse SQL et MySQL, vous pouvez ajouter un paramètre d'affectation dans la section Node Context Parameters > Output Parameters of This Node.
Cas d'usage : Traiter différents formats de données
Scénario 1 : Traitement d'un tableau unidimensionnel
Sortie du nœud d'affectation : 2025-11-01,2025-11-02,2025-11-03
Nombre d'itérations : 3
-
Lors de la deuxième itération :
La valeur de
${dag.foreach.current}est2025-11-02.La valeur de
${dag.loopTimes}est2.
Scénario 2 : Traitement d'un tableau bidimensionnel
-
Sortie du nœud d'affectation (MaxCompute SQL) :
+-----+----------+ | id | city | +-----+----------+ | 101 | beijing | | 102 | shanghai | +-----+----------+ Nombre d'itérations : 2
-
Lors de la deuxième itération :
La valeur de
${dag.foreach.current}est102,shanghai.La valeur de
${dag.loopTimes}est2.La valeur de
${dag.foreach.current[0]}est102.La valeur de
${dag.foreach.current[1]}estshanghai.
Scénario : Traitement par lots des données de tables partitionnées sur plusieurs lignes métier
Cet exemple montre comment utiliser un nœud d'affectation et un nœud for-each pour traiter par lots les données de comportement utilisateur sur plusieurs lignes métier, en automatisant le traitement des données avec une logique unique servant plusieurs lignes de produits.
Contexte
Supposons que vous soyez ingénieur en développement de données dans une entreprise Internet globale, responsable du traitement des données de trois lignes métier principales : e-commerce (ecom), finance (finance) et logistique (logistics), avec la possibilité d'en ajouter d'autres à l'avenir. Vous devez exécuter quotidiennement la même logique d'agrégation sur les journaux de comportement utilisateur de ces trois lignes métier afin de calculer le nombre quotidien de pages vues (PV) par utilisateur et stocker les résultats dans une table d'agrégation unifiée.
-
Tables sources en amont (couche DWD) :
dwd_user_behavior_ecom_d: Table de comportement utilisateur e-commerce.dwd_user_behavior_finance_d: Table de comportement utilisateur finance.dwd_user_behavior_logistics_d: Table de comportement utilisateur logistique.dwd_user_behavior_${business_line}_d: Tables de comportement utilisateur pour d'autres lignes métier potentielles futures.Ces tables partagent le même schéma et sont partitionnées par jour (
dt).
-
Table cible en aval (couche DWS) :
dws_user_summary_d: Table d'agrégation utilisateur.Cette table est doublement partitionnée par ligne métier (
biz_line) et par jour (dt) pour stocker les résultats agrégés de toutes les lignes métier de manière unifiée.
La création d'une tâche distincte pour chaque ligne métier entraîne des coûts de maintenance élevés et multiplie les risques d'erreur. Avec un nœud for-each, vous maintenez une seule logique de traitement, et le système itère automatiquement sur toutes les lignes métier pour effectuer le calcul.
Préparation des données
Commencez par créer les tables d'exemple et insérer des données de test (en utilisant la date métier 20251010 comme exemple).
Associez une ressource de calcul au workspace.
Accédez à Data Studio pour le développement de données et créez un nœud MaxCompute SQL.
-
Créez les tables sources (couche DWD) : Ajoutez le code suivant au nœud MaxCompute SQL et exécutez-le.
-- 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 'Event timestamp in milliseconds (Unix)' ) COMMENT 'E-commerce user behavior log detail 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 e-commerce user behavior table created successfully 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 'Event timestamp in milliseconds (Unix)' ) COMMENT 'Finance user behavior log detail 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 finance user behavior table created successfully 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 'Event timestamp in milliseconds (Unix)' ) COMMENT 'Logistics user behavior log detail 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 logistics user behavior table created successfully SELECT * FROM dwd_user_behavior_logistics_d where dt='20251010'; -
Créez la table cible (couche DWS) : Ajoutez le code suivant au nœud MaxCompute SQL et exécutez-le.
CREATE TABLE IF NOT EXISTS dws_user_summary_d ( user_id STRING COMMENT 'User ID', pv BIGINT COMMENT 'Daily activity count' ) COMMENT 'User daily activity summary table' PARTITIONED BY ( dt STRING COMMENT 'Date partition, format yyyymmdd', biz_line STRING COMMENT 'Business line partition, e.g. ecom, finance, logistics' );ImportantSi le workspace utilise le mode standard, vous devez déployer ce nœud dans l'environnement de production et effectuer un backfill des données.
Implémentation du workflow
Créez un workflow. Dans la section Scheduling Parameters située à droite, définissez le paramètre de planification bizdate sur le jour précédent :
$[yyyymmdd-1].-
Dans le workflow, créez un nœud d'affectation nommé get_biz_list et écrivez le code suivant en MaxCompute SQL. Ce nœud produit la liste des lignes métier à traiter :
-- 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; -
Configurer le nœud for-each
Retournez à la page du workflow et créez un nœud for-each en aval du nœud d'affectation get_biz_list.
Ouvrez la page des paramètres du nœud for-each. Dans la section sous schedule settings à droite, liez le paramètre loopDataArray aux outputs du nœud get_biz_list.
-
Dans le corps de boucle du nœud for-each, cliquez sur Create Internal Node et créez un nœud MaxCompute SQL. Rédigez la logique de traitement à l'intérieur du corps de boucle.
RemarqueCe script est piloté par le nœud for-each et s'exécute une fois pour chaque ligne métier.
La variable intégrée ${dag.foreach.current} est remplacée dynamiquement par le nom de la ligne métier courante à chaque itération. Les valeurs d'itération attendues sont : 'ecom', 'finance', '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;
-
Ajouter un nœud de vérification
Retournez au workflow. Cliquez sur Create Downstream sur le nœud for-each pour créer un nœud MaxCompute SQL et ajoutez le code suivant.
SELECT * FROM dws_user_summary_d WHERE dt='20251010' ORDER BY biz_line, user_id;
Déploiement et résultats
Déployez le workflow dans l'environnement de production. Accédez à la page dans Operation Center. Localisez le workflow cible et effectuez un test fumée avec la date métier définie sur '20251010'.
Une fois l'exécution terminée, consultez le journal d'exécution dans l'instance de test. La sortie attendue du nœud final est la suivante :
|
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 |
Avantages
Haute extensibilité : Pour ajouter une nouvelle ligne métier, il suffit d'ajouter une ligne SQL dans le nœud d'affectation, sans modifier la logique de traitement.
Maintenance simplifiée : Toutes les lignes métier partagent une seule logique de traitement. Une modification unique s'applique à l'ensemble.
FAQ
-
Q : Pourquoi ne puis-je pas exécuter directement un nœud for-each dans Data Studio pour le tester ?
R : C'est une limitation inhérente à sa conception. Le nœud nécessite un environnement de planification complet pour résoudre le contexte du nœud et ses dépendances ; il ne prend donc pas en charge l'exécution directe dans Data Studio. Vous devez déployer la tâche dans Operation Center et la tester en utilisant le backfill ou en déclenchant une exécution planifiée.
-
Q : Pourquoi un test fumée sur un nœud for-each individuel échoue-t-il ou ne produit-il aucun résultat ?
R : Les données de boucle d'un nœud for-each proviennent de son paramètre d'entrée
loopDataArray, qui doit être lié au paramètreoutputsd'un nœud d'affectation en amont. Si vous exécutez le nœud for-each seul, il échouera ou sera ignoré car il ne pourra pas recevoir de jeu de résultats en entrée. -
Q : Pourquoi ma boucle ne s'exécute-t-elle qu'une seule fois ?
R : Cela se produit généralement lorsque la sortie du nœud d'affectation en amont est interprétée comme un élément unique. Vérifiez votre sortie :
S'agit-il d'une chaîne unique sans délimiteur ?
Si vous comptez itérer sur plusieurs éléments, assurez-vous qu'ils sont séparés par des virgules (
,). Par exemple,'item1,item2,item3'entraîne trois itérations, tandis que'item1 item2 item3'n'en produit qu'une seule.