底层 API
对于无需在 ClickHouse 数据与原生或第三方数据类型及结构之间进行转换的场景,ClickHouse Connect 客户端提供了一些方法,可直接使用 ClickHouse 连接。客户端 raw_query 方法
Client.raw_query 方法允许通过客户端连接直接使用 ClickHouse HTTP 查询接口。返回值是一个未经处理的 bytes 对象。它通过极简接口提供了便捷封装,支持参数绑定、错误处理、重试和 settings 管理:
如何处理返回的
bytes 对象由调用方负责。请注意,Client.query_arrow 只是对此方法的一个轻量封装,使用的是 ClickHouse Arrow 输出格式。
Client raw_stream 方法
同步 Client.raw_stream 方法与 raw_query 的 API 相同,但会返回一个由字节块组成的 io.IOBase 流。处理完成后请关闭该流。AsyncClient.raw_stream 需要使用 await 等待,并会返回一个可与 async with 和 async for 搭配使用的异步 StreamContext。
Client raw_insert 方法
Client.raw_insert 方法允许通过客户端连接直接插入 bytes 对象或 bytes 对象生成器。由于它不会对插入载荷进行任何处理,因此性能非常高。该方法提供了用于指定设置和插入格式的选项:
调用方有责任确保
insert_block 采用指定的格式,并使用指定的压缩方法。ClickHouse Connect 会将这些原始插入用于文件上传和 PyArrow 表,并将解析工作交由 ClickHouse 服务器 处理。
将查询结果保存为文件
你可以使用raw_stream 方法直接将文件从 ClickHouse 流式写入本地文件系统。例如,如果你想将某个查询的结果保存为 CSV 文件,可以使用以下代码片段:
output.csv 文件,内容如下:
多线程、多进程和异步/事件驱动用例
ClickHouse Connect 非常适合用于多线程、多进程以及事件循环驱动/异步应用程序。所有查询和插入处理都在单个线程内完成,因此这些操作通常是线程安全的。 (未来可能会在较低层级为某些操作引入并行处理,以克服单线程带来的性能损耗;但即便如此,线程安全性仍会得到保证。) 由于每个已执行的查询或插入都会分别在各自的QueryContext 或 InsertContext 对象中维护状态,这些辅助对象本身并不是线程安全的,因此不应在多个处理流之间共享。有关上下文对象的更多讨论,请参见 QueryContexts 和 InsertContexts 部分。
此外,对于同时有两个或更多查询和/或插入“在进行中”的应用程序,还需要注意另外两个方面。第一是与查询/插入关联的 ClickHouse“会话”,第二是 ClickHouse Connect Client 实例使用的 HTTP 连接池。
AsyncClient
ClickHouse Connect 为 asyncio 应用程序提供了一个基于 aiohttp 的原生客户端。使用前,请先安装可选依赖项:get_async_client 完成以创建并初始化客户端。query、command 和 insert 等 I/O 方法都是协程:
await client._initialize(),再发起请求。如果所属的循环已经关闭,请在当前循环中依次调用 await client.close() 和 await client._initialize()。如果清理操作在所属循环关闭之后才开始,aiohttp 仍可能报告存在未关闭的 transport,因此请尽量在迁移之前先关闭客户端。
在进入返回的上下文之前,需要先 await 异步流式方法:
get_async_client 默认禁用自动生成会话 ID,以便并发协程共享同一个客户端。只有在你需要会话状态,且能够避免在该会话中并发执行查询时,才应显式传入 session_id 或设置 autogenerate_session_id=True。
管理 ClickHouse 会话 ID
每个 ClickHouse 查询都在 ClickHouse “会话”的上下文中执行。目前,会话 主要有两个用途:- 将特定的 ClickHouse settings 关联到多个查询 (请参见用户 settings) 。ClickHouse
SET命令用于在用户 会话 范围内更改 settings。 - 跟踪临时表
Client 会使用自动生成的 会话 ID。只有当来自该客户端的请求到达同一个 ClickHouse 服务器进程时,SET 语句和临时表才会在这些请求之间保留。async factory 默认不会生成 会话 ID。命名会话的状态以及同一会话的重叠检查仅在本进程内有效;如果客户端在发送请求之前检测到本地重叠,会引发 ProgrammingError。在 ClickHouse Cloud 或其他采用负载均衡的部署中,不要依赖固定的 session_id 作为分布式状态或分布式互斥锁。如果重叠会带来问题,请在向 ClickHouse 发送请求之前自行对请求进行序列化。请使用以下模式之一:
- 为每个需要 会话 隔离的线程/进程/event handler 创建单独的
Client实例。这样可以保留每个客户端各自的 会话 状态 (临时表和SET值) 。 - 如果不需要共享 会话 状态,请在调用
query、command或insert时,通过settingsargument 为每个查询指定唯一的session_id。 - 通过在创建客户端之前设置
autogenerate_session_id=False,禁用共享客户端上的 会话 (或者直接将其传递给get_client) 。
autogenerate_session_id=False 传递给 get_client(...)。
在这种情况下,ClickHouse Connect 不会发送 session_id;服务器不会将不同的请求视为属于同一会话。临时表和会话级设置不会在请求之间保留。
自定义 HTTP 连接池
ClickHouse Connect 使用urllib3 连接池来管理与服务器之间的底层 HTTP 连接。默认情况下,同一进程中的所有同步客户端实例共享同一个连接池,这对于绝大多数使用场景已经足够。每个 multiprocessing 工作进程都有各自的进程本地默认连接池,并在该进程内创建的所有客户端之间复用。在 fork 之前创建的客户端会保留父进程的连接池,不应在子进程中使用。该默认连接池会为应用程序使用的每个 ClickHouse 服务器维护最多 8 个 HTTP Keep-Alive 连接。
默认的套接字选项会启用 TCP keepalive 和 TCP_NODELAY。套接字发送缓冲区和接收缓冲区的大小由操作系统管理。
对于大型多线程应用,使用独立的连接池可能更合适。可以通过主 clickhouse_connect.get_client 函数的 pool_mgr 关键字参数提供自定义连接池:
urllib3 PoolManager 文档。
要设置套接字选项,请将 socket_options 传给 httputil.get_pool_manager 或 httputil.get_pool_manager_options。这会替换整个默认选项列表,包括 keepalive 选项和 TCP_NODELAY。传入 [] 或 None 则表示不设置任何显式套接字选项。
async 客户端使用的是 aiohttp 连接池,而不是 urllib3。可通过 get_async_client 上的 connector_limit、connector_limit_per_host 和 keepalive_timeout 对其进行配置。调用 await async_client.close_connections() 会轮换连接池,而不会中断正在处理中的请求。
对于异步查询和插入,等待空闲连接池槽位时不设超时。请完整读取或关闭流式响应,以释放其占用的连接池槽位。connect_timeout 在获得可用槽位后才开始计时,涵盖 DNS 解析、TCP 与 TLS 连接建立以及代理协商。send_receive_timeout 用于限制套接字读取的时间。若要为整个操作 (包括等待连接池的时间) 设置截止时间,请使用 asyncio.wait_for,例如 await asyncio.wait_for(client.query("SELECT 13"), timeout=30)。