すべてのプロダクト
Search
ドキュメントセンター

Migration Hub:Airflow から DataWorks への移行

最終更新日:Jun 17, 2026

このトピックでは、LHM スケジューリング移行ツールを使用して Airflow のタスクフローを DataWorks のワークフローに移行する方法について説明します。このプロセスには、Airflow タスクのエクスポート、定期タスクの変換、および DataWorks へのインポートが含まれます。

1. Airflow タスクフローのエクスポート

Airflow タスクフローからメタデータをエクスポートするために、Airflow Python ライブラリは DAG フォルダをロードします。次に、DAG、その内部タスク、およびそれらの依存関係に関する情報を取得します。最後に、このデータを整理してエクスポート用の JSON ファイルにまとめます。

この移行ツールは、Airflow 2.x からのタスクフローのエクスポートをサポートしています。

1. 実行環境

移行ツールのエクスポート機能は、MigrationX Airflow Reader (MigrationX) を使用します。この機能は、バージョン 3.9.0 以降の Python 環境で実行する必要があります。

エクスポートツールには、以下のいずれかの推奨実行環境を選択してください:

1. Airflow のスケジューリングが実行されるのと同じ Python 環境でツールを実行します。この環境は Python 3.9.0 以降を使用する必要があります。

2. 新しい Python 環境を準備します。ご利用の本番環境と同じバージョンの Airflow Python ライブラリをインストールします。Airflow ライブラリと Python バージョンの互換性にご注意ください。

2 番目のオプションを選択した場合は、次の手順に従って新しい Python 環境を構成できます。この例では、Linux ECS インスタンスを使用します:

## Conda をインストールし、Python 3.9 環境を作成します

# Conda をダウンロードしてインストールします
wget https://repo.anaconda.com/miniconda/Miniconda3-latest-Linux-x86_64.sh
sh Miniconda3-latest-Linux-x86_64.sh

# インストール中に、Conda はインストールディレクトリの指定を求めます。インストール後、次のコマンドを実行して初期化します。
cd <conda_installation_directory>
conda init
conda config --set auto_activate_base false
. ~/.bashrc

# Python 3.9.0 環境を作成します
conda create -n airflow python=3.9.0
conda activate airflow

# Conda 環境に Airflow をインストールします (例:バージョン 2.5.3)
pip install apache-airflow==2.5.3

2. 操作手順

2.1. エクスポートツールパッケージ airflow-workflow-parser.zip のダウンロード

ダウンロードリンク: airflow-workflow-parser

2.2. airflow-workflow-parser.zip の解凍

unzip airflow-workflow-parser.zip

2.3. PYTHONPATH の設定

Airflow 環境でエクスポートツールを実行する場合は、PYTHONPATH を設定します。Conda 環境でツールを実行する場合は、このステップをスキップできます。

PYTHONPATH を Airflow Python ライブラリのディレクトリに設定します。例:

export PYTHONPATH=/usr/local/lib/python3.6/site-packages

export AIRFLOW_HOME=/var/lib/airflow

## 注:新しい Airflow インスタンスに airflow.cfg ファイルがない場合は、「airflow config list --defaults」を実行して生成できます。
export AIRFLOW_CONFIG=/var/run/cloudera-scm-agent/process/2531-airflow-AIRFLOW_SCHEDULER/airflow.cfg

2.4. Airflow タスクのエクスポート

パーサは、指定されたディレクトリからすべての DAG を取得します。各 DAG に対して、DataWorks Spec ファイルを生成します。このファイルは、DataWorks タスクフローを定義する JSON ファイルです。

DataWorks タスクフロー定義: dataworks-spec
https://github.com/aliyun/dataworks-spec/tree/master
python ./parser.py -d /root/airflow/airflow-workflow/dags -o ./result -m ./flowspec-airflowV2-transformer-config.json

パラメーターの説明:

-d: Airflow DAG ファイルが格納されているディレクトリを指定します。

-o: エクスポートされた JSON ファイルを格納するパスを指定します。

-m: 設定ファイルのパスを指定します。このファイルでノードマッピングと変換ルールを構成できます。設定項目の詳細については、セクション 2.4.1 をご参照ください。

エクスポートが失敗する場合、通常は DAG 内の特殊な依存関係が原因です。エラーメッセージに従って pip コマンドを使用し、不足している依存関係をインストールしてください。

2.4.1. 変換設定テンプレート

エクスポートツールはノードタイプを変換できます。設定項目を使用して、Airflow オペレーターを DataWorks のノードタイプにマッピングできます。以下のテンプレートはその方法を示しています。

DataWorks on MaxCompute への移行用設定テンプレート (例):

