Apache Airflow は、ワークフローをコードとして記述し、スケジュールし、監視するためのオープンソースプラットフォームです。ワークフローは、Python で記述されたタスクの有向非巡回グラフ (DAG) として定義されます。
apache-airflow-providers-clickhousedb プロバイダーは Airflow を ClickHouse に接続し、DAG の一部としてクエリの実行、テーブルの作成、データの読み込みを行えるようにします。HTTP インターフェイス 経由で clickhouse-connect クライアントを使用して接続し、Airflow の共通 SQL フレームワークを通じて ClickHouse を利用できるようにするため、標準の SQLExecuteQueryOperator で DDL、DML、分析クエリを処理でき、ClickHouse 専用のオペレーターは不要です。
プロバイダーをインストールする
Airflow のスケジューラとワーカーが実行される環境に、プロバイダーをインストールします。
pip install apache-airflow-providers-clickhousedbこのプロバイダーは apache-airflow-providers-common-sql と clickhouse-connect に依存しており、これらもあわせてインストールされます。クエリ結果を pandas または polars の DataFrame に渡すには、オプションの extras をインストールしてください:
pip install 'apache-airflow-providers-common-sql[pandas,polars]'ClickHouse 接続を作成する
このプロバイダーは、clickhouse という接続タイプを登録します。Airflow UI の Admin > Connections から接続を作成するか、CLI または環境変数で定義できます。
UI では、接続タイプとして ClickHouse を選択し、各フィールドに入力します。
| フィールド | 説明 | デフォルト |
|---|---|---|
| Host | ClickHouse server の hostname (例: abc123.clickhouse.cloud) |
localhost |
| Port | HTTP(S) のポート | 8123 (平文), 8443 (TLS) |
| Login | ClickHouse の username | default |
| Password | ClickHouse ユーザーの password | (空) |
| Database | 接続の default database。UI では Database と表示されます。URI または JSON で接続を定義する場合、これは schema フィールドに該当します。 |
default |
ClickHouse Cloud または TLS が有効なセルフホスト クラスターでは、Extra フィールドで secure を true に設定し、TLS ポート (8443) を使用します。
追加の接続オプション
このプロバイダーでは、接続フォーム内に専用フィールドとして追加オプションが用意されています。代わりに URI、JSON、または環境変数で接続を定義する場合は、これらを extra JSON オブジェクト内のキーとして指定してください。いずれも任意です。
extra キー |
UI フィールド | デフォルト | 説明 |
|---|---|---|---|
secure |
TLS (HTTPS) を使用 | false |
HTTPS/TLS を有効にします。 |
verify |
SSL 証明書を検証 | true |
secure が true の場合、サーバーの TLS 証明書を検証します。自己署名証明書を使用する場合は false に設定します。 |
connect_timeout |
接続タイムアウト (秒) | 10 |
HTTP connection のタイムアウト時間 (秒) です。 |
send_receive_timeout |
クエリタイムアウト (秒) | 300 |
クエリの読み取り/書き込みのタイムアウト時間 (秒) です。長時間実行される分析クエリでは、この値を増やしてください。 |
compress |
LZ4 圧縮を有効化 | true |
LZ4 による結果の圧縮を有効にします。 |
client_name |
Client Name | (empty) | ClickHouse の User-Agent および system.query_log の client_name カラムで、Airflow のバージョン識別子に付加されるラベルです。 |
session_settings |
セッション設定 (JSON) | (empty) | 接続上のすべてのクエリに適用される ClickHouse セッション設定 です。たとえば {"max_execution_time": 300, "max_threads": 8} などです。 |
client_kwargs |
Client kwargs (JSON) | (empty) | clickhouse_connect.get_client() に渡される追加のキーワード引数です。たとえば http_proxy などがあります。 |
UI を使わずに接続を定義する
環境変数で接続を設定します。URI 形式には、ホスト、認証情報、データベースが含まれます。
export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='clickhouse://default:password@localhost:8123/my_database'URI のすべての部分は URL エンコードする必要があります。TLS、タイムアウト、またはセッション設定については、Extra フィールドを利用できる JSON 形式を使用します。
export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='{
"conn_type": "clickhouse",
"host": "abc123.clickhouse.cloud",
"port": 8443,
"login": "default",
"password": "secret",
"schema": "my_database",
"extra": {
"secure": true,
"session_settings": {
"max_execution_time": 300,
"max_memory_usage": 10000000000
}
}
}'すべてのフックとオペレーターは、別の接続 ID を指定しない限り、接続 ID clickhouse_default を使用します。
SQLExecuteQueryOperatorでクエリを実行する
オペレーターのconn_idをClickHouse接続に設定します。次のDAGは、テーブルを作成し、行を挿入し、それらを読み出してから、テーブルを削除します。
from datetime import datetime
from airflow import DAG
from airflow.providers.common.sql.hooks.sql import fetch_all_handler
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
CLICKHOUSE_CONN_ID = "clickhouse_default"
CLICKHOUSE_TABLE = "airflow_example"
with DAG(
dag_id="example_clickhouse",
start_date=datetime(2021, 1, 1),
default_args={"conn_id": CLICKHOUSE_CONN_ID},
schedule="@once",
catchup=False,
) as dag:
create_table = SQLExecuteQueryOperator(
task_id="create_table",
sql=f"""
CREATE TABLE IF NOT EXISTS {CLICKHOUSE_TABLE} (
id UInt32,
name String,
ts DateTime DEFAULT now()
) ENGINE = MergeTree()
ORDER BY id
""",
)
insert_rows = SQLExecuteQueryOperator(
task_id="insert_rows",
sql=f"""
INSERT INTO {CLICKHOUSE_TABLE} (id, name) VALUES
(1, 'Alice'),
(2, 'Bob'),
(3, 'Charlie')
""",
)
read_rows = SQLExecuteQueryOperator(
task_id="read_rows",
sql=f"SELECT id, name FROM {CLICKHOUSE_TABLE} ORDER BY id",
handler=fetch_all_handler,
)
drop_table = SQLExecuteQueryOperator(
task_id="drop_table",
sql=f"DROP TABLE IF EXISTS {CLICKHOUSE_TABLE}",
)
create_table >> insert_rows >> read_rows >> drop_tableクエリ結果は、デフォルトのhandler (fetch_all_handler) を使って取得されます。結果セット全体以外を返したい場合は、別のハンドラーを渡します。たとえば、最初の1行だけを返すにはfetch_one_handlerを使用します。
タスクごとに異なるデータベースを対象にする
1 つの接続先がクラスターで、各タスクが異なるデータベースにクエリを実行する場合は、接続を個別に作成するのではなく、hook_params でデータベースを上書きします。
read_rows = SQLExecuteQueryOperator(
task_id="read_rows",
conn_id=CLICKHOUSE_CONN_ID,
sql="SELECT count() FROM events",
hook_params={"database": "analytics"},
)フックを直接使用する
SQLオペレーターでは対応しにくい処理 — 一括挿入、ストリーミング、または ClickHouse 固有のクライアント呼び出し — には、Python タスク内で ClickHouseHook を使用します。
フックの bulk_insert_rows メソッドは、clickhouse-connect のネイティブな列指向の挿入パスを使用します。これは、大規模なデータセットを1行ずつ挿入するよりもはるかに高速です。非常に大きな入力でピークメモリを抑えるには、batch_size を設定します:
from airflow.providers.clickhousedb.hooks.clickhouse import ClickHouseHook
hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default")
hook.bulk_insert_rows(
table="events",
rows=[("user1", "click"), ("user2", "view")],
column_names=["user_id", "action"],
batch_size=1000,
)フックからは直接利用できない機能が必要な場合は、get_client() を呼び出して基盤となる clickhouse-connect クライアントにアクセスします:
client = hook.get_client()
total = client.query("SELECT count() FROM events").result_rows[0][0]セッション設定を適用
フックの構築時に、セッション設定 を直接、またはオペレーターの hook_params を通じて渡します。コンストラクターに渡した設定は、接続の Extra フィールドで定義された session_settings の上にマージされ、同じキーがある場合はコンストラクター側の値が優先されます。
hook = ClickHouseHook(
clickhouse_conn_id="clickhouse_default",
session_settings={"max_execution_time": 60, "max_threads": 4},
)