QueryContexts
ClickHouse Connect 在 QueryContext 中执行标准查询。QueryContext 包含用于针对 ClickHouse 数据库构建查询的关键结构,以及用于将结果处理为 QueryResult 或其他响应数据结构的配置。其中包括查询本身、参数、settings、读取格式以及其他属性。
可以使用客户端的 create_query_context 方法获取 QueryContext。此方法接受与核心查询方法相同的参数。随后,可将此查询上下文作为 context 关键字参数传递给 query、query_df 或 query_np 方法,以替代这些方法中的部分或全部其他参数。请注意,在方法调用中额外指定的参数会覆盖 QueryContext 中的相应属性。
QueryContext 最清晰的用例是使用不同的绑定参数值发送同一个查询。调用 QueryContext.set_parameters 方法并传入一个字典,即可更新所有参数值;也可以调用 QueryContext.set_parameter,传入所需的 key、value 对,来更新任意单个值。
qc = client.create_query_context(
query="SELECT {k:Int32}",
parameters={"k": 13},
)
result = client.query(context=qc)
assert result.first_row == (13,)
qc.set_parameter("k", 79)
result = client.query(context=qc)
assert result.first_row == (79,)请注意,QueryContext 不是线程安全的;但在多线程环境中,可以通过调用 QueryContext.updated_copy 方法获取其副本。
流式查询
ClickHouse Connect 客户端提供了多种以流形式检索数据的方法 (实现为 Python 生成器) :
query_column_block_stream– 使用原生 Python 对象,以列序列形式按块返回查询数据query_row_block_stream– 使用原生 Python 对象,以行块形式返回查询数据query_rows_stream– 使用原生 Python 对象,以行序列形式返回查询数据query_np_stream– 将每个 ClickHouse 查询数据块作为 NumPy array 返回query_df_stream– 将每个 ClickHouse 查询数据块作为 Pandas DataFrame 返回query_arrow_stream– 将查询数据作为 PyArrowRecordBatch对象返回query_df_arrow_stream– 将每个 Arrow 批次作为 Pandas 或 Polars DataFrame 返回,由dataframe_library选择
每个方法都会返回一个 StreamContext,必须使用 with 语句打开。async 客户端的流式方法需要先使用 await 等待,再通过 async with 打开。
数据块
ClickHouse Connect 会将主要 query 方法返回的所有数据,作为从 ClickHouse 服务器 接收的块流进行处理。这些块以自定义的 “Native” 格式在客户端与 ClickHouse 之间传输。“块”本质上就是一系列二进制数据列,其中每一列都包含数量相同、且具有指定数据类型的数据值。 (作为列式数据库,ClickHouse 也以类似的形式存储这些数据。) 查询返回的块大小由两个用户设置控制,这两个设置可在多个级别上设定 (用户 profile、用户、会话或查询) 。它们是:
- max_block_size – 以行为单位的最大块大小。
- preferred_block_size_bytes – 以字节为单位的首选块大小。
无论 preferred_block_size_bytes 如何设置,块都不会超过 max_block_size 行。实际大小可能更小,不应将其视为稳定不变。
使用 Client 的 query_*_stream 方法之一时,结果会按块返回。ClickHouse Connect 一次只加载一个块。这样就能在无需将整个大型结果集全部加载到内存中的情况下处理大量数据。请注意,应用程序应能够处理任意数量的块,并且无法精确控制每个块的大小。
用于慢速处理的 HTTP 数据缓冲区
如果应用程序消费块的速度远慢于服务器生成块的速度,HTTP 连接可能会在处理完成前关闭。如果应用程序有足够的内存来缓冲更多响应数据,请增大通用的 http_buffer_size 设置。默认值为 10 MiB。lz4 和 zstd 的响应字节在此缓冲区中会保持压缩状态,从而提高其有效容量。
StreamContexts
每个 query_*_stream 方法 (例如 query_row_block_stream) 都会返回一个 ClickHouse StreamContext 对象,它结合了 Python 的上下文管理器和生成器。基本用法如下:
with client.query_row_block_stream(
"SELECT pickup, dropoff, pickup_longitude, pickup_latitude FROM taxi_trips"
) as stream:
for block in stream:
for row in block:
process_trip(row)请注意,如果不使用 with 语句就尝试使用 StreamContext,将会引发错误。使用 Python 上下文可以确保 stream (此处指流式 HTTP 响应) 被正确关闭,即使没有消费完所有数据和/或在处理过程中引发了异常也是如此。此外,StreamContext 只能用于消费一次 stream。在 StreamContext 退出后再次尝试使用它,会产生 StreamClosedError。
如果在读取结果时 connection 失败,则会引发 StreamFailureError,而不是静默返回截断的结果。其消息遵循 client 的 show_clickhouse_errors 设置。
你可以使用 StreamContext 的 source 属性来访问父结果对象,其中包含列名和类型。对于大多数 stream,该对象是 QueryResult;但 query_np_stream 和 query_df_stream 方法返回的则是 NumpyResult。
stream 类型
query_column_block_stream 方法将块作为序列返回,其中包含以原生 Python 数据类型存储的列数据。使用上面的 taxi_trips 查询时,返回的数据将是一个列表,其中每个元素都是另一个列表 (或元组) ,包含对应列的全部数据。因此,block[0] 会是一个只包含字符串的元组。面向列的格式最常用于对某一列中的所有值执行聚合操作,例如计算总车费。
query_row_block_stream 方法将块作为行序列返回,类似于传统的关系型数据库。对于出租车行程,返回的数据将是一个列表,其中每个元素都是另一个表示一行数据的列表。因此,block[0] 将按顺序包含第一条出租车行程的所有字段,block[1] 将包含第二条出租车行程的所有字段,依此类推。面向行的结果通常用于显示或转换处理。
query_rows_stream 方法会自动切换到下一个块,并且每次生成一行。它是 query_row_block_stream 的逐行对应方法。
query_np_stream 方法将每个块作为 NumPy array 返回。当所有结果列共享同一种 NumPy dtype 时,该数组是一个形态为 (rows, columns) 的二维数组。混合结果将作为一维 structured array 返回,或使用 object dtype。
query_df_stream 方法将每个 ClickHouse 块作为二维 Pandas DataFrame 返回。下面的示例展示了 StreamContext 对象可以以延迟方式用作上下文 (但只能使用一次) 。
df_stream = client.query_df_stream("SELECT * FROM hits")
column_names = df_stream.source.column_names
with df_stream:
for df in df_stream:
process_dataframe(df)query_df_arrow_stream 方法会将 Arrow 批次转换为 Pandas 或 Polars DataFrame。使用 dataframe_library 选择库,默认为 "pandas"。
最后,query_arrow_stream 会将 ClickHouse ArrowStream 响应封装为 StreamContext。每次迭代都会返回一个 PyArrow RecordBatch。
流式示例
流式传输行数据
import clickhouse_connect
client = clickhouse_connect.get_client()
# Stream large result sets row by row
with client.query_rows_stream("SELECT number, number * 2 as doubled FROM system.numbers LIMIT 100000") as stream:
for row in stream:
print(row) # Process each row
# Output:
# (0, 0)
# (1, 2)
# (2, 4)
# Additional rows follow流式传输行数据块
import clickhouse_connect
client = clickhouse_connect.get_client()
# Stream in blocks of rows (more efficient than row-by-row)
with client.query_row_block_stream("SELECT number, number * 2 FROM system.numbers LIMIT 100000") as stream:
for block in stream:
print(f"Received block with {len(block)} rows")流式传输 Pandas DataFrames
import clickhouse_connect
client = clickhouse_connect.get_client()
# Stream query results as Pandas DataFrames
with client.query_df_stream("SELECT number, toString(number) AS str FROM system.numbers LIMIT 100000") as stream:
for df in stream:
# Process each DataFrame block
print(f"Received DataFrame with {len(df)} rows")
print(df.head(3))流式传输 Arrow 批次数据
import clickhouse_connect
client = clickhouse_connect.get_client()
# Stream query results as Arrow record batches
with client.query_arrow_stream("SELECT * FROM large_table") as stream:
for arrow_batch in stream:
# Process each Arrow batch
print(f"Received Arrow batch with {arrow_batch.num_rows} rows")异步流式返回行
import asyncio
import clickhouse_connect
async def main():
async_client = await clickhouse_connect.get_async_client()
async with await async_client.query_rows_stream(
"SELECT number FROM numbers(100000)"
) as stream:
async for row in stream:
print(row)
asyncio.run(main())NumPy、Pandas 和 Arrow 查询
ClickHouse Connect 提供了专门的查询方法,用于处理 NumPy、Pandas 和 Arrow 数据结构。借助这些方法,您可以直接以这些常用数据格式获取查询结果,无需手动转换。
NumPy 查询
query_np 方法返回的是 NumPy 数组形式的查询结果,而不是 ClickHouse Connect 的 QueryResult。
import clickhouse_connect
client = clickhouse_connect.get_client()
# Query returns a NumPy array
np_array = client.query_np("SELECT number, number * 2 AS doubled FROM system.numbers LIMIT 5")
print(type(np_array))
# Output:
# <class 'numpy.ndarray'>
print(np_array)
# Output:
# [[0 0]
# [1 2]
# [2 4]
# [3 6]
# [4 8]]Pandas 查询
query_df 方法会将查询结果以 Pandas DataFrame 的形式返回,而不是返回 ClickHouse Connect 的 QueryResult。
import clickhouse_connect
client = clickhouse_connect.get_client()
# Query returns a Pandas DataFrame
df = client.query_df("SELECT number, number * 2 AS doubled FROM system.numbers LIMIT 5")
print(type(df))
# Output: <class 'pandas.core.frame.DataFrame'>
print(df)
# Output:
# number doubled
# 0 0 0
# 1 1 2
# 2 2 4
# 3 3 6
# 4 4 8PyArrow 查询
query_arrow 方法会直接使用 ClickHouse 的 Arrow 输出格式返回一个 PyArrow Table。它接受 query、parameters、settings、external_data 和 transport_settings。use_strings 选项用于控制 ClickHouse String 列输出为 Arrow 字符串还是二进制值。
import clickhouse_connect
client = clickhouse_connect.get_client()
# Query returns a PyArrow Table
arrow_table = client.query_arrow("SELECT number, toString(number) AS str FROM system.numbers LIMIT 3")
print(type(arrow_table))
# Output:
# <class 'pyarrow.lib.Table'>
print(arrow_table)
# Output:
# pyarrow.Table
# number: uint64 not null
# str: string not null
# ----
# number: [[0,1,2]]
# str: [["0","1","2"]]基于 Arrow 的 DataFrames
ClickHouse Connect 支持通过 query_df_arrow 和 query_df_arrow_stream 根据 Arrow 结果高效创建 DataFrame。这些方法避免了经由 Python 行对象进行转换,并在目标库支持的情况下复用 Arrow 缓冲区:
query_df_arrow:使用 ClickHouse 的Arrow输出格式 执行查询,并返回一个 DataFrame。dataframe_library="pandas":返回使用pd.ArrowDtype的 Pandas 2.0 或更高版本 DataFrame。dataframe_library="polars":返回通过pl.from_arrow创建的 Polars DataFrame。
query_df_arrow_stream:将 Arrow 批次以 Pandas 或 Polars DataFrame 的形式流式返回。
查询到基于 Arrow 的 DataFrame
import clickhouse_connect
client = clickhouse_connect.get_client()
# Query returns a Pandas DataFrame with Arrow dtypes (requires pandas 2.x)
df = client.query_df_arrow(
"SELECT number, toString(number) AS str FROM system.numbers LIMIT 3",
dataframe_library="pandas"
)
print(df.dtypes)
# Output:
# number uint64[pyarrow]
# str string[pyarrow]
# dtype: object
# Or use Polars
polars_df = client.query_df_arrow(
"SELECT number, toString(number) AS str FROM system.numbers LIMIT 3",
dataframe_library="polars"
)
print(polars_df.dtypes)
# Output:
# [UInt64, String]
# Streaming into batches of DataFrames (polars shown)
with client.query_df_arrow_stream(
"SELECT number, toString(number) AS str FROM system.numbers LIMIT 100000", dataframe_library="polars"
) as stream:
for df_batch in stream:
print(f"Received {type(df_batch)} batch with {len(df_batch)} rows and dtypes: {df_batch.dtypes}")注意事项和局限性
- Arrow schema 由 ClickHouse 控制。对于没有直接 Arrow 表示的类型,可以使用兼容的物理类型返回,包括二进制字段。在执行特定于应用程序的转换之前,请先检查
table.schema或 DataFrame 的 dtype。 - 基于 Arrow 的 Pandas 结果需要 Pandas 2.0 或更高版本。
use_strings用于控制在服务端支持output_format_arrow_string_as_string时,ClickHouseString列使用 Arrow 字符串字段还是二进制字段。- 基于 Arrow 的查询方法暂不支持
tz_mode="schema"。它们会发出警告,并保留 Arrow 响应提供的时区元数据。
读取格式
读取格式控制 query、query_np 和 query_df 返回的值。它们不适用于 raw 或 Arrow 方法,因为这些方法会直接使用服务器输出格式。例如,将 UUID 的读取格式设置为 "string" 会返回 UUID 字符串,而不是 uuid.UUID 对象。
任何格式化函数的“data type”参数都可以包含通配符。format 是一个单独的小写字符串。Array、Nullable 和 LowCardinality 等容器包装器会为其元素类型保留所选 format。
读取格式可以在多个级别设置:
- 全局设置:使用
clickhouse_connect.datatypes.format包中定义的方法。这将控制所有查询中已配置数据类型的格式。
from clickhouse_connect.datatypes.format import set_read_format
# Return both IPv6 and IPv4 values as strings
set_read_format("IPv*", "string")
# Return all Date types as the underlying epoch second or epoch day
set_read_format("Date*", "int")- 对于整个查询,可以使用可选的
query_formats字典参数。在这种情况下,任何属于指定数据类型的列 (或子列) 都会使用已配置的格式。
# Return any UUID column as a string
client.query(
"SELECT user_id, user_uuid, device_uuid FROM users",
query_formats={"UUID": "string"},
)- 对于特定的结果列,可使用可选的
column_formats字典。每个键都是返回的列名。其值可以是格式字符串,也可以是从 ClickHouse 类型名称到格式的嵌套映射,这对 Tuple、Map 及其他容器类型特别有用。
# Return IPv6 values in the `dev_address` column as strings
client.query(
"SELECT device_id, dev_address, gw_address FROM devices",
column_formats={"dev_address": "string"},
)读取格式选项 (Python 类型)
| ClickHouse 类型 | 原生 Python 类型 | 读取格式 | 说明 |
|---|---|---|---|
| Int[8-64], UInt[8-32] | int | string | |
| UInt64 | int | signed | Superset 目前还无法处理较大的无符号 UInt64 值 |
| [U]Int[128,256] | int | string | Pandas 和 NumPy 的 int 值最大仅支持 64 位,因此这些值可以作为字符串返回 |
| BFloat16 | float | - | 所有 Python float 在内部都是 64 位 |
| Float32 | float | string | 所有 Python float 在内部都是 64 位 |
| Float64 | float | string | |
| Decimal | decimal.Decimal | - | |
| String | str | bytes | ClickHouse 的 String 列本身不带编码信息,因此也可用于存储可变长度的二进制数据 |
| FixedString | bytes | string | FixedString 是固定大小的字节数组,但有时也会被当作 Python 字符串处理 |
| Enum[8,16] | str | int | 原生格式返回标签;int 返回底层整数值。 |
| Date | datetime.date | int | 整数格式返回自 1970-01-01 起的天数。 |
| Date32 | datetime.date | int | 整数格式返回更大范围的有符号天数偏移。 |
| DateTime | datetime.datetime | int | 整数格式返回纪元秒。 |
| DateTime64 | datetime.datetime | int | 整数格式返回该列精度下的 ticks。Python datetime 仅支持到微秒。 |
| Time | datetime.timedelta | int, string, time | 整数格式返回秒数。time 格式仅适用于可放入 datetime.time 的值。 |
| Time64 | datetime.timedelta | int, string, time | 整数格式返回该列精度下的 ticks。Python timedelta 仅支持到微秒。 |
| IPv4 | ipaddress.IPv4Address |
string, int | IP 地址可以读取为字符串或整数。 |
| IPv6 | ipaddress.IPv6Address |
string | IP 地址可以读取为字符串;格式正确时,也可以作为 IP 地址插入 |
| Tuple | dict or tuple | tuple, dict, json | 命名元组默认返回字典;未命名元组返回元组。 |
| Map | dict | - | |
| Nested | Sequence[dict] | - | |
| UUID | uuid.UUID | string | UUID 可以读取为符合 RFC 4122 格式的字符串 |
| JSON | dict | string | 默认返回 Python 字典。string 格式会返回 JSON string |
| Variant | object | typed | typed 返回 TypedVariant(value, type_name),从而保留原始成员类型。 |
| Dynamic | object | - | 返回与该值中存储的 ClickHouse 数据类型相对应的 Python 类型 |
| QBit | list[float] | - | 安装 NumPy 后,会自动使用它来加快位转置。 |
外部数据
ClickHouse 查询可以接收任何受支持输入格式的外部数据。客户端会将数据作为请求的一部分发送,查询可将其作为临时外部表引用。请参阅 ClickHouse 外部数据文档。客户端查询方法可通过 external_data 参数接收一个 clickhouse_connect.driver.external.ExternalData 对象。
| Name | Type | Description |
|---|---|---|
| file_path | str | 用于读取外部数据的本地文件路径。必须提供 file_path 或 data 之一 |
| file_name | str | 外部数据“文件”的名称。如未提供,则取自 file_path 中的文件名部分。外部表名为去掉扩展名后的文件名 |
| data | bytes | 以二进制形式提供的外部数据 (而不是从文件中读取) 。必须提供 data 或 file_path 之一 |
| fmt | str | 数据的 ClickHouse 输入格式。默认为 TSV |
| types | str or seq of str | 外部数据中各列的数据类型列表。如果是字符串,各类型之间应以逗号分隔。必须提供 types 或 structure 之一 |
| structure | str or seq of str | 数据中“列名 + 数据类型”的列表 (参见示例) 。必须提供 structure 或 types 之一 |
| mime_type | str | 文件数据的可选 MIME 类型。目前 ClickHouse 会忽略这个 HTTP 子请求头 |
以下示例将一个外部 CSV 文件与存储在服务器上的 directors 表进行 JOIN:
import clickhouse_connect
from clickhouse_connect.driver.external import ExternalData
client = clickhouse_connect.get_client()
ext_data = ExternalData(
file_path="/data/movies.csv",
fmt="CSV",
structure=[
"movie String",
"year UInt16",
"rating Decimal32(3)",
"director String",
],
)
result = client.query(
"SELECT name, avg(rating) "
"FROM directors INNER JOIN movies ON directors.name = movies.director "
"GROUP BY directors.name",
external_data=ext_data,
).result_rows可以使用 add_file 方法向初始的 ExternalData 对象添加额外的外部数据文件;该方法接受的参数与构造函数相同。对于 HTTP,所有外部数据都会作为 multi-part/form-data 文件上传的一部分进行传输。
chDB backend 不支持外部数据。
时区
ClickHouse DateTime 和 DateTime64 值会以基于纪元的数值形式传输。ClickHouse Connect 会根据列元数据、查询 override 以及客户端的时区策略,将它们转换为 Python datetime 对象。
客户端有两个彼此独立的时区选项:
tz_source为没有显式时区元数据的列选择回退时区:"auto"是默认值。如果客户端能够在夏令时切换期间安全解析服务器时区,则使用服务器时区;否则使用本地时区。"server"始终使用服务器时区。"local"始终使用本地进程时区。
tz_mode控制是否包含时区信息:"naive_utc"是默认值。出于向后兼容考虑,UTC 和与 UTC 等效的结果会以不带时区信息的datetime对象形式返回。"aware"会保留 UTCtzinfo,并返回带时区信息的 UTC 值。"schema"仅在列类型声明了时区时返回带时区信息的值,而对未声明时区的DateTime/DateTime64列则返回不带时区信息的值。
对于常规的 "naive_utc" 和 "aware" 查询,生效时区按以下顺序选择:
- 按列指定的
column_tzsoverride。 - ClickHouse 列类型上的时区元数据。
- 查询级别的
query_tzoverride。 - HTTP 响应中返回的时区信息。
- 由
tz_source选择的回退时区。
tz_mode="schema" 会忽略查询时区和回退时区,但显式指定的 column_tzs override 仍然具有优先次序。
result = client.query(
"SELECT "
"toDateTime('2026-01-15 12:00:00', 'UTC') AS utc_time, "
"toDateTime('2026-01-15 12:00:00', 'America/Denver') AS denver_time",
tz_mode="aware",
)
assert result.first_row[0].tzinfo is not None
assert result.first_row[1].tzinfo is not None时区名称通过标准库 zoneinfo 模块解析。Windows 安装会自动获取 tzdata。对于没有 IANA 时区数据库的精简 Linux 容器镜像,请安装 clickhouse-connect[tzdata]。
Pandas 结果会保留每种 ClickHouse 类型的原生精度,例如 DateTime 对应 datetime64[s],DateTime64(3) 对应 datetime64[ms]。基于 Arrow 的 DataFrame 方法 query_df_arrow 和 query_df_arrow_stream 尚未实现 tz_mode="schema",请求该模式时会发出警告。query_arrow 和 query_arrow_stream 会原样返回 Arrow 响应中的时区元数据。