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

Data Management:DTSLakeInjectionOperator

最終更新日:May 12, 2026

このトピックでは、DTSLakeInjectionOperator の設定について説明します。

概要

このオペレーターは、Data Transmission Service (DTS) を使用して、Data Management Service (DMS) が管理するデータベースから Object Storage Service (OSS) へデータを同期します。

パラメーター

説明

bucket_namedb_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_iddag は 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 オペレーター」をご参照ください。