Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

高级用法

底层 API

对于无需在 ClickHouse 数据与原生或第三方数据类型及结构之间进行转换的场景,ClickHouse Connect 客户端提供了一些方法,可直接使用 ClickHouse 连接。

客户端 raw_query 方法

Client.raw_query 方法允许通过客户端连接直接使用 ClickHouse HTTP 查询接口。返回值是一个未经处理的 bytes 对象。它通过极简接口提供了便捷封装,支持参数绑定、错误处理、重试和 settings 管理:

参数 类型 默认值 说明
query str 必填 任何有效的 ClickHouse 查询。
parameters dict or sequence None 请参见Parameters 参数
settings dict None 请参见Settings 参数
fmt str None ClickHouse 输出格式。未指定格式时,ClickHouse 使用 TSV。
use_database bool True 包含客户端上配置的数据库。
external_data ExternalData None 外部文件或二进制数据。请参见外部数据
transport_settings dict None 添加到此请求的 HTTP 请求头。

如何处理返回的 bytes 对象由调用方负责。请注意,Client.query_arrow 只是对此方法的一个轻量封装,使用的是 ClickHouse Arrow 输出格式。

Client raw_stream 方法

同步 Client.raw_stream 方法与 raw_query 的 API 相同,但会返回一个由字节块组成的 io.IOBase 流。处理完成后请关闭该流。AsyncClient.raw_stream 需要使用 await 等待,并会返回一个可与 async withasync for 搭配使用的异步 StreamContext

Client raw_insert 方法

Client.raw_insert 方法允许通过客户端连接直接插入 bytes 对象或 bytes 对象生成器。由于它不会对插入载荷进行任何处理,因此性能非常高。该方法提供了用于指定设置和插入格式的选项:

Parameter Type Default Description
table str 必填 简单表名或带数据库限定名的目标表。
column_names Sequence[str] None 插入块的列名。当 fmt 不包含列名时为必填项。
insert_block str, bytes, generator, or BinaryIO 必填 要插入的数据。字符串会使用客户端编码进行编码。
settings dict None 请参见 Settings 参数
fmt str None insert_block 载荷的 ClickHouse 输入格式。未指定格式时使用 Native
compression str None 已应用于 insert_block 的压缩,例如 "gzip""lz4""zstd"
transport_settings dict None 添加到此请求的 HTTP 请求头。

调用方有责任确保 insert_block 采用指定的格式,并使用指定的压缩方法。ClickHouse Connect 会将这些原始插入用于文件上传和 PyArrow 表,并将解析工作交由 ClickHouse 服务器 处理。

将查询结果保存为文件

你可以使用 raw_stream 方法直接将文件从 ClickHouse 流式写入本地文件系统。例如,如果你想将某个查询的结果保存为 CSV 文件,可以使用以下代码片段:

import clickhouse_connect

if __name__ == "__main__":
    client = clickhouse_connect.get_client()
    query = (
        "SELECT number, toString(number) AS number_as_str "
        "FROM system.numbers LIMIT 5"
    )
    stream = client.raw_stream(query=query, fmt="CSVWithNames")
    try:
        with open("output.csv", "wb") as file:
            for chunk in stream:
                file.write(chunk)
    finally:
        stream.close()
        client.close()

上述代码会生成一个 output.csv 文件,内容如下:

"number","number_as_str"
0,"0"
1,"1"
2,"2"
3,"3"
4,"4"

同样,你也可以将数据保存为 TabSeparated 等其他格式。有关所有可用格式选项的概述,请参阅 输入和输出数据格式

多线程、多进程和异步/事件驱动用例

ClickHouse Connect 非常适合用于多线程、多进程以及事件循环驱动/异步应用程序。所有查询和插入处理都在单个线程内完成,因此这些操作通常是线程安全的。 (未来可能会在较低层级为某些操作引入并行处理,以克服单线程带来的性能损耗;但即便如此,线程安全性仍会得到保证。)

由于每个已执行的查询或插入都会分别在各自的 QueryContextInsertContext 对象中维护状态,这些辅助对象本身并不是线程安全的,因此不应在多个处理流之间共享。有关上下文对象的更多讨论,请参见 QueryContextsInsertContexts 部分。

此外,对于同时有两个或更多查询和/或插入“在进行中”的应用程序,还需要注意另外两个方面。第一是与查询/插入关联的 ClickHouse“会话”,第二是 ClickHouse Connect Client 实例使用的 HTTP 连接池。

