apache-airflow-providers-clickhousedb 提供商可将 Airflow 连接到 ClickHouse,让你能够在 DAG 中运行查询、创建表和加载数据。它通过 HTTP 接口 使用 clickhouse-connect 客户端 进行连接,并通过 Airflow’s 通用 SQL 框架接入 ClickHouse,因此标准的 SQLExecuteQueryOperator 就能处理 DDL、DML 和分析查询,无需专门的 ClickHouse 特有 operator。
已经在使用社区版
airflow-clickhouse-plugin?它使用的是不同的 driver、protocol 和端口。请按照
迁移指南 将现有的 DAG 和 connection 迁移到该提供商。安装提供程序
将该提供程序安装到 Airflow 调度器和工作线程所在的环境中: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) 。
额外连接选项
该提供商在连接表单中将其他选项显示为专用字段。如果你改为通过 URI、JSON 或环境变量来定义连接,则需将这些选项作为extra JSON 对象中的键传入。以下选项均为可选:
无需通过 UI 定义连接
通过环境变量设置连接。URI 格式包含主机、凭据和数据库信息:clickhouse_default,除非你另有指定。
使用 SQLExecuteQueryOperator 运行查询
将该 operator 的conn_id 设置为你的 ClickHouse 连接。以下 DAG 会创建一个表、插入行、读回这些行,并删除该表:
handler (fetch_all_handler) 拉取。若要返回完整结果集以外的内容,请传入其他 handler,例如只返回第一行的 fetch_one_handler。
为每个任务指定不同的数据库
当一个连接指向某个集群,而各个任务需要查询不同的数据库时,请通过hook_params 覆盖数据库配置,而不要单独创建连接:
直接使用 hook
对于不适合通过 SQL Operator 完成的工作——批量插入、流式处理或 ClickHouse 特有的客户端调用——请在 Python 任务中使用ClickHouseHook。
该 hook 的 bulk_insert_rows 方法使用 clickhouse-connect 中原生的列式插入路径;对于大型数据集,这种方式比逐行插入快得多。对于超大输入,可设置 batch_size 来限制峰值内存占用:
get_client() 访问底层的 clickhouse-connect 客户端,以便使用该 hook 未直接暴露的功能:
应用会话设置
在构造 hook 时传入会话设置,既可以直接传入,也可以通过 operator 的hook_params 传入。传给构造函数的设置会与连接的 Extra 字段中定义的任何 session_settings 合并,并在同名键冲突时以构造函数中的值为准: