apache-airflow-providers-clickhousedb。若想了解该
provider 本身的工作原理,请参阅将 Apache Airflow 连接到 ClickHouse。
为什么要用原生 provider?
Airflow 正是通过 provider 与第三方系统集成的。provider 与 Airflow 生态的其他部分一同发布、测试并编写文档,而这个 provider 由 Airflow 社区与 ClickHouse 团队共同维护。它构建在 clickhouse-connect 之上——这是 ClickHouse 官方开发并提供支持的 Python 客户端,而非社区维护的 driver。因此,服务端的新特性和修复都能通过受支持的路径传递给 Airflow 用户。迁移到该 provider 后,你将获得一个有官方归属的软件包、标准的common.sql operator 和 sensor,以及一种和其他数据库一样出现在 Airflow UI 中的连接类型。
这两个软件包的差异不只是 import 路径不同。插件使用 clickhouse-driver,通过 原生 TCP 协议 与 ClickHouse 通信;而 provider 使用 clickhouse-connect,通过 HTTP(S) 通信,并接入通用的 common.sql operator,而不是自带一套 ClickHouse 专用的 operator。
在做任何改动之前,请先把本指南完整读一遍。尤其是连接方式的变更,会同时影响所有 DAG。
概览
步骤 1:检查前置条件并安装
该提供商要求 Airflow 2.11 或更高版本,以及apache-airflow-providers-common-sql 1.32.0 或更高版本。如果你使用的是更旧的版本,请先升级 Airflow。
步骤 2:更新连接
这一步如果跳过,就会出问题。现有的每一个 ClickHouse 连接都指向 native 端口,而 provider 需要的是 HTTP 端口。
如果 ClickHouse 部署在防火墙或负载均衡器之后,请在切换前确认工作线程能够访问 HTTP 端口。确认服务器已启用 HTTP 接口 (服务器配置中的
http_port 或 https_port) 。ClickHouse Cloud 仅在 8443 端口上提供 HTTPS。
以 URI 形式存储的连接 (clickhouse://user:pass@host:9000/db?secure=true) 已经是 clickhouse 连接类型,因为 Airflow 会从 URI 的 scheme 推导出类型。这类连接只需修改端口和 extra 中的键;查询字符串中的值会按 JSON 解析,因此 secure=true 仍是布尔值。
连接额外参数
该插件会将extra 中的每个键直接传给 clickhouse_driver.Client,因此连接中可以携带任意 clickhouse-driver 的 keyword argument。而 provider 只读取一组固定的键,其余内容通过 client_kwargs 转发。请按如下方式进行转换:
修改前:
步骤 3:替换导入
带
ClickHouse 前缀的 common.sql wrapper 存在的唯一目的,就是注入插件的 hook。该 provider 会注册 clickhouse 连接类型,因此不带前缀的 common.sql 类可以自行从 connection 中解析出 hook。如果你之前用的是这些 wrapper,迁移通常只需改动 import 行、去掉 ClickHouse 前缀,并显式传入 conn_id (参见
步骤 7) 。
第 4 步:ClickHouseOperator 迁移到 SQLExecuteQueryOperator
迁移前:
多语句结果
该插件只将最后一条语句的结果推送到 XCom。当sql 为列表时,SQLExecuteQueryOperator 会为每条语句各返回一个结果,
因此上面的示例推送的是 [[], [(12345.0,)]],而插件推送的是 [(12345.0,)]。可采用以下任一方式处理:
- 修改下游的
xcom_pull,使其只取最后一个元素。 - 将多条语句合并为以
;分隔的单个字符串传入,并设置split_statements=True。此时在默认的return_last=True下,该 operator 只会推送最后一条语句的行,与插件行为一致。
parameters 插入一组行的 ClickHouseOperator 会变为
SQLInsertRowsOperator:
columns;否则该 operator 会通过 SQLAlchemy 去查找该表。
保留列类型
with_column_types=True 会返回 (rows, [(name, type), ...])。可以借助 handler 实现同样的效果;clickhouse-connect 的 cursor 会在 cursor.description 中给出 ClickHouse 类型名称:
第 5 步:从 ClickHouseHook.execute 迁移到 DbApiHook 方法
插件的 hook 只暴露了一个方法 execute,与 clickhouse_driver.Client.execute 一一对应。而 provider 的 hook 是一个 DbApiHook,因此它具备其他所有 SQL provider 都有的标准方法:
run、get_records、get_first、get_pandas_df、get_df、insert_rows 和 test_connection。
构造函数参数 clickhouse_conn_id 和 database 保持不变。
fetch_all_handler 及其他 handler 均从
airflow.providers.common.sql.hooks.handlers 导入。
插件时代的代码中最常见的 hook 写法是 bulk insert。它无法用普通的 run 调用来实现,因为 DB-API 的 cursor 会尝试把这些行格式化拼接进 SQL 字符串。请改用 native insert:
之前:
bulk_insert_rows 需要 column_names。batch_size 为可选参数,可在输入量非常大时限制内存占用。通用的 insert_rows(table, rows, target_fields=[...], executemany=True) 最终同样会走 native insert,但仅在 executemany=True 时如此;默认行为是每行发送一次 HTTP request。
对于 DB-API 层未覆盖的功能,get_client() 会返回依据 Airflow connection 配置好的原始 clickhouse-connect client。它取代了该插件此前暴露的所有 clickhouse-driver 专有 argument:
步骤 6:ClickHouseSensor 迁移为 SqlSensor
这是唯一一处可调用对象输入会发生变化的替换。
由于该插件传递的是完整结果集,现有的传感器代码会对其进行索引取值。迁移时请移除这些索引操作:
迁移前:
selector=lambda row: row。如果它需要所有行,则应将该 check 写成 SQL,让查询返回单个布尔值或计数值。
步骤 7:common.sql 包装类家族
使用 ClickHouseSQLExecuteQueryOperator、ClickHouseSqlSensor 以及其他以
ClickHouse 为前缀的 wrapper 的代码改动最少:
- 将 import 改为
common.sqlmodule,并去掉 class 名称中的ClickHouse前缀。 - 显式传入
conn_id。这些 wrapper 会把缺失或为None的conn_id视为clickhouse_default;而common.sql中的 class 没有默认值,缺少该参数就会失败。 设置default_args={"conn_id": "clickhouse_default"}即可覆盖整个 DAG。 - operator 上的
database=与 sensor 上的hook_params={"schema": ...}仍然有效; 该 provider 的 hook 会将schema视为database的别名。 ClickHouseDbApiHook改为ClickHouseHook。其构造函数参数schema仍可作为别名使用;database是 ClickHouse 的 native 写法,两者同时给出时以database为准。- connection 仍需完成步骤 2 中的端口与 extras 变更,因为这些 wrapper 同样使用原生协议。
需要注意的行为差异
即使代码能够编译通过,仍有一些行为在运行时会有所不同。 INSERT 任务的 XCom 值。 插件会把clickhouse-driver 的返回值原样推送出去,对于带参数的 VALUES 插入来说就是插入的行数。而 provider 对于不返回行的语句会推送一个空结果集。那些从 XCom 中读取行数的下游任务必须改用其他方式获取,例如追加一个 SELECT count() 查询。
会话设置与 SET 语句。 两个包都会在单个连接上执行多语句列表,而 clickhouse-connect 默认为每个客户端创建一个 session,因此列表中较早出现的 SET 语句仍应对后续语句生效。不过仍建议优先使用 session_settings:它更加显式、支持模板化,并且无论服务端或中间 proxy 是否保持 session,行为都一致。如果你的 DAG 依赖 SET,请在自己的环境中确认其行为。
异常。 错误现在是 clickhouse_connect.driver.exceptions.DatabaseError、OperationalError 或 ProgrammingError,而不再是 clickhouse_driver.errors.ServerException 和 NetworkError。请更新 except 子句、on_failure_callback 代码以及依据异常类型判断的重试逻辑。
压缩。 除非设置了 compression,clickhouse-driver 不会启用压缩;而 clickhouse-connect 默认启用 HTTP 响应压缩,并与服务端协商算法。在连接的 extra 中设置 "compress": false 即可恢复原有行为。插件为原生压缩所需的 clickhouse-cityhash 包已不再需要;lz4 作为 clickhouse-connect 的依赖仍会安装。
类型映射。 两个 driver 都返回原生 Python 类型,但二者是不同的代码库。请检查那些依赖精确类型的任务,尤其是涉及带时区的 DateTime64、Decimal、UUID、Nullable 列以及嵌套的 Array 或 Map 值,特别是结果会推送到 XCom 并由下游消费的情况。
system.query_log 中的查询识别。 查询现在通过 HTTP 接口到达,因此其 interface = 2 而不是 1,并且 system.query_log 的 http_user_agent 列会带上 Airflow 和 provider 的版本号,以及 (若已设置) client_name extra。任何基于原生协议或 clickhouse-driver 客户端名称进行过滤的监控都需要更新。不返回行的 SELECT 会产生第二条记录:DB-API 游标会执行 SELECT * FROM (...) LIMIT 0 以获取列的元数据。
HTTP 上的超时。 send_receive_timeout 现在对应 HTTP 读取超时,而 worker 与 ClickHouse 之间的任何 proxy 或负载均衡器都会对请求施加各自的空闲超时。此前在原生协议上运行数分钟的语句可能需要调高这些限制。
连接处理。 hook 会为每次 run 或 get_client 调用创建一个 clickhouse-connect 客户端,这与插件为每次 execute 打开一个新原生连接的做法类似。客户端共享进程级的 HTTP 连接池,因此对 get_client() 返回的客户端调用 close() 只是良好习惯,而非必需;任务进程退出时连接池会被释放。该客户端本身是上下文管理器,因此 with hook.get_client() as client: 是最简洁的写法。
检查清单
- Airflow 版本为 2.11 或更高。
- 工作线程可访问 HTTP 端口;HTTP 端点使用的 TLS 证书有效。
- 每个 ClickHouse connection:类型为
clickhouse,端口为8123或8443,extras 已按第 2 步完成转换,并通过airflow connections test验证。 - 已按第 3 步替换 import;
common.sqlwrapper 上的ClickHouse前缀已去除。 - operator 和 sensor 上的
clickhouse_conn_id已重命名为conn_id,且此前依赖插件默认值的每个 task 均已设置conn_id。 settings=已迁移至hook_params={"session_settings": ...}。hook.execute("INSERT ... VALUES", rows)已替换为bulk_insert_rows;ClickHouseOperator(parameters=rows)已替换为SQLInsertRowsOperator。- sensor 的可调用对象已从
result[0][0]调整为直接使用单元值。 - 使用
with_column_types、external_tables、columnar、query_id和types_check的代码已改用handler或get_client()重写。 - 已检查多语句 task 和 INSERT task 所产生 XCom 的下游消费者。
- 捕获
clickhouse_driver异常的代码已更新。 - 已卸载
airflow-clickhouse-plugin和clickhouse-driver。