apache-airflow-providers-clickhousedb 提供商可将 Airflow 连接到 ClickHouse,让你能够在 DAG 中运行查询、创建表和加载数据。它通过 HTTP 接口 使用 clickhouse-connect 客户端 进行连接,并通过 Airflow’s 通用 SQL 框架接入 ClickHouse,因此标准的 SQLExecuteQueryOperator 就能处理 DDL、DML 和分析查询,无需专门的 ClickHouse 特有 operator。
安装提供程序
apache-airflow-providers-common-sql 和 clickhouse-connect,安装时会一并装上它们。若要将查询结果传递给 pandas 或 polars DataFrame,请安装以下可选扩展:
创建 ClickHouse 连接
clickhouse 连接类型。你可以在 Airflow UI 的 Admin > Connections 下创建连接,也可以通过命令行客户端或环境变量定义连接。
在 UI 中,选择 ClickHouse 作为连接类型,并填写以下字段:
对于 ClickHouse Cloud 或任何启用 TLS 的 self-hosted cluster,请在 Extra 字段中将
secure 设置为 true,并使用 TLS 端口 (8443) 。
额外连接选项
extra JSON 对象中的键传入。以下选项均为可选:
无需通过 UI 定义连接
clickhouse_default,除非你另有指定。
使用 SQLExecuteQueryOperator 运行查询
conn_id 设置为你的 ClickHouse 连接。以下 DAG 会创建一个表、插入行、读回这些行,并删除该表:
handler (fetch_all_handler) 拉取。若要返回完整结果集以外的内容,请传入其他 handler,例如只返回第一行的 fetch_one_handler。
为每个任务指定不同的数据库
hook_params 覆盖数据库配置,而不要单独创建连接:
直接使用 hook
ClickHouseHook。
该 hook 的 bulk_insert_rows 方法使用 clickhouse-connect 中原生的列式插入路径;对于大型数据集,这种方式比逐行插入快得多。对于超大输入,可设置 batch_size 来限制峰值内存占用:
get_client() 访问底层的 clickhouse-connect 客户端,以便使用该 hook 未直接暴露的功能:
应用会话设置
hook_params 传入。传给构造函数的设置会与连接的 Extra 字段中定义的任何 session_settings 合并,并在同名键冲突时以构造函数中的值为准: