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-sql 和 clickhouse-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 |
当 secure 为 true 时,验证 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_log 的 client_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},
)