Apache Airflowは、バッチ処理向けのワークフローを開発、スケジューリング、監視するためのオープンソースプラットフォームです。AirflowのすべてのワークフローはPythonコードで定義でき、Webインターフェースでその状態を管理できます。
前提条件
ログインアカウントが、各監視対象プロジェクトで必要なロール(プロジェクト管理者|インスタンス管理者|データ読み書き)に設定されていること。アカウント権限の詳細については、メンバー管理をご参照ください。
Apache Airflowがインストール済みであること。詳細については、Apache Airflow公式サイトをご参照ください。
OB Cloudデータベース接続文字列の取得
インスタンス一覧ページで、対象のトランザクション型インスタンスの情報を展開し、対象テナントで、接続 > 接続文字列を取得 をクリックします。
ポップアップウィンドウで、パブリックネットワークを使用 を選択します。
パブリックIPを使用してデータベースに接続 ページで以下の設定を完了し、接続文字列を生成します:
パラメータ説明IPアドレスの追加 追加をクリックし、エグジットIPをホワイトリストに追加します。 証明書のダウンロード (オプション)証明書をダウンロードをクリックし、CA証明書をダウンロードして認証を完了します。 テナントへの接続 データベース ドロップダウンボックスをクリックした後、+ データベースの作成をクリックし、表示された手順に従ってデータベースの作成を完了します。 アカウント ドロップダウンボックスをクリックした後、+ アカウントの作成をクリックし、表示された手順に従ってアカウントの作成を完了します。 接続方式 MySQL CLIを接続方式として選択します。 注意
アカウント作成後、作成時に生成されたパスワードを適切に保管してください。
AirflowにOB Cloudデータソースを追加して接続する
AirflowのWeb UIを開きます。
Admin -> Connectionsに移動します。
「+」記号をクリックして新しい接続を追加します。
以下のフィールドを入力します:
パラメータ説明Connection Id ob(任意の識別子で構いません)。 Connection Type MySQL Host 接続文字列の -hパラメータから取得したOB Cloudデータベースの接続アドレス。例:t5******.aws-ap-southeast-1.oceanbase.cloud。Schema 接続文字列の -Dパラメータから取得した、アクセスするデータベース名。Login 接続文字列の -uパラメータから取得したアカウント名。例:test。Password 接続文字列の -pパラメータから取得したアカウントのパスワード。Port 接続文字列の -Pパラメータから取得したOB Cloudデータベースの接続ポート。作成が成功すると、Airflowタスク内で
obというConnection Idを参照することで、OceanBaseデータベースにアクセスできるようになります。
Airflowタスクの例
AirflowにOceanBaseデータベースを追加した後、以下の内容を記述することで、AirflowがOceanBaseデータベースからデータを読み取り、出力する処理を実現できます。
Airflowのインストールディレクトリ配下のdagsフォルダに、query.pyファイルを新規作成し、以下の内容を編集します。
from airflow import DAG from airflow.utils.dates import days_ago from airflow.providers.mysql.hooks.mysql import MySqlHook from airflow.operators.python import PythonOperator default_args = { 'owner': 'airflow', 'retries': 0, } def fetch_and_print_data(): hook = MySqlHook(mysql_conn_id='ob') sql = "SELECT * FROM person LIMIT 1;" connection = hook.get_conn() cursor = connection.cursor() cursor.execute(sql) rows = cursor.fetchall() for row in rows: print(row) with DAG( dag_id='sql_query', default_args=default_args, schedule_interval='@daily', start_date=days_ago(1), catchup=False, ) as dag: run_and_print = PythonOperator( task_id='run_and_print', python_callable=fetch_and_print_data, ) run_and_printairflow tasks test sql_query run_and_printを実行すると、personテーブルの最初のデータが出力されます。