Daft 写入 ClickHouse:write_clickhouse 连接指南与 DataSink 源码解析
2026/9/17 5:42:21 网站建设 项目流程

Daft 写入 ClickHouse:write_clickhouse 连接指南与 DataSink 源码解析

【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft

本文围绕 Daft 的 ClickHouse 写入能力展开:从df.write_clickhouse()的安装依赖、基础用法、连接参数到client_kwargs/write_kwargs高级配置,完整覆盖 官方 ClickHouse 连接器文档 的全部实操要点,并结合 ClickHouseDataSink 的实现源码,讲清楚写入结果的统计 DataFrame 是如何产生、连接参数如何合并、以及 DataSink 在分布式执行中的调用流程,帮助你在真实分析管线中把 Daft DataFrame 可靠地落库到 ClickHouse。

功能概述

ClickHouse,可以将 DataFrame 直接写入 ClickHouse 表。该 API 是DataFrame类上的一个公开方法,内部通过构造ClickHouseDataSink并调用通用的write_sink()机制完成写入:

# daft/dataframe/dataframe.py(节选) def write_clickhouse( self, table: str, *, host: str, port: int | None = None, user: str | None = None, password: str | None = None, database: str | None = None, client_kwargs: dict[str, Any] | None = None, write_kwargs: dict[str, Any] | None = None, ) -> "DataFrame": from daft.io.clickhouse.clickhouse_data_sink import ClickHouseDataSink sink = ClickHouseDataSink( table, host=host, port=port, user=user, password=password, database=database, client_kwargs=client_kwargs, write_kwargs=write_kwargs, ) return self.write_sink(sink)

从源码结构看,write_clickhouse()本身不做任何网络 I/O,它只负责把参数打包成ClickHouseDataSink,真正的写入由 Daft 的 DataSink 执行框架在任务调度阶段完成。

安装依赖

ClickHouse 支持需要clickhouse-connect包:

pip install clickhouse-connect

需要注意的是,仓库的 pyproject.toml 中为 ClickHouse 定义的依赖组版本约束为clickhouse_connect<1.1.0,即测试环境锁定在 1.1.0 之前的版本。如果你通过 extras 安装,可以使用daft[clickhouse];自行安装时建议与这一约束保持一致,以确保行为经过仓库测试的验证。

基础用法

创建一个 DataFrame 并写入 ClickHouse 表:

import daft # 创建 DataFrame df = daft.from_pydict({ "id": [1, 2, 3, 4], "name": ["Alice", "Bob", "Charlie", "Diana"], "value": [100.5, 200.3, 150.7, 300.2], }) # 写入 ClickHouse result = df.write_clickhouse( table="my_table", host="localhost", port=8123, user="default", password="", ) result.show()

调用write_clickhouse()后返回的是一个新的 DataFrame,内容为本次写入的统计信息(见下文“输出 Schema”),可以用result.show()直接查看写入是否成功及规模。

连接参数

完整的参数表如下(与 write_clickhouse 方法签名 一致):

参数类型必填说明
tablestr目标 ClickHouse 表名
hoststrClickHouse 服务器主机名
portintClickHouse HTTP 端口(默认 8123)
userstrClickHouse 用户名
passwordstrClickHouse 密码
databasestrClickHouse 数据库名
client_kwargsdict透传给 ClickHouse 客户端构造函数的额外参数
write_kwargsdict透传给写入操作的额外参数

一个容易忽略的细节在 ClickHouseDataSink.init中:portuserpassworddatabase的默认值都是None,只有显式传入时才会放入连接参数字典;port缺省为None时,最终由clickhouse-connect客户端自身的默认值(HTTP 端口 8123)生效。文档中写的“默认 8123”即来源于此。

此外,显式传入的连接参数具有最高优先级——源码中先展开用户提供的client_kwargs,再覆盖以显式参数合并:

# 将用户提供的 client_kwargs 与显式连接参数合并 self._client_kwargs = {**(client_kwargs or {}), **connection_params}

也就是说,如果在client_kwargs里也写了hostport,会以host=/port=的显式入参为准。

输出 Schema:写入统计 DataFrame

写操作返回一个带写入统计信息的 DataFrame:

列名类型说明
total_written_rowsint64写入的总行数
total_written_bytesint64写入的总字节数

