> ## Documentation Index
> Fetch the complete documentation index at: https://private-7c7dfe99-vortex-format.mintlify.site/llms.txt
> Use this file to discover all available pages before exploring further.

> 将现有的 Airflow DAG 从社区的 airflow-clickhouse-plugin 迁移到官方的 apache-airflow-providers-clickhousedb provider

# 从 airflow-clickhouse-plugin 迁移到 ClickHouse provider

在[官方 ClickHouse provider](/zh/integrations/connectors/data-ingestion/etl-tools/airflow-and-clickhouse) 推出之前，大多数 Airflow 用户都是通过社区包
[airflow-clickhouse-plugin](https://github.com/bryzgaloff/airflow-clickhouse-plugin) 来连接
ClickHouse。本指南将介绍如何将现有部署迁移到 `apache-airflow-providers-clickhousedb`。若想了解该
provider 本身的工作原理，请参阅[将 Apache Airflow 连接到 ClickHouse](/zh/integrations/connectors/data-ingestion/etl-tools/airflow-and-clickhouse)。

<h2 id="why-a-native-provider">
  为什么要用原生 provider？
</h2>

Airflow 正是通过 provider 与第三方系统集成的。provider 与 Airflow 生态的其他部分一同发布、测试并编写文档，而这个 provider 由 Airflow 社区与 ClickHouse 团队共同维护。它构建在 [clickhouse-connect](/zh/integrations/language-clients/python/index) 之上——这是 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。

<Note>
  在做任何改动之前，请先把本指南完整读一遍。尤其是连接方式的变更，会同时影响所有 DAG。
</Note>

<h2 id="at-a-glance">
  概览
</h2>

| | `airflow-clickhouse-plugin` | `apache-airflow-providers-clickhousedb` |
| - | - | - |
| 导入根路径 | `airflow_clickhouse_plugin` | `airflow.providers.clickhousedb` |
| 驱动 | `clickhouse-driver` | `clickhouse-connect` |
| 协议 / 默认端口 | 原生 TCP，`9000` (启用 TLS 时为 `9440`) | HTTP，`8123` (启用 TLS 时为 `8443`) |
| 连接类型 | 未注册，任意类型均可 | `clickhouse` |
| 连接额外参数 | 原样传递给 `clickhouse_driver.Client` | 固定的键集合，参见[额外连接选项](/zh/integrations/connectors/data-ingestion/etl-tools/airflow-and-clickhouse#extra-connection-options) |
| 操作符与传感器 | `ClickHouseOperator`、`ClickHouseSensor`，以及以 `ClickHouse` 为前缀的 `common.sql` 包装类 | 直接使用 `common.sql` 的操作符与传感器 |
| Hook | `ClickHouseHook` (`BaseHook`) 和 `ClickHouseDbApiHook` (`DbApiHook`) | 仅一个 `ClickHouseHook` (`DbApiHook`) |
| 最低 Airflow 版本 | 2.0 | 2.11 |

<h2 id="step-1-check-prerequisites-and-install">
  步骤 1：检查前置条件并安装
</h2>

该提供商要求 Airflow 2.11 或更高版本，以及 `apache-airflow-providers-common-sql` 1.32.0 或更高版本。如果你使用的是更旧的版本，请先升级 Airflow。

```bash theme={null}
pip install apache-airflow-providers-clickhousedb
```

这两个包位于不同的 Python 命名空间中，因此可以并存安装，便于你逐个迁移 DAG。当不再有任何代码导入该插件时，即可将其移除：

```bash theme={null}
pip uninstall airflow-clickhouse-plugin clickhouse-driver
```

<h2 id="step-2-update-connections">
  步骤 2：更新连接
</h2>

这一步如果跳过，就会出问题。现有的每一个 ClickHouse 连接都指向 native 端口，而 provider 需要的是 HTTP 端口。

| 字段 | 插件 | Provider |
| - | - | - |
| 连接类型 | 任意 (通常是 `sqlite` 或 `generic`) | `clickhouse` |
| 端口 (明文) | `9000` | `8123` |
| 端口 (TLS) | `9440` | `8443` |
| Login | 驱动默认值 `default` | 相同 |
| Schema | 数据库 | 相同 |

如果 ClickHouse 部署在防火墙或负载均衡器之后，请在切换前确认工作线程能够访问 HTTP 端口。确认服务器已启用 HTTP 接口 (服务器配置中的 `http_port` 或 `https_port`) 。[ClickHouse Cloud](/zh/products/cloud/getting-started/intro) 仅在 `8443` 端口上提供 HTTPS。

以 URI 形式存储的连接 (`clickhouse://user:pass@host:9000/db?secure=true`) 已经是 `clickhouse` 连接类型，因为 Airflow 会从 URI 的 scheme 推导出类型。这类连接只需修改端口和 extra 中的键；查询字符串中的值会按 JSON 解析，因此 `secure=true` 仍是布尔值。

<h3 id="connection-extras">
  连接额外参数
</h3>

该插件会将 `extra` 中的每个键直接传给 `clickhouse_driver.Client`，因此连接中可以携带任意 `clickhouse-driver` 的 keyword argument。而 provider 只读取一组固定的键，其余内容通过 `client_kwargs` 转发。请按如下方式进行转换：

| 插件额外参数 (`clickhouse-driver`) | Provider 额外参数 | 说明 |
| - | - | - |
| `secure` | `secure` | 不变。记得同时修改端口。 |
| `verify` | `verify` | 不变。 |
| `settings` | `session_settings` | 内容相同，键名不同。 |
| `compression` | `compress` | 布尔值。不同 driver 支持的算法名称有差异；使用 `true` 最为稳妥。 |
| `connect_timeout` | `connect_timeout` | 不变。 |
| `send_receive_timeout` | `send_receive_timeout` | 不变。 |
| `client_name` | `client_name` | 语义有变化。provider 始终发送 `apache-airflow/<version> apache-airflow-providers-clickhousedb/<version>`，并将你设置的值作为 label 追加在后面。 |
| `ca_certs` | `client_kwargs.ca_cert` | CA bundle 的路径。 |
| `certfile` / `keyfile` | `client_kwargs.client_cert` / `client_kwargs.client_cert_key` | 双向 TLS。 |
| `server_hostname` | `client_kwargs.server_host_name` | TLS SNI override。 |
| `alt_hosts`、`round_robin` | 无对应项 | 客户端只会连接到连接配置中指定的单个 host。如果你原先依赖故障转移列表，请将连接指向 load balancer 或 ClickHouse Cloud 端点。 |
| `sync_request_timeout`、`tcp_keepalive`、`compress_block_size` | 删除 | 仅适用于原生协议。 |
| `ssl_version`、`ciphers`、`use_numpy`、`client_revision`、`settings_is_important`、`opentelemetry_traceparent`、`opentelemetry_tracestate` | 删除 | 在 `clickhouse-connect` 中没有对应项。 |

修改前：

```json theme={null}
{
    "conn_type": "sqlite",
    "host": "ch.example.com",
    "port": 9440,
    "login": "airflow",
    "password": "secret",
    "schema": "analytics",
    "extra": {
        "secure": true,
        "settings": {"max_execution_time": 300},
        "compression": true
    }
}
```

修改后：

```json theme={null}
{
    "conn_type": "clickhouse",
    "host": "ch.example.com",
    "port": 8443,
    "login": "airflow",
    "password": "secret",
    "schema": "analytics",
    "extra": {
        "secure": true,
        "session_settings": {"max_execution_time": 300},
        "compress": true
    }
}
```

在改动 DAG 之前，请先逐一验证已迁移的 connection：

```bash theme={null}
airflow connections test clickhouse_default
```

<h2 id="step-3-replace-imports">
  步骤 3：替换导入
</h2>

| 插件类 | Provider 替代项 |
| - | - |
| `airflow_clickhouse_plugin.hooks.clickhouse.ClickHouseHook` | `airflow.providers.clickhousedb.hooks.clickhouse.ClickHouseHook` |
| `airflow_clickhouse_plugin.hooks.clickhouse_dbapi.ClickHouseDbApiHook` | `airflow.providers.clickhousedb.hooks.clickhouse.ClickHouseHook` |
| `airflow_clickhouse_plugin.operators.clickhouse.ClickHouseOperator` | `airflow.providers.common.sql.operators.sql.SQLExecuteQueryOperator` |
| `airflow_clickhouse_plugin.sensors.clickhouse.ClickHouseSensor` | `airflow.providers.common.sql.sensors.sql.SqlSensor` |
| `airflow_clickhouse_plugin.operators.clickhouse_dbapi.ClickHouseSQLExecuteQueryOperator` | `airflow.providers.common.sql.operators.sql.SQLExecuteQueryOperator` |
| `...clickhouse_dbapi.ClickHouseSQLCheckOperator` | `...common.sql.operators.sql.SQLCheckOperator` |
| `...clickhouse_dbapi.ClickHouseSQLValueCheckOperator` | `...common.sql.operators.sql.SQLValueCheckOperator` |
| `...clickhouse_dbapi.ClickHouseSQLIntervalCheckOperator` | `...common.sql.operators.sql.SQLIntervalCheckOperator` |
| `...clickhouse_dbapi.ClickHouseSQLThresholdCheckOperator` | `...common.sql.operators.sql.SQLThresholdCheckOperator` |
| `...clickhouse_dbapi.ClickHouseSQLColumnCheckOperator` | `...common.sql.operators.sql.SQLColumnCheckOperator` |
| `...clickhouse_dbapi.ClickHouseSQLTableCheckOperator` | `...common.sql.operators.sql.SQLTableCheckOperator` |
| `...clickhouse_dbapi.ClickHouseBranchSQLOperator` | `...common.sql.operators.sql.BranchSQLOperator` |
| `airflow_clickhouse_plugin.sensors.clickhouse_dbapi.ClickHouseSqlSensor` | `airflow.providers.common.sql.sensors.sql.SqlSensor` |

带 `ClickHouse` 前缀的 `common.sql` wrapper 存在的唯一目的，就是注入插件的 hook。该 provider 会注册 `clickhouse` 连接类型，因此不带前缀的 `common.sql` 类可以自行从 connection 中解析出 hook。如果你之前用的是这些 wrapper，迁移通常只需改动 import 行、去掉 `ClickHouse` 前缀，并显式传入 `conn_id` (参见
[步骤 7](#step-7-the-commonsql-wrapper-family)) 。

<h2 id="step-4-clickhouseoperator-to-sqlexecutequeryoperator">
  第 4 步：`ClickHouseOperator` 迁移到 `SQLExecuteQueryOperator`
</h2>

| `ClickHouseOperator` 参数 | `SQLExecuteQueryOperator` 等价写法 | 说明 |
| - | - | - |
| `sql` | `sql` | 可以是字符串、字符串列表或 `.sql` 文件路径。两者都支持模板渲染。 |
| `clickhouse_conn_id` | `conn_id` | 插件默认使用 `clickhouse_default`。`SQLExecuteQueryOperator` 没有默认值：需在每个 task 上传入 `conn_id`，或在 `default_args` 中统一设置一次。 |
| `database` | `database` | 保持不变。 |
| 用于 `SELECT` 的 `parameters` | `parameters` | `%(name)s` 占位符仍然可用。也支持 `{name:Type}` 形式的服务端绑定。 |
| 用于 `INSERT` 的 `parameters` (行列表) | `SQLInsertRowsOperator` | 将这些行 (或某个 XCom) 作为 `rows` 传入，用 `columns` 指定列名，并设置 `insert_args={"executemany": True}`，让 `clickhouse-connect` 使用其 native insert。也可以在 `@task` 中调用 `ClickHouseHook.bulk_insert_rows`，参见[第 5 步](#step-5-clickhousehookexecute-to-dbapihook-methods)。 |
| `settings` | `hook_params={"session_settings": {...}}` | 通过 `hook_params` 支持模板渲染。这些值会合并并覆盖 connection 的 `session_settings`。 |
| `do_xcom_push` | `do_xcom_push` | 标志相同，但在多 statement task 中结果形态不同。插件只推送最后一条 statement 的结果；而传入 statement 列表时，provider 会推送一个列表，每条 statement 对应一个条目。参见[多 statement 结果](#multi-statement-results)。 |
| `query_id` | `hook_params={"session_settings": {"query_id": "..."}}` | 通过 `hook_params` 支持模板渲染。与插件一样，在多 statement task 中每条 statement 都会发送相同的 id。未设置时，`clickhouse-connect` 会为每条 statement 生成唯一 id；可通过 `handler` 从 `cursor.summary[-1]["query_id"]` 读取最后一个 id。 |
| `with_column_types` | `handler` | 传入一个 handler，在读取行的同时读取 `cursor.description`。参见[保留列类型](#keeping-column-types)。 |
| `external_tables` | `ClickHouseHook.get_client()` | 配合 `client.query` 使用 `clickhouse_connect.driver.external.ExternalData`。 |
| `columnar` | `ClickHouseHook.get_client()` | `client.query(...).result_columns`。 |
| `types_check` | 移除 | 仅适用于原生协议。 |

迁移前：

```python theme={null}
from airflow_clickhouse_plugin.operators.clickhouse import ClickHouseOperator

update_income_aggregate = ClickHouseOperator(
    task_id="update_income_aggregate",
    clickhouse_conn_id="clickhouse_test",
    database="default",
    sql=(
        """
        INSERT INTO aggregate
        SELECT eventDt, sum(price * qty) AS income FROM sales
        WHERE eventDt = '{{ ds }}' GROUP BY eventDt
        """,
        """
        SELECT sum(income) FROM aggregate
        WHERE eventDt BETWEEN
            '{{ data_interval_start | ds }}' AND '{{ data_interval_end | ds }}'
        """,
    ),
    settings={"max_execution_time": 600},
)
```

之后:

```python theme={null}
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator

update_income_aggregate = SQLExecuteQueryOperator(
    task_id="update_income_aggregate",
    conn_id="clickhouse_test",
    database="default",
    sql=[
        """
        INSERT INTO aggregate
        SELECT eventDt, sum(price * qty) AS income FROM sales
        WHERE eventDt = '{{ ds }}' GROUP BY eventDt
        """,
        """
        SELECT sum(income) FROM aggregate
        WHERE eventDt BETWEEN
            '{{ data_interval_start | ds }}' AND '{{ data_interval_end | ds }}'
        """,
    ],
    hook_params={"session_settings": {"max_execution_time": 600}},
)
```

<h3 id="multi-statement-results">
  多语句结果
</h3>

该插件只将**最后一条**语句的结果推送到 XCom。当 `sql` 为列表时,`SQLExecuteQueryOperator` 会为每条语句各返回一个结果,
因此上面的示例推送的是 `[[], [(12345.0,)]]`,而插件推送的是 `[(12345.0,)]`。可采用以下任一方式处理:

* 修改下游的 `xcom_pull`,使其只取最后一个元素。
* 将多条语句合并为以 `;` 分隔的单个字符串传入,并设置 `split_statements=True`。此时在默认的
  `return_last=True` 下,该 operator 只会推送最后一条语句的行,与插件行为一致。

```python theme={null}
SQLExecuteQueryOperator(
    task_id="update_income_aggregate",
    conn_id="clickhouse_test",
    sql="""
        INSERT INTO aggregate
        SELECT eventDt, sum(price * qty) AS income FROM sales
        WHERE eventDt = '{{ ds }}' GROUP BY eventDt;
        SELECT sum(income) FROM aggregate WHERE eventDt = '{{ ds }}'
    """,
    split_statements=True,
)
```

通过 `parameters` 插入一组行的 `ClickHouseOperator` 会变为
`SQLInsertRowsOperator`：

```python theme={null}
from airflow.providers.common.sql.operators.sql import SQLInsertRowsOperator

load_rows = SQLInsertRowsOperator(
    task_id="load_rows",
    conn_id="clickhouse_default",
    table_name="some_ch_table",
    columns=["id", "name"],
    rows=extract_task.output,
    insert_args={"executemany": True},
)
```

请始终传入 `columns`；否则该 operator 会通过 SQLAlchemy 去查找该表。

<h3 id="keeping-column-types">
  保留列类型
</h3>

`with_column_types=True` 会返回 `(rows, [(name, type), ...])`。可以借助 handler 实现同样的效果；`clickhouse-connect` 的 cursor 会在 `cursor.description` 中给出 ClickHouse 类型名称：

```python theme={null}
def fetch_with_column_types(cursor):
    return cursor.fetchall(), [(col[0], col[1]) for col in cursor.description]


SQLExecuteQueryOperator(
    task_id="typed_query",
    conn_id="clickhouse_default",
    sql="SELECT id, name FROM users LIMIT 10",
    handler=fetch_with_column_types,
)
```

<h2 id="step-5-clickhousehookexecute-to-dbapihook-methods">
  第 5 步：从 `ClickHouseHook.execute` 迁移到 `DbApiHook` 方法
</h2>

插件的 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` 保持不变。

| 插件 | Provider |
| - | - |
| `hook.execute("SELECT ...")` | `hook.get_records("SELECT ...")` |
| `hook.execute("SELECT ...", params={"d": ds})` | `hook.get_records("SELECT ...", parameters={"d": ds})` |
| `hook.execute("SELECT count() ...")[0][0]` | `hook.get_first("SELECT count() ...")[0]` |
| `hook.execute("INSERT INTO t VALUES", rows)` | `hook.bulk_insert_rows("t", rows, column_names=[...])` |
| `hook.execute(["SET ...", "INSERT ...", "SELECT ..."])` | `hook.run([...], handler=fetch_all_handler, return_last=True)` |
| `hook.execute(..., settings={...})` | `ClickHouseHook(session_settings={...})` |
| `hook.execute(..., external_tables=..., columnar=..., query_id=...)` | `hook.get_client().query(...)` |
| `hook.get_conn()` 返回 `clickhouse_driver.Client` | `hook.get_client()` 返回 `clickhouse_connect` 的 `Client` |

`fetch_all_handler` 及其他 handler 均从
`airflow.providers.common.sql.hooks.handlers` 导入。

插件时代的代码中最常见的 hook 写法是 bulk insert。它无法用普通的 `run` 调用来实现，因为 DB-API 的 cursor 会尝试把这些行格式化拼接进 SQL 字符串。请改用 native insert：

之前：

```python theme={null}
from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook


def sqlite_to_clickhouse():
    records = SqliteHook().get_records("SELECT id, name FROM some_sqlite_table")
    ClickHouseHook().execute("INSERT INTO some_ch_table VALUES", records)
```

修改后：

```python theme={null}
from airflow.providers.clickhousedb.hooks.clickhouse import ClickHouseHook


def sqlite_to_clickhouse():
    records = SqliteHook().get_records("SELECT id, name FROM some_sqlite_table")
    ClickHouseHook().bulk_insert_rows(
        "some_ch_table", records, column_names=["id", "name"], batch_size=100_000
    )
```

`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：

```python theme={null}
from clickhouse_connect.driver.external import ExternalData

hook = ClickHouseHook()
with hook.get_client() as client:
    ext = ExternalData(file_name="ids", structure="id UInt64", data=b"1\n2\n3\n")
    result = client.query(
        "SELECT * FROM events WHERE id IN ids",
        external_data=ext,
        settings={"query_id": "my-traceable-id"},
    )
    columns = result.result_columns
```

<h2 id="step-6-clickhousesensor-to-sqlsensor">
  步骤 6：`ClickHouseSensor` 迁移为 `SqlSensor`
</h2>

这是唯一一处可调用对象输入会发生变化的替换。

| | `ClickHouseSensor` | `SqlSensor` |
| - | - | - |
| 成功可调用对象 | `is_success(result)` | `success(cell)` |
| 失败可调用对象 | `is_failure(result)` | `failure(cell)` |
| 可调用对象接收的内容 | 最后一条语句的完整结果，即由行元组组成的 `list` | 默认为第一行的第一列 (`selector=itemgetter(0)`) |
| 默认成功条件 | `bool(result)`，只要有行返回即为真 | 有行返回时为 `bool(first cell)`，无行返回时为 `False` |
| `sql` | 字符串或语句列表，以及全部 `ClickHouseOperator` 参数 | 单个字符串。使用 `hook_params` 指定数据库和会话设置。 |
| 连接 id | `clickhouse_conn_id`，默认为 `clickhouse_default` | `conn_id`，必填 |
| 未返回任何行 | `is_success([])`，默认情况下为 `False` | `False`，若设置 `fail_on_empty=True` 则报错 |

由于该插件传递的是完整结果集，现有的传感器代码会对其进行索引取值。迁移时请移除这些索引操作：

迁移前：

```python theme={null}
from airflow_clickhouse_plugin.sensors.clickhouse import ClickHouseSensor

ClickHouseSensor(
    task_id="poke_events_count",
    database="monitor",
    sql="SELECT count() FROM warnings WHERE eventDate = '{{ ds }}'",
    is_success=lambda result: result[0][0] > 10000,
)
```

之后:

```python theme={null}
from airflow.providers.common.sql.sensors.sql import SqlSensor

SqlSensor(
    task_id="poke_events_count",
    conn_id="clickhouse_default",
    hook_params={"database": "monitor"},
    sql="SELECT count() FROM warnings WHERE eventDate = '{{ ds }}'",
    success=lambda cnt: cnt > 10000,
)
```

如果可调用对象需要整行数据，可传入 `selector=lambda row: row`。如果它需要所有行，则应将该 check 写成 SQL，让查询返回单个布尔值或计数值。

<h2 id="step-7-the-commonsql-wrapper-family">
  步骤 7：`common.sql` 包装类家族
</h2>

使用 `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 同样使用原生协议。

<h2 id="behavior-differences-to-review">
  需要注意的行为差异
</h2>

即使代码能够编译通过，仍有一些行为在运行时会有所不同。

**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`](/zh/reference/system-tables/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:` 是最简洁的写法。

<h2 id="checklist">
  检查清单
</h2>

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`。

<h2 id="using-an-ai-coding-assistant">
  使用 AI 编码助手
</h2>

上述 mapping 有意设计得较为机械，以便编码助手能够将其应用于 DAG repository。以下 prompt 的实际效果不错：

```text theme={null}
Migrate this repository from airflow-clickhouse-plugin to
apache-airflow-providers-clickhousedb following
https://clickhouse.com/docs/integrations/airflow/migrating-from-airflow-clickhouse-plugin

- Replace every airflow_clickhouse_plugin import per the class mapping table.
- Rename clickhouse_conn_id to conn_id on operators and sensors, not on hooks; add
  conn_id="clickhouse_default" wherever a task relied on the plugin's default.
- Move settings= into hook_params={"session_settings": ...}.
- Replace hook.execute(...) with get_records, get_first, run or bulk_insert_rows
  according to the hook table; never pass a list of rows as parameters. Replace
  ClickHouseOperator(parameters=<rows>) with SQLInsertRowsOperator.
- Rewrite sensor callables to receive the first cell instead of the full result.
- Flag, but do not silently rewrite, any use of with_column_types, external_tables,
  columnar, query_id or types_check, and any XCom consumer of an INSERT task.
- List every Airflow connection that must change port and extras; do not edit
  connections yourself.
Show the diff and a summary of items flagged for human review.
```

请仔细检查差异内容。connection 的改动、sensor 语义以及 XCom 消费者，正是自动化 rewrite 最容易出错的地方。

<h2 id="related-content">
  相关内容
</h2>

* [将 Apache Airflow 连接到 ClickHouse](/zh/integrations/connectors/data-ingestion/etl-tools/airflow-and-clickhouse)
* [`clickhouse-connect` Python 客户端](/zh/integrations/language-clients/python/index)
* [`apache-airflow-providers-clickhousedb` 参考文档](https://airflow.apache.org/docs/apache-airflow-providers-clickhousedb/)
* [GitHub 上的 airflow-clickhouse-plugin](https://github.com/bryzgaloff/airflow-clickhouse-plugin)