{
  "workflowPathPrefix": "Airflow_Import_V3/",
  "typeMapping": {
    "EmptyOperator": "VIRTUAL",
    "DummyOperator": "VIRTUAL",
    "ExternalTaskSensor": "VIRTUAL",
    "BashOperator": "DIDE_SHELL",
    "HiveToMySqlTransfer": "DI",
    "PrestoToMySqlTransfer": "DI",
    "PythonOperator": "PYTHON",
    "HiveOperator": "ODPS_SQL",
    "SqoopOperator": "DI",
    "SparkSqlOperator": "ODPS_SQL",
    "SparkSubmitOperator": "ODPS_SPARK",
    "SQLExecuteQueryOperator": "MySQL",
    "PostgresOperator": "Postgresql",
    "MySqlOperator": "MySQL",
    "default": "PYTHON"
  },
  "settings": {
    "workflow.converter.target.schedule.resGroupIdentifier": "Serverless_res_group_651147510078336_710919704558240"
  }
}

DataWorks on EMR への移行用設定テンプレート (例):

{
  "workflowPathPrefix": "Airflow_Import_V3/",
  "typeMapping": {
    "EmptyOperator": "VIRTUAL",
    "DummyOperator": "VIRTUAL",
    "ExternalTaskSensor": "VIRTUAL",
    "BashOperator": "EMR_SHELL",
    "HiveToMySqlTransfer": "DI",
    "PrestoToMySqlTransfer": "DI",
    "PythonOperator": "PYTHON",
    "HiveOperator": "EMR_HIVE",
    "SqoopOperator": "EMR_SQOOP",
    "SparkSqlOperator": "EMR_SPARK_SQL",
    "SparkSubmitOperator": "EMR_SPARK",
    "SQLExecuteQueryOperator": "MySQL",
    "PostgresOperator": "Postgresql",
    "MySqlOperator": "MySQL",
    "default": "PYTHON"
  },
  "settings": {
    "workflow.converter.target.schedule.resGroupIdentifier": "Serverless_res_group_651147510078336_710919704558240"
  }
}

パラメーターの説明:

パラメーター名

説明

workflowPathPrefix

DataWorks にインポートされた後のタスクフローのストレージパスです。

typeMapping

Airflow オペレーターと DataWorks ノードタイプ間のマッピングルールです。

DataWorks のノードタイプの詳細については、こちらの列挙クラスをご参照ください: https://github.com/aliyun/dataworks-spec/blob/b0f4a4fd769215d5f81c0bbe990addd7498df5f4/spec/src/main/java/com/aliyun/dataworks/common/spec/domain/dw/types/CodeProgramType.java#L180

workflow.converter.target.schedule.resGroupIdentifier

DataWorks リソースグループの ID です。

ご利用の DataWorks ワークスペースの製品ページに移動します。左側のナビゲーションウィンドウで、[リソースグループ] を選択します。リソースグループをアタッチし、その ID を取得します。

2.5. 例

ツールをダウンロードして解凍すると、次のディレクトリ構造が得られます:

.
├── airflow_workflow
│   ├── common
│   ├── connections.py
│   ├── converter
│   ├── dag_parser.py
│   ├── downloader.py
│   ├── get_dags.py
│   ├── __init__.py
│   ├── miscs
│   ├── models
│   ├── patch
│   ├── __pycache__
│   └── test
├── airflow-workflow.tgz
├── dags
│   ├── example.py
│   └── __pycache__
├── flowspec-airflowV2-transformer-config.json
├── parser.py
├── README.MD
└── result

コマンドを実行します。

python ./parser.py -d /root/airflow/airflow-workflow/dags -o ./result -m ./flowspec-airflowV2-transformer-config.json

運用ログは次のとおりです:

変換結果を表示します:

2.6. よくある質問

2.6.1. エラー: TIMEZONE = pendulum.tz.timezone("UTC") TypeError: 'module' object is not callable

Airflow 2.x の場合、pendulum のバージョンが 2.1.2 以前であることを確認してください。

Airflow 3.x の場合、pendulum は不要です。Airflow 3.x では pendulum への依存関係が削除されました。

解決策:

pip uninstall pendulum -y
pip install pendulum==2.1.2
2.6.2. DAG ファイルの解析における依存関係エラー

解析が失敗する場合、通常は DAG に特殊な依存関係があることが原因です。エラーメッセージに従って pip コマンドを使用し、不足している依存関係をインストールしてください。

2. DataWorks へのインポート

Airflow エクスポートツールは、スケジューリングタスクフローの JSON 記述ファイルを出力します。このファイルのデータ構造は DataWorks Spec に準拠しています。

ファイルを DataWorks にインポートするには、「カスタム DataWorks Spec を DataWorks にインポートする」をご参照ください。