这个统计 DataFrame 的产生过程可以在 finalize 方法 中确认:所有分片写入完成后,finalize()会累加每个WriteResultrows_writtenbytes_written,生成一个包含上述两列的单行 MicroPartition。其中“写入字节数”在 write 方法 中是通过df.memory_usage().sum()按 DataFrame 内存占用估算的,因此该数值是客户端侧的估算口径,而非服务端实际落盘字节数。QuerySummary(服务端查询摘要)作为WriteResult.result被携带,但不会直接出现在返回的统计 DataFrame 中。

高级配置

客户端选项(client_kwargs)

client_kwargs会原样透传给clickhouse_connect.get_client()构造函数,可用于 TLS、超时等客户端级配置:

result = df.write_clickhouse( table="my_table", host="localhost", client_kwargs={ "secure": True, "verify": True, "connect_timeout": 30, }, )

写入选项(write_kwargs)

write_kwargs则透传给客户端的insert_df()写入方法,可以控制服务端写入行为,例如启用异步插入:

result = df.write_clickhouse( table="my_table", host="localhost", write_kwargs={ "settings": { "async_insert": 1, "wait_for_async_insert": 1, }, }, )

在源码中可以看到insert_df的调用方式——每个微分区先转为 pandas DataFrame,再连同write_kwargs一起提交:

query_summary = ck_client.insert_df(self._table, df, **self._write_kwargs)

因此 测试用例 中也验证了write_kwargs={"column_names": ["col1"]}这类透传参数会正确地出现在insert_df的调用参数里。

执行机制:DataSink 框架下的写入流程

ClickHouseDataSink实现了 Daft 通用的 DataSink 接口,该接口的执行序列为:

  1. 写入开始时调用 sink 的start()(本 sink 未重写,默认空实现);
  2. DataFrame 执行完毕后,其输出被切分为微分区(MicroPartition);
  3. sink 的write()方法对每个微分区被调用,且可能并行、分布地在多个任务或 worker 上执行;
  4. 所有写入完成后,各分片的WriteResult汇总到单节点;
  5. finalize()被调用,基于全部写入结果产出最终的统计 MicroPartition。

正因为write()会在不同 worker 上执行,源码 中有一个关键设计:clickhouse-connect客户端内部持有 socket 连接,无法跨进程序列化,所以客户端不是在__init__中创建,而是在write()内部通过get_client(**self._client_kwargs)新建,并在finally块中调用close()释放——测试用例 test_client_cleanup 专门验证了即使写入抛出异常,close()也必然被调用。这一点对在 Ray 等分布式 runner 下使用该连接器尤为重要:每个写任务都会拥有独立的客户端连接。

实战场景

分析管线:聚合后落库供仪表盘使用

import daft from daft import col # 处理数据 df = daft.read_parquet("s3://bucket/events/*.parquet") aggregated = df.groupby("event_type").agg( col("value").sum().alias("total_value"), col("user_id").count().alias("event_count"), ) # 将聚合结果写入 ClickHouse 供仪表盘查询 aggregated.write_clickhouse( table="event_aggregates", host="clickhouse.example.com", database="analytics", user="writer", password="secret", )

典型模式是:Daft 负责从对象存储(如 S3)读取 Parquet、做过滤与聚合等重计算,ClickHouse 负责承接低延迟的交互式查询——两者通过write_clickhouse()无缝衔接。

实时数据接入

import daft # 读取流式批次 df = daft.read_json("/data/batch/*.json") # 转换并加载到 ClickHouse df.write_clickhouse( table="events", host="localhost", write_kwargs={ "settings": {"async_insert": 1}, }, )

配合async_insert设置,写入请求可以异步提交以提升吞吐;如果希望写入返回前确认数据已落地,可同时设置wait_for_async_insert: 1

注意事项与限制

  • 目标表必须已存在write_clickhouse()只执行插入(INSERT),不会自动建表,写入前请先在 ClickHouse 中创建好目标表;
  • Schema 对齐:DataFrame 的列名和类型应与目标表 schema 匹配;ClickHouse 在类型兼容时会执行隐式类型转换(type coercion),但超出转换能力的列会写入失败;
  • 版本约束:依赖clickhouse-connect,仓库测试环境锁定clickhouse_connect<1.1.0
  • 参数透传边界client_kwargs影响客户端构造(TLS、超时等),write_kwargs影响单次插入行为(settings、column_names 等),两者的生效位置在 write 方法与客户端构造 中分别体现,配置时不要混用。

相关资源

  • 官方文档:ClickHouse 连接器、连接器总览
  • API 入口:DataFrame.write_clickhouse
  • Sink 实现:ClickHouseDataSink
  • 通用写入框架:DataSink 接口
  • 测试用例:test_clickhouse_writes.py

【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询