AsyncClient

ClickHouse Connect 为 asyncio 应用程序提供了一个基于 aiohttp 的原生客户端。使用前,请先安装可选依赖项:

pip install "clickhouse-connect[async]"

等待 get_async_client 完成以创建并初始化客户端。querycommandinsert 等 I/O 方法都是协程:

import asyncio

import clickhouse_connect


async def main():
    async with await clickhouse_connect.get_async_client() as client:
        result = await client.query(
            "SELECT name FROM system.databases ORDER BY name LIMIT 1"
        )
        print(result.result_rows)


asyncio.run(main())

异步 client 与同步 client 遵循相同的 查询、insert、raw、Arrow 和 流式约定。它使用 aiohttp 进行网络 I/O。受 CPU 限制的 Native 格式解析可能会在 executor 中运行,以免阻塞事件循环。

在进入返回的上下文之前,需要先 await 异步流式方法:

async with await client.query_rows_stream(
    "SELECT number FROM numbers(100000)"
) as stream:
    async for row in stream:
        process(row)

与同步工厂不同,get_async_client 默认禁用自动生成会话 ID,以便并发协程共享同一个客户端。只有在你需要会话状态,且能够避免在该会话中并发执行查询时,才应显式传入 session_id 或设置 autogenerate_session_id=True

管理 ClickHouse 会话 ID

每个 ClickHouse 查询都在 ClickHouse “会话”的上下文中执行。目前,会话 主要有两个用途:

  • 将特定的 ClickHouse settings 关联到多个查询 (请参见用户 settings) 。ClickHouse SET 命令用于在用户 会话 范围内更改 settings。
  • 跟踪临时表

默认情况下,同步 Client 会使用自动生成的 会话 ID。因此,SET 语句和临时表会在来自该客户端的请求之间保留。async factory 默认不会生成 会话 ID。ClickHouse 不允许在同一个 会话 中执行并发查询,如果尝试这样做,客户端会引发 ProgrammingError,因此请使用以下模式之一:

  1. 为每个需要 会话 隔离的线程/进程/event handler 创建单独的 Client 实例。这样可以保留每个客户端各自的 会话 状态 (临时表和 SET 值) 。
  2. 如果不需要共享 会话 状态,请在调用 querycommandinsert 时,通过 settings argument 为每个查询指定唯一的 session_id
  3. 通过在创建客户端之前设置 autogenerate_session_id=False,禁用共享客户端上的 会话 (或者直接将其传递给 get_client) 。
import clickhouse_connect

from clickhouse_connect import common

common.set_setting("autogenerate_session_id", False)
client = clickhouse_connect.get_client(
    host="somehost.com",
    username="dbuser",
    password="password",
)

或者,直接将 autogenerate_session_id=False 传递给 get_client(...)

在这种情况下,ClickHouse Connect 不会发送 session_id;服务器不会将不同的请求视为属于同一会话。临时表和会话级设置不会在请求之间保留。

自定义 HTTP 连接池

ClickHouse Connect 使用 urllib3 连接池来管理与服务器之间的底层 HTTP 连接。默认情况下,所有客户端实例共享同一个连接池,这对于绝大多数使用场景已经足够。该默认连接池会为应用程序使用的每个 ClickHouse 服务器维护最多 8 个 HTTP Keep-Alive 连接。

对于大型多线程应用,使用独立的连接池可能更合适。可以通过主 clickhouse_connect.get_client 函数的 pool_mgr 关键字参数提供自定义连接池:

import clickhouse_connect

from clickhouse_connect.driver import httputil

big_pool_mgr = httputil.get_pool_manager(maxsize=16, num_pools=12)

client1 = clickhouse_connect.get_client(pool_mgr=big_pool_mgr)
client2 = clickhouse_connect.get_client(pool_mgr=big_pool_mgr)

客户端既可以共享同一个池管理器,也可以让每个客户端使用单独的管理器。更多详情,请参阅 urllib3 PoolManager 文档

async 客户端使用的是 aiohttp 连接池,而不是 urllib3。可通过 get_async_client 上的 connector_limitconnector_limit_per_hostkeepalive_timeout 对其进行配置。调用 await async_client.close_connections() 会轮换连接池,而不会中断正在处理中的请求。

Navigation