apache-airflow-providers-clickhousedb プロバイダーは Airflow を ClickHouse に接続し、DAG の一部としてクエリの実行、テーブルの作成、データの読み込みを行えるようにします。HTTP インターフェイス 経由で clickhouse-connect クライアントを使用して接続し、Airflow の共通 SQL フレームワークを通じて ClickHouse を利用できるようにするため、標準の SQLExecuteQueryOperator で DDL、DML、分析クエリを処理でき、ClickHouse 固有のオペレーターは不要です。
すでにコミュニティ版の
airflow-clickhouse-plugin を使用していますか?こちらは異なるドライバー、プロトコル、ポートを使用します。
既存の DAG と接続をプロバイダーへ移行するには、移行ガイド に従ってください。プロバイダーをインストールする
Airflow のスケジューラとワーカーが実行される環境に、プロバイダーをインストールします。apache-airflow-providers-common-sql と clickhouse-connect に依存しており、これらもあわせてインストールされます。クエリ結果を pandas または polars の DataFrame に渡すには、オプションの extras をインストールしてください:
ClickHouse 接続を作成する
このプロバイダーは、clickhouse という接続タイプを登録します。Airflow UI の Admin > Connections から接続を作成するか、CLI または環境変数で定義できます。
UI では、接続タイプとして ClickHouse を選択し、各フィールドに入力します。
ClickHouse Cloud または TLS が有効なセルフホスト クラスターでは、Extra フィールドで
secure を true に設定し、TLS ポート (8443) を使用します。
追加の接続オプション
このプロバイダーでは、接続フォーム内に専用フィールドとして追加オプションが用意されています。代わりに URI、JSON、または環境変数で接続を定義する場合は、これらをextra JSON オブジェクト内のキーとして指定してください。いずれも任意です。
UI を使わずに接続を定義する
環境変数で接続を設定します。URI 形式には、ホスト、認証情報、データベースが含まれます。clickhouse_default を使用します。
SQLExecuteQueryOperatorでクエリを実行する
オペレーターのconn_idをClickHouse接続に設定します。次のDAGは、テーブルを作成し、行を挿入し、それらを読み出してから、テーブルを削除します。
handler (fetch_all_handler) を使って取得されます。結果セット全体以外を返したい場合は、別のハンドラーを渡します。たとえば、最初の1行だけを返すにはfetch_one_handlerを使用します。
タスクごとに異なるデータベースを対象にする
1 つの接続先がクラスターで、各タスクが異なるデータベースにクエリを実行する場合は、接続を個別に作成するのではなく、hook_params でデータベースを上書きします。
フックを直接使用する
SQLオペレーターでは対応しにくい処理 — 一括挿入、ストリーミング、または ClickHouse 固有のクライアント呼び出し — には、Python タスク内でClickHouseHook を使用します。
フックの bulk_insert_rows メソッドは、clickhouse-connect のネイティブな列指向の挿入パスを使用します。これは、大規模なデータセットを1行ずつ挿入するよりもはるかに高速です。非常に大きな入力でピークメモリを抑えるには、batch_size を設定します:
get_client() を呼び出して基盤となる clickhouse-connect クライアントにアクセスします:
セッション設定を適用
フックの構築時に、セッション設定 を直接、またはオペレーターのhook_params を通じて渡します。コンストラクターに渡した設定は、接続の Extra フィールドで定義された session_settings の上にマージされ、同じキーがある場合はコンストラクター側の値が優先されます。