このトピックでは、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.32. 操作手順
2.1. エクスポートツールパッケージ airflow-workflow-parser.zip のダウンロード
ダウンロードリンク: airflow-workflow-parser
2.2. airflow-workflow-parser.zip の解凍
unzip airflow-workflow-parser.zip2.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.cfg2.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.22.6.2. DAG ファイルの解析における依存関係エラー
解析が失敗する場合、通常は DAG に特殊な依存関係があることが原因です。エラーメッセージに従って pip コマンドを使用し、不足している依存関係をインストールしてください。
2. DataWorks へのインポート
Airflow エクスポートツールは、スケジューリングタスクフローの JSON 記述ファイルを出力します。このファイルのデータ構造は DataWorks Spec に準拠しています。
ファイルを DataWorks にインポートするには、「カスタム DataWorks Spec を DataWorks にインポートする」をご参照ください。