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's 通用 SQL 框架接入 ClickHouse,因此标准的 SQLExecuteQueryOperator 就能处理 DDL、DML 和分析查询,无需专门的 ClickHouse 特有 operator。

安装提供程序

将该提供程序安装到 Airflow 调度器和工作线程所在的环境中:

pip install apache-airflow-providers-clickhousedb

该提供商依赖 apache-airflow-providers-common-sqlclickhouse-connect,安装时会一并装上它们。若要将查询结果传递给 pandas 或 polars DataFrame,请安装以下可选扩展:

pip install 'apache-airflow-providers-common-sql[pandas,polars]'

创建 ClickHouse 连接

该提供商注册了一种 clickhouse 连接类型。你可以在 Airflow UI 的 Admin > Connections 下创建连接,也可以通过命令行客户端或环境变量定义连接。

在 UI 中,选择 ClickHouse 作为连接类型,并填写以下字段:

字段 说明 默认值
Host ClickHouse server 主机名,例如 abc123.clickhouse.cloud localhost
Port HTTP(S) 端口 8123 (明文) ,8443 (TLS)
Login ClickHouse 用户名 default
Password ClickHouse 用户密码 (空)
Database 该连接的默认 数据库。UI 中此字段显示为 Database;通过 URI 或 JSON 定义连接时,它对应 schema 字段。 default

对于 ClickHouse Cloud 或任何启用 TLS 的 self-hosted cluster,请在 Extra 字段中将 secure 设置为 true,并使用 TLS 端口 (8443) 。

额外连接选项

该提供商在连接表单中将其他选项显示为专用字段。如果你改为通过 URI、JSON 或环境变量来定义连接,则需将这些选项作为 extra JSON 对象中的键传入。以下选项均为可选:

extra key UI 字段 默认值 说明
secure 使用 TLS (HTTPS) false 启用 HTTPS/TLS。
verify 验证 SSL 证书 true securetrue 时,验证 server 的 TLS certificate。对于自签名 certificate,请将其设为 false
connect_timeout 连接超时 (秒) 10 HTTP connection timeout,单位为秒。
send_receive_timeout 查询超时 (秒) 300 查询读写超时,单位为秒。对于长时间运行的分析查询,请适当增大该值。
compress 启用 LZ4 压缩 true 启用 LZ4 结果压缩。
client_name 客户端名称 (空) 追加到 ClickHouse User-Agent 中的 Airflow 版本标识后的标签,同时也会写入 system.query_logclient_name 列。
session_settings 会话设置 (JSON) (空) 应用于此连接上每个查询的 ClickHouse 会话设置,例如 {"max_execution_time": 300, "max_threads": 8}
client_kwargs 客户端 kwargs (JSON) (空) 额外的关键字参数,会转发给 clickhouse_connect.get_client(),例如 http_proxy

无需通过 UI 定义连接

通过环境变量设置连接。URI 格式包含主机、凭据和数据库信息:

export AIRFLOW_CONN_CLICKHOUSE_DEFAULT='clickhouse://default:password@localhost:8123/my_database'

URI 的所有组成部分都必须进行 URL 编码。对于 TLS、超时或会话设置,请使用 JSON 格式,因为这种格式会显示 Extra 字段:

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
        }
    }
}'

所有 Hook 和 Operator 默认使用连接 ID clickhouse_default,除非你另有指定。

使用 SQLExecuteQueryOperator 运行查询

将该 operator 的 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) 拉取。若要返回完整结果集以外的内容,请传入其他 handler,例如只返回第一行的 fetch_one_handler

为每个任务指定不同的数据库

当一个连接指向某个集群,而各个任务需要查询不同的数据库时,请通过 hook_params 覆盖数据库配置,而不要单独创建连接:

read_rows = SQLExecuteQueryOperator(
    task_id="read_rows",
    conn_id=CLICKHOUSE_CONN_ID,
    sql="SELECT count() FROM events",
    hook_params={"database": "analytics"},
)

直接使用 hook

对于不适合通过 SQL Operator 完成的工作——批量插入、流式处理或 ClickHouse 特有的客户端调用——请在 Python 任务中使用 ClickHouseHook

该 hook 的 bulk_insert_rows 方法使用 clickhouse-connect 中原生的列式插入路径;对于大型数据集,这种方式比逐行插入快得多。对于超大输入,可设置 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 客户端,以便使用该 hook 未直接暴露的功能:

client = hook.get_client()
total = client.query("SELECT count() FROM events").result_rows[0][0]

应用会话设置

在构造 hook 时传入会话设置,既可以直接传入,也可以通过 operator 的 hook_params 传入。传给构造函数的设置会与连接的 Extra 字段中定义的任何 session_settings 合并,并在同名键冲突时以构造函数中的值为准:

hook = ClickHouseHook(
    clickhouse_conn_id="clickhouse_default",
    session_settings={"max_execution_time": 60, "max_threads": 4},
)
Navigation