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-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},
)
Navigation