Skip to main content
在官方 ClickHouse provider 推出之前,大多数 Airflow 用户都是通过社区包 airflow-clickhouse-plugin 来连接 ClickHouse。本指南将介绍如何将现有部署迁移到 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。
这两个包位于不同的 Python 命名空间中,因此可以并存安装,便于你逐个迁移 DAG。当不再有任何代码导入该插件时,即可将其移除:

步骤 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 转发。请按如下方式进行转换: 修改前:
修改后:
在改动 DAG 之前,请先逐一验证已迁移的 connection:

步骤 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.sql module,并去掉 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: 是最简洁的写法。

检查清单

  1. Airflow 版本为 2.11 或更高。
  2. 工作线程可访问 HTTP 端口;HTTP 端点使用的 TLS 证书有效。
  3. 每个 ClickHouse connection:类型为 clickhouse,端口为 8123 或 8443,extras 已按第 2 步完成转换,并通过 airflow connections test 验证。
  4. 已按第 3 步替换 import;common.sql wrapper 上的 ClickHouse 前缀已去除。
  5. operator 和 sensor 上的 clickhouse_conn_id 已重命名为 conn_id,且此前依赖插件默认值的每个 task 均已设置 conn_id。
  6. settings= 已迁移至 hook_params={"session_settings": ...}。
  7. hook.execute("INSERT ... VALUES", rows) 已替换为 bulk_insert_rows;ClickHouseOperator(parameters=rows) 已替换为 SQLInsertRowsOperator。
  8. sensor 的可调用对象已从 result[0][0] 调整为直接使用单元值。
  9. 使用 with_column_types、external_tables、columnar、query_id 和 types_check 的代码已改用 handler 或 get_client() 重写。
  10. 已检查多语句 task 和 INSERT task 所产生 XCom 的下游消费者。
  11. 捕获 clickhouse_driver 异常的代码已更新。
  12. 已卸载 airflow-clickhouse-plugin 和 clickhouse-driver。

使用 AI 编码助手

上述 mapping 有意设计得较为机械,以便编码助手能够将其应用于 DAG repository。以下 prompt 的实际效果不错:
请仔细检查差异内容。connection 的改动、sensor 语义以及 XCom 消费者,正是自动化 rewrite 最容易出错的地方。
最后修改于 2026年9月26日