for-each ノード は、ファイル名やパーティションのリストなど、上流の結果セットを反復処理し、要素ごとに同じサブタスクを実行します。これにより、個別のタスクを手動で作成する必要がなくなり、動的で自動化されたワークフローを実現できます。
ユースケース
for-each ノード を使用すると、異なる事業部門、製品ライン、または設定項目に同じ分析や処理ロジックを適用する必要がある場合に、パラメーター化した実行が可能になります。たとえば、会社に複数の製品ラインがあり、それぞれに対して個別のデイリーレポートを生成する必要がある場合、処理ロジックは同じで、対象データのみが異なります。
プログラミング言語の for ループと同様に、for-each ノード はテーブル名、パーティション名、ファイル名などのリストを反復処理し、各項目に対して事前に定義されたサブワークフローを実行します。
使用上の注意
バージョン要件:DataWorks Standard Edition 以降でのみ利用可能です。
権限:ご使用の RAM アカウントを対象の ワークスペース に追加し、開発者 または ワークスペース管理者 のロールを割り当てる必要があります。詳細については、「ワークスペースへのメンバー追加」をご参照ください。
仕組み
for-each ノードは、「ループ本体」と呼ばれるカスタマイズ可能なサブワークフローをカプセル化するコンテナとして機能します。次のように動作します。
-
データ入力: for-each ノードは、上流の代入ノードまたはその他の代入可能なノード (EMR Hive ノードなど) に依存し、
loopDataArrayパラメーターにバインドすることで配列形式の結果セットを取得します。 -
ループ実行: ノードが開始すると、結果セット内の各要素を順番に反復処理します。各要素に対して、内部のループ本体を
StartノードからEndノードまで一度完全に実行します。説明Start ノードと End ノードは編集できません。これらはループ本体の開始と終了を示すだけです。
-
データ受け渡し: 各反復中に、現在の要素の値が組み込み変数を介してループ本体内のノードに渡されます。内部のビジネスノードは
${dag.foreach.current}を使用して、処理中のデータ項目にアクセスします。
組み込みパラメーター
${...} 形式の変数は、DataWorks 固有のテンプレート構文です。DataWorks はこれらのパラメーターを直接解析し、実行前にその値に置き換えます。
for-each ループ本体内のノードは、次の組み込み変数を使用してループのステータスとデータにアクセスできます。
|
組み込みパラメーター |
説明 |
For ループの類推 |
|
|
上流の割り当てノードから渡された完全な結果セット。 |
例えば、次のような for ループコードがあるとします。
|
|
|
現在の反復で処理中のデータ項目。 |
|
|
|
現在のループのオフセット (0 から始まるインデックス)。 |
|
|
|
現在のループ回数 (1 から始まるインデックス)。 |
上流の出力が SQL クエリ結果のような二次元配列である場合、次の構文を使用して特定の値にアクセスすることもできます。
|
その他のパラメーター |
説明 |
|
|
現在のデータ行 (一次元配列) の要素をカンマ |
|
|
現在のデータ行の |
|
|
結果セット全体の for-each ノードは現在、ネストされたループをサポートしていません。この例は、値の取得方法を説明するためのものです。 |
制限事項
-
実行メカニズム: ループは直列実行と並列実行の両方をサポートします。反復が互いに独立している場合は、並列実行を選択できます。
-
ループ制限: デフォルトの最大ループ数は 128 で、最大 1024 まで調整可能です。
-
デバッグの制約: Data Studio で for-each ノードを直接実行することはできません。タスクをデプロイし、オペレーションセンターでスモークテスト機能を使用してテストする必要があります。
-
実行の制約: for-each ノードを単独で実行することはできません。スモークテスト、バックフィル、手動実行のいずれの場合も同様です。
-
ループ本体内のフロー制御: for-each ループ本体内で分岐ノードを使用する場合、すべての分岐が
Endノードに接続する前に、最終的に単一の合流ノードで合流させる必要があります。これにより、ループ本体の論理的な整合性が保証されます。 -
再実行の制約: ノードがデプロイされた後、失敗時の自動再実行は失敗した時点から再開します。ただし、手動での再実行は、for-each ノード全体の完全な再実行をトリガーします。
操作手順
この手順では、上流ノードとして割り当てノードを使用し、ループ本体内で Shell ノードを使用して結果を出力します。完全な for-each タスクの設定手順を説明します。
-
上流データの準備 (割り当てノードの設定)
割り当てノードを作成および設定して、下流の for-each ノードに反復処理用の結果セットを提供します。
-
ワークフローで、割り当てノード (たとえば、
assign) を作成し、for-each ノードの上流に配置します。 -
割り当てノードをダブルクリックし、Python 2 環境を選択します。たとえば、
Python 2を使用して 4 つの要素を持つ配列を出力します。ノードは、最後の出力行をカンマを区切り文字として配列に自動的に分割し、下流のノードに [10,20,30,40] を出力します。
print "10,20,30,40" -
割り当てノードは、その結果セットを表す
outputsという名前の出力パラメーターを自動的に生成します。 -
割り当てノードを保存します。
-
-
データを消費する for-each ノードの設定
上流のデータを受信し、それをループ本体内で使用するように for-each ノードを設定します。
-
for-each ノードをダブルクリックして、その内部キャンバスを開きます。
-
右側の[スケジューリング]パネルで、[スケジューリングパラメーター]の下にある
loopDataArrayパラメーターを見つけ、[バインド]をクリックします。assign ノードの outputs パラメーターを選択してバインディングを作成します。バインディングが完了すると、loopDataArray パラメーターの値にそのバインド状態が反映されます。
-
表示されるダイアログボックスで、[値ソース] を上流の [割り当てノード] (
assign) に設定し、そのoutputsパラメーターを選択します。この操作により、2 つのノード間の依存関係が自動的に作成されます。 -
[for-each ループ本体] で、[内部ノードの作成] をクリックし、
Shellノードを作成します。実際のシナリオでは、任意のタイプのノードを設定できます。
-
新しい Shell ノードをダブルクリックし、コード内で組み込み変数を使用してループに関する情報を取得し、出力します。
#!/bin/bash # ${dag.loopTimes} を使用して現在のループ数を取得 echo "Current loop number is: ${dag.loopTimes}" # ${dag.foreach.current} を使用して現在の反復のデータ項目を取得 echo "Current item is: ${dag.foreach.current}" -
(任意) 右側のスケジューリング設定パネルで、スケジューリングポリシーのプロパティを設定します。
-
[最大ループ回数]: デフォルトは 128 で、最大は 1024 です。
重要このパラメーターは、ループ本体の最大反復回数を決定します。上流のデータ項目数が多い場合は、すべての項目を処理できるようにこの値を増やしてください。
-
[実行ポリシー]:この例では、シリアルを選択します。
-
[シリアル]: 反復を順次実行します。
-
[並列]: ループの反復を同時に実行して、タスクの効率を向上させます。並列モードでは、1 つの反復が失敗しても、他の反復には影響しません。スケジューラはすべての反復の完了を試みます。デフォルトの並列度は 5、最大は 20 です。
-
-
-
Shell ノードを保存します。
-
-
デプロイ、実行、検証
ワークフローをオペレーションセンターにデプロイして実行し、for-each ノードの結果を検証します。
-
メインのワークフローキャンバスに戻り、ツールバーの [デプロイ] ボタンをクリックしてワークフロー全体を公開します。
-
に移動し、対象のワークフローでスモークテストを実行します。
重要for-each ノードを個別にスモークテストしないでください。for-each ノードは上流の割り当てノードの出力に依存するため、データリネージが完全であることを確認するために、テストを割り当てノードから開始する必要があります。
-
テストインスタンスが正常に実行された後、リストから for-each ノードインスタンスを見つけて開き、右クリックして [内部ノードを表示] を選択します。
-
内部ノードビューで、各ループが生成した Shell ノードインスタンスを確認します。任意のインスタンスの実行ログを開いてその反復の出力を表示し、出力が正しいことを確認します。
左側のパネルに、4 つのループ反復がすべて完了したことが表示されます。4 回目の反復の実行ログには
Current loop number is: 4とCurrent item is: 40が出力され、Shell コマンドはコード 0 で終了し、正常に実行されたことを示しています。
-
従来の上流ノードとして割り当てノードを使用するだけでなく、for-each ノードは、上流 SQL ノードの割り当てパラメーター機能を通じて同じ反復効果を実現することも可能です。EMR Hive、Hologres SQL、EMR Spark SQL、AnalyticDB for PostgreSQL、ClickHouse SQL、MySQL など、割り当てパラメーターをサポートするノードタイプの場合、[ノードコンテキストパラメーター] > [このノードの出力パラメーター] セクションで割り当てパラメーターを追加できます。
ユースケース:異なるデータ形式の処理
シナリオ 1:一次元配列の処理
-
割り当てノードの出力: 2025-11-01,2025-11-02,2025-11-03
-
反復回数: 3
-
2 回目の反復では:
-
${dag.foreach.current}の値は2025-11-02です。 -
${dag.loopTimes}の値は2です。
-
シナリオ 2:二次元配列の処理
-
割り当てノード (MaxCompute SQL) の出力:
+-----+----------+ | id | city | +-----+----------+ | 101 | beijing | | 102 | shanghai | +-----+----------+ -
反復回数: 2
-
2 回目の反復では:
-
${dag.foreach.current}の値は102,shanghaiです。 -
${dag.loopTimes}の値は2です。 -
${dag.foreach.current[0]}の値は102です。 -
${dag.foreach.current[1]}の値はshanghaiです。
-
シナリオ:複数のビジネスラインにまたがるパーティションテーブルデータの一括処理
この例では、割り当てノードと for-each ノードを使用して、複数のビジネスラインにわたるユーザー行動データを一括処理する方法を示します。これにより、単一のロジックで複数の製品ラインに対応するデータ処理を自動化します。
背景情報
ある総合インターネット企業でデータ開発エンジニアであると仮定します。あなたは3つの主要ビジネスライン (eコマース (ecom)、金融 (finance)、物流 (logistics)) のデータを処理しており、将来的にはさらに追加される可能性があります。毎日、これら 3 つのビジネスラインのユーザー行動ログに対して同じ集計ロジックを実行し、ユーザーごとの日次ページビュー (PV) を計算し、その結果を統一された集計テーブルに保存する必要があります。
上流のソーステーブル (DWD レイヤー):
dwd_user_behavior_ecom_d: E コマースのユーザー行動テーブル。dwd_user_behavior_finance_d: 金融ユーザー行動テーブル。dwd_user_behavior_logistics_d:物流ユーザー行動テーブル。dwd_user_behavior_${business_line}_d:将来、さらにビジネスラインが追加される可能性に対応するユーザー行動テーブル。これらのテーブルはスキーマが同じで、日単位 (
dt) でパーティション化されています。
下流のターゲットテーブル (DWS レイヤー):
dws_user_summary_d:ユーザー集計テーブル。このテーブルは、すべてのビジネスラインからの集計結果を一元的に保存するために、ビジネスライン(
biz_line)と日(dt)で二重にパーティション化されています。
ビジネスラインごとに個別のタスクを作成すると、メンテナンスコストが高くなり、エラーが発生しやすくなります。for-each ノードを使用すると、単一の処理ロジックを維持するだけで、システムが自動的にすべてのビジネスラインを反復処理して計算を完了します。
データ準備
まず、サンプルテーブルを作成し、テストデータを挿入します (業務日 20251010 を例として使用)。
ワークスペースにコンピューティングリソースを関連付けます。
DataStudio に移動してデータ開発を行い、MaxCompute SQL ノードを作成します。
ソーステーブル (DWD レイヤー) の作成:次のコードを MaxCompute SQL ノードに追加して実行します。
-- eコマースユーザー行動テーブル CREATE TABLE IF NOT EXISTS dwd_user_behavior_ecom_d ( user_id STRING COMMENT 'ユーザー ID', action_type STRING COMMENT 'アクションタイプ', event_time BIGINT COMMENT 'イベントタイムスタンプ (ミリ秒単位の Unix)' ) COMMENT 'eコマースユーザー行動ログ詳細テーブル' PARTITIONED BY (dt STRING COMMENT '日付パーティション、形式 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 -- eコマースユーザー行動テーブルが正常に作成されたことを確認します SELECT * FROM dwd_user_behavior_ecom_d where dt='20251010'; -- 金融ユーザー行動テーブル CREATE TABLE IF NOT EXISTS dwd_user_behavior_finance_d ( user_id STRING COMMENT 'ユーザー ID', action_type STRING COMMENT 'アクションタイプ', event_time BIGINT COMMENT 'イベントタイムスタンプ (ミリ秒単位の Unix)' ) COMMENT '金融ユーザー行動ログ詳細テーブル' PARTITIONED BY (dt STRING COMMENT '日付パーティション、形式 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 -- 金融ユーザー行動テーブルが正常に作成されたことを確認します SELECT * FROM dwd_user_behavior_finance_d where dt='20251010'; -- 物流ユーザー行動テーブル CREATE TABLE IF NOT EXISTS dwd_user_behavior_logistics_d ( user_id STRING COMMENT 'ユーザー ID', action_type STRING COMMENT 'アクションタイプ', event_time BIGINT COMMENT 'イベントタイムスタンプ (ミリ秒単位の Unix)' ) COMMENT '物流ユーザー行動ログ詳細テーブル' PARTITIONED BY (dt STRING COMMENT '日付パーティション、形式 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 -- 物流ユーザー行動テーブルが正常に作成されたことを確認します SELECT * FROM dwd_user_behavior_logistics_d where dt='20251010';ターゲットテーブル (DWS レイヤー) の作成:次のコードを MaxCompute SQL ノードに追加して実行します。
CREATE TABLE IF NOT EXISTS dws_user_summary_d ( user_id STRING COMMENT 'ユーザー ID', pv BIGINT COMMENT '日次アクティビティ数' ) COMMENT 'ユーザー日次アクティビティサマリーテーブル' PARTITIONED BY ( dt STRING COMMENT '日付パーティション、形式 yyyymmdd', biz_line STRING COMMENT 'ビジネスラインパーティション、例:ecom, finance, logistics' );重要ワークスペースが標準モードを使用している場合、このノードを本番環境にデプロイし、データをバックフィルする必要があります。
ワークフローの実装
ワークフローを作成します。右側のスケジューリングパラメータセクションで、スケジューリングパラメータ bizdate に、前日を表す
$[yyyymmdd-1]を設定します。ワークフローで、get_biz_list という名前の割り当てノードを作成し、MaxCompute SQL で次のコードを記述します。このノードは、処理対象のビジネスラインのリストを出力します。
-- 処理するすべてのビジネスラインを出力 SELECT 'ecom' AS biz_line UNION ALL SELECT 'finance' AS biz_line UNION ALL SELECT 'logistics' AS biz_line;for-each ノードの設定
ワークフローページに戻り、割り当てノード get_biz_list の下流に for-each ノードを作成します。
for-each ノード設定ページを開きます。右側の スケジューリング設定 の下にある セクションで、loopDataArray パラメーターを get_biz_list ノードの outputs にバインドします。
for-each ノードのループ本体で、[Create Inner Node] をクリックし、MaxCompute SQL ノードを作成します。ループ本体内に処理ロジックを記述します。
説明このスクリプトは for-each ノードによって駆動され、ビジネスラインごとに 1 回実行されます。
組み込み変数
${dag.foreach.current}は、各反復で現在のビジネスライン名に動的に置き換えられます。期待される反復値は '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;
検証ノードの追加
ワークフローに戻ります。for-each ノードで [Create Downstream] をクリックして MaxCompute SQL ノードを作成し、次のコードを追加します。
SELECT * FROM dws_user_summary_d WHERE dt='20251010' ORDER BY biz_line, user_id;
デプロイと結果
ワークフローを本番環境にデプロイします。運用センターの ページに移動し、対象のワークフローを見つけ、ビジネス日付を '20251010' に設定してスモークテストを実行します。
実行完了後、テストインスタンスで実行ログを表示します。最終ノードの期待される出力は次のとおりです。
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 |
メリット
高いスケーラビリティ:新しいビジネスラインを追加するには、処理ロジックを変更することなく、割り当てノードに SQL を 1 行追加するだけです。
容易なメンテナンス:すべてのビジネスラインが 1 つの処理ロジックを共有します。1 回の変更がすべてに適用されます。
よくある質問
-
Q: Data Studio で for-each ノードを直接実行してテストできないのはなぜですか?
A: これは仕様です。ノードは、ノードコンテキストとその依存関係を解決するために完全なスケジューリング環境を必要とするため、Data Studio での直接実行はサポートしていません。タスクをオペレーションセンターにデプロイし、バックフィルまたはスケジュールされた実行をトリガーしてテストする必要があります。
-
Q: 個別の for-each ノードでスモークテストが失敗するか、何も実行されないのはなぜですか?
A: for-each ノードのループデータは、その
loopDataArray入力パラメーターから取得されます。これは、上流の割り当てノードのoutputsパラメーターにバインドされている必要があります。for-each ノードを単独で実行すると、入力の結果セットを受信できないため、失敗するかスキップされます。 -
Q: ループが 1 回しか実行されないのはなぜですか?
A: これは通常、上流の割り当てノードからの出力が単一の要素として解析されるために発生します。出力を確認してください。
-
1. 区切り文字のない単一の文字列ではありませんか?
-
2. 複数の項目を反復処理する場合は、それらがカンマ (
,) で区切られていることを確認してください。たとえば、'item1,item2,item3'は 3 回のループになりますが、'item1 item2 item3'は 1 回しかループしません。
-