このトピックでは、LHM スケジュール移行ツールを使用して、カスタム DataWorks Spec を DataWorks にインポートする方法について説明します。
1. DataWorks Spec の構築と解析
LHM 移行ツールを使用すると、DataWorks Spec と呼ばれる標準の DataWorks ワークフロー構造に基づいて記述ファイルを作成し、インポートできます。インポートの前に、LHM は DataWorks Spec を解析し、標準の移行パッケージを生成します。
標準の移行パッケージが生成された後は、それ以上の変換は必要ありません。インポートツールを直接実行して、ワークフローを DataWorks に書き込むことができます。
1.1 DataWorks Spec の定義
https://github.com/aliyun/dataworks-spec/tree/master
1.2 入力パッケージのビルド
DataWorks Spec の定義に基づいてワークフロー定義を作成し、JSON ファイルとして保存します。
すべてのワークフロー JSON ファイルを、サブフォルダのない単一のフォルダに配置します。フォルダを .zip ファイルに圧縮します。
{
"version": "1.1.0",
"kind": "CycleWorkflow",
"spec": {
"nodes": [],
"workflows": [
{
"id": "3451387436863448",
"outputs": {
"nodeOutputs": [
{
"artifactType": "NodeOutput",
"data": "3451387436863448",
"refTableName": "example_hive_operator"
}
]
},
"nodes": [
{
"id": "2105761738722077",
"name": "run_first",
"instanceMode": "T+1",
"rerunMode": "FailureAllowed",
"rerunTimes": 0,
"rerunInterval": 300000,
"trigger": {
"type": "Scheduler",
"startTime": "2025-03-30 00:00:00",
"timezone": "UTC"
},
"runtimeResource": {
"resourceGroup": "Serverless_res_group_651147510078336_710919704558240"
},
"script": {
"path": "run_first",
"runtime": {
"command": "EMR_HIVE"
},
"parameters": [],
"content": "\n create database if not exists airflow;\n use airflow;\n drop table if exists test_hive;\n create table test_hive(name string);\n insert into test_hive values('studio');\n "
},
"outputs": {
"nodeOutputs": [
{
"data": "2105761738722077"
},
{
"artifactType": "NodeOutput",
"data": "example_hive_operator.run_first.run_second",
"refTableName": "run_first"
}
]
},
"type": "EMR_HIVE"
},
{
"id": "3326558787158921",
"name": "run_second",
"instanceMode": "T+1",
"rerunMode": "FailureAllowed",
"rerunTimes": 0,
"rerunInterval": 300000,
"trigger": {
"type": "Scheduler",
"startTime": "2025-03-30 00:00:00",
"timezone": "UTC"
},
"runtimeResource": {
"resourceGroup": "Serverless_res_group_651147510078336_710919704558240"
},
"script": {
"path": "run_second",
"runtime": {
"command": "EMR_HIVE"
},
"parameters": [],
"content": "\n use airflow;\n add jar oss://emr-studio-example/hive-udf-1.0-SNAPSHOT.jar;\n create temporary function simpleudf AS 'com.aliyun.emr.hive.udf.SimpleUDFExample';\n show functions like '*udf';\n select simpleudf(name) from test_hive;\n "
},
"outputs": {
"nodeOutputs": [
{
"data": "3326558787158921"
}
]
},
"type": "EMR_HIVE"
}
],
"dependencies": [
{
"nodeId": "3326558787158921",
"depends": [
{
"type": "Normal",
"output": "2105761738722077",
"refTableName": "run_first"
}
]
}
],
"script": {
"path": "Airflow_Import_V3/example_hive_operator",
"runtime": {
"command": "WORKFLOW"
},
"parameters": []
},
"name": "example_hive_operator",
"trigger": {
"type": "Scheduler",
"cron": "0 0 0 * * ?",
"timezone": "UTC",
"delaySeconds": 0
},
"type": "CycleWorkflow",
"strategy": {
"timeout": 0,
"instanceMode": "T+1",
"rerunMode": "Allowed",
"rerunTimes": 0,
"rerunInterval": 0,
"failureStrategy": "Continue"
}
}
],
"flow": [
{
"nodeId": "3451387436863448",
"depends": []
}
]
}
}圧縮コマンドの例
zip -q -r -m -o <PackageName>.zip PackageName1.3 設定項目のビルド
以下の設定項目を使用します。内容は変更しないでください。スケジュール情報はカスタム DataWorks Spec ですでに記述されているため、多くの設定項目を追加する必要はありません。
{
"schedule_datasource": {
"name": "MySpec",
"type": "DataWorks",
"operaterType": "MANUAL"
},
"conf": {}
}1.4 スケジュール解析ツールの実行
コマンドラインから次のコマンドを使用して解析ツールを実行します。
sh ./bin/run.sh read \
-c ./conf/<your_config_file>.json \
-f ./data/0_OriginalPackage/<input_package>.zip \
-o ./data/1_ReaderOutput/<standard_migration_package>.zip \
-t dw-newidea-reader-c パラメーターは設定ファイルのパスを指定します。-f パラメーターは入力パッケージのパスを指定します。-o パラメーターは標準の移行パッケージが生成されるパスを指定します。-t パラメーターは検出プラグインの名前を指定します。
たとえば、プロジェクト A を解析するには、次のようにします。
sh ./bin/run.sh read \
-c ./conf/projectA_read.json \
-f ./data/0_OriginalPackage/projectA_DataworksSpec.zip \
-o ./data/1_ReaderOutput/projectA_ReaderOutput.zip \
-t dw-newidea-reader検出ツールは実行中にプロセス情報を表示します。プロセス中に発生したエラーを確認できます。
1.5 解析結果の表示
./data/1_ReaderOutput/ ディレクトリで生成された ReaderOutput.zip パッケージを開いて、解析結果をプレビューできます。
統計レポートには、DataWorks Spec からのワークフロー、ノード、リソース、関数、およびデータソースに関する基本情報の概要が記載されています。
data/project フォルダには、DataWorks Spec からのスケジュール情報の標準化されたデータ構造が含まれています。
統計レポートには、2 つの特別な機能があります:
1. レポートでワークフローとノードの一部のプロパティを変更できます。編集可能なフィールドは青色で表示されます。ツールは、データを DataWorks にインポートするときにこれらのプロパティの変更を適用します。
2. レポートのワークフロー子テーブルから対応する行を削除することで、インポート中にワークフローをスキップできます。これはワークフローのブラックリストとして機能します。注: ワークフローが相互に依存している場合は、同じバッチでインポートする必要があります。ブラックリストを使用してそれらを分離しないでください。エラーの原因となる可能性があります。
詳細については、「スケジュール移行概要レポートを使用したスケジュールプロパティの更新」をご参照ください。
2. DataWorks へのインポート
インポートツールは複数の書き込み操作をサポートし、上書きモードでワークフローを自動的に作成または更新します。
2.1 前提条件
2.1.1 DataWorks Spec が正常に解析される
DataWorks Spec が解析され、ReaderOutput.zip ファイルが生成されます。
(任意ですが推奨) 解析出力パッケージを開き、統計レポートを表示して、移行範囲が正しく解析されたことを確認できます。
2.1.2 DataWorks の設定
DataWorks で次の操作を実行します。
1. ワークスペースを作成します。
2. AccessKey ペア (AK と SK) を作成し、そのペアにワークスペースの管理者権限を付与します。書き込み操作のトラブルシューティングを簡素化するために、アカウントにバインドされた AccessKey ペアを作成することを強くお勧めします。
3. ワークスペースで、データソースを作成し、計算リソースをアタッチし、リソースグループを作成します。
4. ワークスペースで、ファイルリソースをアップロードし、ユーザー定義関数 (UDF) を作成します。
2.1.3 ネットワーク接続の確認
DataWorks エンドポイントへの接続を確認します。
エンドポイントのリスト:
ping dataworks.aliyuncs.com2.2 インポートの設定
プロジェクトディレクトリの `conf` フォルダに、`writer.json` などの JSON 形式の設定ファイルを作成します。
ファイルを使用する前に、JSON コードからコメントを削除してください。
{
"schedule_datasource": {
"name": "YourDataWorks", // DataWorks データソースに名前を付けます。
"type": "DataWorks",
"properties": {
"endpoint": "dataworks.cn-hangzhou.aliyuncs.com", // サービスエンドポイント
"project_id": "YourProjectId", // ワークスペース ID
"project_name": "YourProject", // ワークスペース名
"ak": "************", // AK
"sk": "************" // SK
},
"operaterType": "MANUAL"
},
"conf": {
"di.resource.group.identifier": "Serverless_res_group_***_***", // スケジュールリソースグループ
"resource.group.identifier": "Serverless_res_group_***_***", // データ統合リソースグループ
"dataworks.node.type.xls": "/Software/bwm-client/conf/CodeProgramType.xls", // DataWorks ノードタイプテーブルへのパス
"qps.limit": 5 // DataWorks への API リクエストの QPS 制限
}
}2.2.1 エンドポイント
DataWorks ワークスペースのリージョンに基づいてエンドポイントを選択します。詳細については、次のドキュメントをご参照ください。
2.2.2 ワークスペース ID と名前
DataWorks コンソールを開き、ワークスペース詳細ページに移動します。右側の [基本情報] セクションでワークスペース ID と名前を確認できます。
2.2.3 AccessKey ペアの作成と承認
ユーザーページで AccessKey ペア (AK と SK) を作成します。AccessKey ペアには、ターゲットの DataWorks ワークスペースに対する管理者の読み取りおよび書き込み権限が必要です。
権限管理は 2 つの場所で処理されます。Resource Access Management (RAM) ユーザーを使用している場合は、まず RAM ユーザーに DataWorks 操作を実行するための権限を付与する必要があります。
アクセスポリシーページ: https://ram.console.alibabacloud.com/policies
次に、DataWorks ワークスペースで、アカウントにワークスペース権限を割り当てることができます。
注: AccessKey にネットワークアクセスポリシーを設定できます。移行ツールを実行しているマシンの IP アドレスがアクセスポリシーに含まれていることを確認してください。
2.2.4 リソースグループ
DataWorks ワークスペース詳細ページで、左側のメニューバーから [リソースグループ] ページに移動します。リソースグループをアタッチし、その ID を取得します。
汎用リソースグループは、ノードのスケジューリングとデータ統合の両方に使用できます。スケジューリングリソースグループ (`resource.group.identifier`) とデータ統合リソースグループ (`di.resource.group.identifier`) の両方を同じ汎用リソースグループに設定できます。
2.2.5 QPS 設定
このツールは、DataWorks API を呼び出すことによってデータをインポートします。DataWorks のバージョンごとに、クエリ/秒 (QPS) と、読み取りおよび書き込み OpenAPI 操作の 1 日あたりの呼び出し回数に独自の制限があります。詳細については、「制限」をご参照ください。
DataWorks Basic、Standard、および Professional Edition の場合、`"qps.limit"` を `5` に設定できます。Enterprise Edition の場合、`"qps.limit"` を `20` に設定できます。
注: 複数のインポートツールを同時に実行しないでください。
2.2.6 DataWorks ノードタイプ ID の設定
DataWorks では、一部のノードタイプにはリージョンごとに異なる TypeIds が割り当てられます。特定の TypeID は、DataWorks データ開発インターフェイスによって決定されます。この属性は、主にデータベースノードに適用されます。詳細については、「データベースノード」をご参照ください。
たとえば、MySQL ノードの NodeTypeId は、中国 (杭州) リージョンでは 1000039、中国 (深圳) リージョンでは 1000041 です。
DataWorks のこれらのリージョン差を処理するために、このツールは、使用するノード TypeId テーブルを設定する方法を提供します。
インポートツールの設定項目を使用してテーブルをインポートできます。
"conf": {
"dataworks.node.type.xls": "/Software/bwm-client/conf/CodeProgramType.xls" // DataWorks ノードタイプテーブルへのパス
}DataWorks データ開発インターフェイスからノードタイプ ID を取得するには、新しいワークフローを作成し、それに新しいノードを追加します。ノードを保存した後、ワークフローの Spec を表示できます。
ノードタイプが正しく設定されていない場合、ワークフローを公開するときに次のエラーが表示されます。
2.3 DataWorks インポートツールの実行
コマンドラインから次のコマンドを使用してインポートツールを実行します。
sh ./bin/run.sh write \
-c ./conf/<your_config_file>.json \
-f ./data/1_ReaderOutput/<parsing_result_output_package>.zip \
-o ./data/4_WriterOutput/<import_result_storage_package>.zip \
-t dw-newide-writer-c パラメーターは設定ファイルのパスを指定します。-f パラメーターは ReaderOutput パッケージのパスを指定します。-o パラメーターは WriterOutput パッケージが保存されるパスを指定します。-t パラメーターは送信プラグインの名前を指定します。
たとえば、プロジェクト A を DataWorks にインポートするには、次のようにします。
sh ./bin/run.sh write \
-c ./conf/projectA_write.json \
-f ./data/1_ReaderOutput/projectA_ReaderOutput.zip \
-o ./data/4_WriterOutput/projectA_WriterOutput.zip \
-t dw-newide-writerインポートツールは実行中にプロセス情報を表示します。エラーがないか確認できます。インポートが完了すると、コマンドラインにインポートの成功と失敗の統計が表示されます。一部のノードのインポートに失敗しても、プロセス全体には影響しないことに注意してください。少数のノードのインポートに失敗した場合は、DataWorks でこれらのノードを手動で変更できます。
2.4 インポート結果の表示
インポートが完了したら、DataWorks で結果を表示できます。各ワークフローがインポートされるときに、インポートプロセスを監視することもできます。問題を発見してインポートを停止する必要がある場合は、`jps` コマンドを実行して `BwmClientApp` を見つけ、`kill -9` コマンドを使用してプロセスを停止できます。
5 Q&A
5.1 開発はソースで進行中です。増分変更を DataWorks に送信するにはどうすればよいですか?
移行ツールは上書きモードで実行されます。ソースから DataWorks に増分変更を送信するには、エクスポート、変換、およびインポートプロセスを再実行できます。ツールは、完全なパスによってワークフローを照合し、新しいワークフローを作成するか、既存のワークフローを更新するかを決定します。変更を正しく移行するには、ワークフローを移動してはなりません。
5.2 開発はソースで進行中であり、DataWorks でもワークフローを変更および管理しています。増分移行によって DataWorks での変更が上書きされますか?
はい、上書きされます。移行ツールは上書きモードで実行されます。移行全体が完了した後にのみ DataWorks で変更を行うことをお勧めします。または、バッチで移行することもできます。ワークフローのバッチが移行され、再度上書きされないことを確認したら、DataWorks でワークフローの変更を開始できます。異なるバッチは互いに影響しません。
5.3 フルパッケージのインポートに時間がかかりすぎます。一部だけをインポートできますか?
はい、できます。インポートパッケージを手動で編集することで、部分的なインポートを実行できます。`data/project/workflow` フォルダで、インポートしたいワークフローのみを保持し、残りを削除します。その後、フォルダを再圧縮してインポートツールを実行できます。相互に依存関係のあるワークフローは一緒にインポートする必要があることに注意してください。そうしないと、ワークフロー間のノードリネージが失われます。