このトピックでは、DTSLakeInjectionOperator の設定について説明します。
概要
このオペレーターは、Data Transmission Service (DTS) を使用して、Data Management Service (DMS) が管理するデータベースから Object Storage Service (OSS) へデータを同期します。
パラメーター
bucket_name、db_list、および reserve パラメーターは Jinja テンプレートをサポートします。
パラメーター | タイプ | 必須 | 説明 |
source_instance | string | はい | ソース DBLink の名前。。 |
source_database | string | はい | ソースデータベースの名前。 |
target_instance | string | はい | ターゲット DBLink の名前。 |
bucket_name | string | はい | OSS バケットの名前。 |
db_list | dict | はい | 同期オブジェクト。詳細については、「同期オブジェクトの説明」をご参照ください。 |
reserve | dict | いいえ 説明 一部のデータベースでは必須です。 | タスクの予約済みパラメーター。詳細については、「予約パラメーターの説明」をご参照ください。 |
polling_interval | int | いいえ | タスクステータスのポーリング間隔 (秒単位) です。デフォルト値は 10 です。0 または負の数に設定した場合、オペレーターは完了を待たずにタスクを送信します。ポーリングプロセスには、組み込みのリトライメカニズムが含まれています。 |
例
task_id と dag は Airflow の特定のパラメーターです。詳細については、「Airflow 公式ドキュメント」をご参照ください。
from airflow import DAG
from airflow.decorators import task
from airflow.models.param import Param
from airflow.operators.empty import EmptyOperator
from airflow.providers.alibabadms.cloud.operators.dms_dts import DTSLakeInjectionOperator
with DAG(
"dms_dts_dblink",
params={
},
) as dag:
dts_operator = DTSLakeInjectionOperator(
task_id="dts_test_dblink",
source_instance='dblink_90',
source_database='student_db',
target_instance='dbl_oss_63',
bucket_name='hansheng-bj',
db_list=json.loads("""
{\"student_db\":{\"name\":\"hansheng_student_db\",\"all\":false,\"Table\":{\"student_info\":{\"name\":\"student_info\",\"all\":true}}}}
"""),
reserve=json.loads("""
{\"fusionOssFileFormat\":\"DELTA\",\"fusionOssFilePath\":\"/student_info.txt\",\"a2aFlag\":\"2.0\",\"autoStartModulesAfterConfig\":\"none\",\"fusionCreatetableAfterCompleted\": \"false\"}
"""),
polling_interval=5,
dag=dag
)
run_this_last = EmptyOperator(
task_id="run_this_last",
dag=dag,
)
dts_operator >> run_this_last
if __name__ == "__main__":
dag.test(
run_conf={}
)すべての DMS Airflow オペレーターは、タスクのキャンセルや自動リトライなどの共通機能をサポートしています。詳細については、「Airflow DMS オペレーター」をご参照ください。