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

Data Management:DMSSqlOperator

最終更新日:May 12, 2026

このドキュメントでは、DMSSqlOperator の設定方法について説明します。

概要

DMSSqlOperator は、DMS で管理されているデータベースインスタンスに対して SQL ステートメントを実行し、その結果を取得します。

パラメーター

説明

パラメーター instancedatabase、および sql は、Jinja テンプレート をサポートしています。

パラメーター

タイプ

必須

説明

instance

string

はい

DMS で管理されるデータベースインスタンスの DBLink 接続名

database

string

はい

データベース名。

sql

string

はい

実行する SQL ステートメント。

説明

複数の SQL ステートメントを実行する場合は、セミコロン (;) で区切ってください。

csv_null_replace_str

string

いいえ

null 値を resultset 内で置き換えるために使用されます。デフォルト値は文字列 "null" です。

callback

function

いいえ

SQL 操作の結果を処理するためのコールバック関数。入力パラメーターは PollAsyncSQLExecuteResult です。

説明

このパラメーターは、polling_interval が 0 より大きい場合にのみ有効になります。

polling_interval

int

いいえ

実行結果をポーリングする間隔(秒単位)。デフォルトは 10 です。この値を 0 以下に設定すると、結果を待たずにタスクが送信されます。ステータスチェックには組み込みのリトライ機構が含まれます。

show_return_value_in_logs

bool

いいえ

コールバック関数の戻り値をログに出力するかどうかを指定します。デフォルトは True です。

PollAsyncSQLExecuteResult

パラメーター

タイプ

説明

Status

string

SQL 操作の実行ステータス。有効な値は以下のとおりです。

  • WAITING

  • RUNNING

  • SUCCESS

  • FAILURE

SQLType

string

SQL 操作のタイプ。有効な値は以下のとおりです。

  • DDL

  • DML

  • DQL

  • UNKNOWN

    説明

    SQL タイプを解析できなかったことを示します。

ResultType

string

SQL 操作の結果タイプ。

説明

空の値は、SQL 操作が完了していないことを示します。

  • PLAINTEXT:ResultContent の内容。通常は DDL および DML 情報です。

  • FILE:結果がファイルであり、OSS からダウンロードする必要があります。デフォルトフォーマットは CSV で、これは通常 DQL 操作に該当します。

ResultContent

JSON

SQL 操作の結果。

説明

resultSetFileLink リンクを使用してファイルをダウンロードする際は、リクエストヘッダー "x-oss-range-behavior:standard" を指定する必要があります。指定しないと、署名検証が失敗します。

{
  "resultSetFileLink" : "https://xxxx", // SQLType が DQL の場合に返されます。
  "resultSetFileType" : "CSV", // DQL 操作の場合に CSV に設定されます。それ以外の場合は、このフィールドは空です。
  "resultSetOption" : {
    ""
  }, 
  "columnMetas": [
    {
      
    }
      
  ],
  "count" : 1,  // DQL 操作の場合に空ではありません。結果セット内の行数を示します。
  "affectRows" : 1, // DML および DQL 操作の場合、影響を受けた行数(例:更新された行数)を示します。
  "errorMessage" : "xxxx" // Status が FAILURE の場合に返されます。
}

説明

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_sql import DMSSqlOperator

import json
import requests

def callback(result):
    print(f"result: {result}")
    if 'Data' in result and 'ResultContent' in result['Data']:
        link = result['Data']['ResultContent']['ResultSetFileLink']
        print(f"get link: {link}")
        http_res = requests.get(link, headers={
            "x-oss-range-behavior": "standard"
        })
        print(f"link res: {http_res.text}")


with DAG(
    "dms_sql_dblink",
    params={
    },
) as dag:

    sql_operator = DMSSqlOperator(
        task_id="sql_test_dblink",
        instance="dblink_90",
        database="student_db",
        sql="show databases;show tables;",
        csv_null_replace_str="null",
        callback=callback,
        polling_interval=5,
        dag=dag
    )

    run_this_last = EmptyOperator(
        task_id="run_this_last",
        dag=dag,
    )

    sql_operator >> run_this_last

if __name__ == "__main__":
    dag.test(
        run_conf={}
    )
説明

すべての DMS Airflow オペレーターは、タスクのキャンセルや自動リトライなどの共通機能をサポートしています。詳細については、「Airflow DMS Operator」をご参照ください。