Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Apache Airflow を ClickHouse に接続

ClickHouse対応

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-sqlclickhouse-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 フィールドで securetrue に設定し、TLS ポート (8443) を使用します。

追加の接続オプション

このプロバイダーでは、接続フォーム内に専用フィールドとして追加オプションが用意されています。代わりに URI、JSON、または環境変数で接続を定義する場合は、これらを extra JSON オブジェクト内のキーとして指定してください。いずれも任意です。

extra キー UI フィールド デフォルト 説明
secure TLS (HTTPS) を使用 false HTTPS/TLS を有効にします。
verify SSL 証明書を検証 true securetrue の場合、サーバーの 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_logclient_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},
)
Navigation