このドキュメントでは、DMSSqlOperator の設定方法について説明します。
概要
DMSSqlOperator は、DMS で管理されているデータベースインスタンスに対して SQL ステートメントを実行し、その結果を取得します。
パラメーター
パラメーター instance、database、および sql は、Jinja テンプレート をサポートしています。
パラメーター | タイプ | 必須 | 説明 |
instance | string | はい | DMS で管理されるデータベースインスタンスの DBLink 接続名 |
database | string | はい | データベース名。 |
sql | string | はい | 実行する SQL ステートメント。 説明 複数の SQL ステートメントを実行する場合は、セミコロン (;) で区切ってください。 |
csv_null_replace_str | string | いいえ |
|
callback | function | いいえ | SQL 操作の結果を処理するためのコールバック関数。入力パラメーターは PollAsyncSQLExecuteResult です。 説明 このパラメーターは、 |
polling_interval | int | いいえ | 実行結果をポーリングする間隔(秒単位)。デフォルトは 10 です。この値を 0 以下に設定すると、結果を待たずにタスクが送信されます。ステータスチェックには組み込みのリトライ機構が含まれます。 |
show_return_value_in_logs | bool | いいえ | コールバック関数の戻り値をログに出力するかどうかを指定します。デフォルトは |
PollAsyncSQLExecuteResult
パラメーター | タイプ | 説明 |
Status | string | SQL 操作の実行ステータス。有効な値は以下のとおりです。
|
SQLType | string | SQL 操作のタイプ。有効な値は以下のとおりです。
|
ResultType | string | SQL 操作の結果タイプ。 説明 空の値は、SQL 操作が完了していないことを示します。
|
ResultContent | JSON | SQL 操作の結果。 説明
|
例
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」をご参照ください。