Pub/Sub 到 ClickHouse 模板是一个流式管道,它从 Pub/Sub 订阅中读取 JSON 编码的消息,并将其写入 ClickHouse 表。 无法解析或无法映射到目标 schema 的消息会被路由到死信目标端:ClickHouse 表、Pub/Sub topic,或两者兼有。
管道要求
- 源 Pub/Sub 订阅必须已存在。
- 发布到该订阅的消息必须是有效的 JSON。
- 目标 ClickHouse 表必须已存在,且其列名必须与 JSON 载荷中的字段名匹配。
- Dataflow 工作线程所在的机器必须能够访问 ClickHouse 主机。
- 必须至少提供一个死信目标端 (
clickHouseDeadLetterTable或deadLetterTopic) 。如果两者都提供,失败的消息将同时路由到这两个目标端。 - 设置
clickHouseDeadLetterTable时,死信表必须已在 ClickHouse 中存在,并且其 schema 必须与死信处理中所示一致。 - 设置
deadLetterTopic时,Pub/Sub topic 必须已存在。
模板参数
| 参数名称 | 参数说明 | 必填 | 说明 |
|---|---|---|---|
inputSubscription |
要从中读取消息的 Pub/Sub 订阅。示例:projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>。 |
✅ | 消息必须采用 JSON 编码。 |
clickHouseUrl |
ClickHouse 端点 URL。SSL 连接 (ClickHouse Cloud) 使用 https://,非 SSL 连接使用 http://。示例:https://<HOST>:8443 或 http://<HOST>:8123。 |
✅ | 对于 ClickHouse Cloud,请使用端口 8443 上的 HTTPS 端点。 |
clickHouseDatabase |
目标表所在的 ClickHouse 数据库名称。示例:default。 |
✅ | |
clickHouseTable |
要写入数据的 ClickHouse 表名称。 | ✅ | 运行该管道前,该表必须已存在。 |
clickHouseUsername |
用于向 ClickHouse 进行身份验证的用户名。 | ✅ | |
clickHousePassword |
用于向 ClickHouse 进行身份验证的密码。 | ✅ | |
clickHouseDeadLetterTable |
用于写入失败消息的 ClickHouse 表。示例:my_table_dead_letter。 |
必须提供 clickHouseDeadLetterTable 或 deadLetterTopic 中的至少一个。该表必须已存在,并且具有死信处理中所示的死信 schema。 |
|
deadLetterTopic |
用于发布失败消息的 Pub/Sub topic。示例:projects/<PROJECT_ID>/topics/<TOPIC_NAME>。 |
必须提供 clickHouseDeadLetterTable 或 deadLetterTopic 中的至少一个。失败载荷会发布到该 topic,并将 errorMessage 和 failedAt 设置为消息 attribute。 |
|
windowSeconds |
基于时间的批处理窗口时长 (秒) 。 | 有关它与 batchRowCount 的相互作用,请参见批处理与窗口。如果两者都未设置,则组合模式默认使用 30s 和 1000 行。 |
|
batchRowCount |
在刷写到 ClickHouse 前要累积的行数。 | 有关它与 windowSeconds 的相互作用,请参见批处理与窗口。 |
|
maxInsertBlockSize |
发送到 ClickHouse 的每条 INSERT 语句的最大行数。默认为 1,000,000。 |
一个 ClickHouseIO 选项。 |
|
maxRetries |
ClickHouse 插入失败后的最大重试次数。默认为 5。 |
一个 ClickHouseIO 选项。 |
|
insertDeduplicate |
是否为复制表中的 INSERT 查询启用去重。默认为 true。 |
一个 ClickHouseIO 选项。 |
|
insertQuorum |
对于复制表中的 INSERT 查询,等待指定数量的副本确认写入,并线性化数据写入。0 会禁用 quorum 写入。 |
一个 ClickHouseIO 选项。在默认服务器设置中禁用。 |
|
insertDistributedSync |
如果启用,写入分布式表的 INSERT 查询会等待数据发送到 cluster 中的所有节点。默认为 true。 |
一个 ClickHouseIO 选项。 |
消息格式与 schema 映射
Pub/Sub 消息必须是 JSON 对象,且其顶层字段名必须与目标 ClickHouse 表的列名完全一致。
为将传入消息映射到目标表,管道会在启动时执行以下操作:
- 拉取目标 ClickHouse 表的 schema。
- 根据该 ClickHouse schema 构建 Beam
Rowschema。 - 对每条传入的 Pub/Sub 消息,解析 JSON 载荷,并根据 ClickHouse schema 中定义的字段组装出一行数据。
类型转换
JSON 值会被转换为对应的 ClickHouse 列类型:
| ClickHouse 类型 | 说明 |
|---|---|
Float32 |
通过 Float.valueOf 进行解析。 |
Float64 |
通过 Double.valueOf 进行解析。 |
Date |
解析为 ISO-8601 日期字符串。 |
DateTime |
解析为 ISO-8601 日期时间字符串 (例如 2026-01-15T12:34:56Z) 。 |
Array(T) |
JSON 数组;每个元素都会转换为元素类型 T。空数组或缺失的数组都会生成空数组。 |
Integer types (Int8/Int16/Int32/Int64, UInt8/UInt16/UInt32/UInt64) |
从 JSON 数值或其字符串表示形式中解析。 |
String |
对于文本字段,直接按原样使用;非文本 JSON 节点会被序列化为其 JSON 字符串形式。 |
批处理和窗口化
由于该管道是流式的,传入的行会先在窗口中累积,然后再刷写到 ClickHouse。窗口策略取决于你提供的参数:
windowSeconds |
batchRowCount |
行为 |
|---|---|---|
| 已设置 | 未设置 | 基于时间的固定窗口,窗口长度为 windowSeconds。 |
| 未设置 | 已设置 | 全局窗口,使用计数触发器;每达到 batchRowCount 行时触发一次。 |
| 均已设置 | 均已设置 | 全局窗口,使用组合触发器;哪个条件先满足 (时间 或 行数) 就先触发。 |
| 均未设置 | 均未设置 | 使用默认值的组合模式:30 秒或 1000 行,以先达到者为准。 |
调整这些值可以在延迟和插入效率之间进行权衡。较小的窗口可降低端到端延迟;较大的窗口则会生成更少但更大的 INSERT 批次。
死信处理
在 JSON 解析、schema 映射或类型强制转换过程中失败的消息,会被路由到已配置的死信目标端。必须至少提供 clickHouseDeadLetterTable 或 deadLetterTopic 之一;如果两者都已设置,则失败的消息会同时发送到这两者。
ClickHouse 死信表
设置 clickHouseDeadLetterTable 后,死信表必须已存在,并且采用以下固定 schema:
| 列 | 类型 | 描述 |
|---|---|---|
raw_message |
String |
原始 Pub/Sub 消息载荷,采用 UTF-8 文本格式。 |
error_message |
String |
说明该行失败原因的异常消息。 |
stack_trace |
String |
失败时捕获的完整 Java 堆栈跟踪。 |
failed_at |
DateTime |
该行进入失败状态时的处理时间戳。 |
单节点部署的最小定义:
CREATE TABLE my_table_dead_letter (
raw_message String,
error_message String,
stack_trace String,
failed_at DateTime
) ENGINE = MergeTree()
ORDER BY failed_at;Pub/Sub 死信 topic
设置 deadLetterTopic 后,每条处理失败的消息都会重新发布到该 topic,并附带:
- 载荷:原始消息字节。
- 属性
errorMessage:失败时捕获的异常消息。 - 属性
failedAt:该行失败时的处理时间戳。
这样一来,在底层 schema 或生产者问题解决后,就可以方便地重放失败消息。
运行模板
可在 Google Cloud Console 中使用 Pub/Sub 到 ClickHouse 模板。
登录 Google Cloud Console 并搜索 Dataflow。
-
点击
CREATE JOB FROM TEMPLATE按钮。
-
打开模板表单后,输入作业名称并选择所需的区域。
-
在
Dataflow Template输入框中,输入ClickHouse或Pub/Sub,然后选择Pub/Sub to ClickHouse模板。 -
选中后,表单会展开。请填写:
- Pub/Sub 输入订阅,格式为
projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>。 - ClickHouse 端点 URL——对于 ClickHouse Cloud,请使用
https://<HOST>:8443。 - ClickHouse 数据库、目标表、用户名和密码。
- 至少一个死信目标端:ClickHouse 表或 Pub/Sub topic (或两者都填) 。
- Pub/Sub 输入订阅,格式为
-
你也可以按需自定义批处理 (
windowSeconds、batchRowCount) 以及ClickHouseIO调优参数,详见模板参数一节。
监控作业
前往 Google Cloud Console 中的 Dataflow Jobs 选项卡 以监控作业状态。你可以在其中查看作业详情,包括进度和错误信息:

该模板还会在 PubSubToClickHouse 命名空间下导出以下自定义指标,可在 Dataflow 作业页面查看:
| 指标 | 类型 | 描述 |
|---|---|---|
messages-received |
Counter | 解析步骤接收到的 Pub/Sub 消息总数。 |
rows-parsed-ok |
Counter | 成功转换为一行并路由到主输出的消息数。 |
rows-parse-failed |
Counter | 解析或 schema 映射失败,并被路由到死信的消息数。 |
message-payload-bytes |
Distribution | 传入 Pub/Sub 消息载荷大小的分布,单位为字节。 |
故障排查
超出内存限制 (总量) 错误 (代码 241)
当 ClickHouse 在处理大批次数据时内存耗尽,就会出现此错误。要解决此问题:
- 增加实例资源:将 ClickHouse server 升级到内存更大的实例,以承载数据处理负载。
- 减小批次大小:在 Dataflow job 配置中调小
batchRowCount(和/或maxInsertBlockSize) ,以向 ClickHouse 发送更小的数据块,从而降低每个批次的内存消耗。
所有消息都被发送到死信目标端
最常见的原因是:
- JSON 字段名与 ClickHouse 列名不完全一致 (匹配区分大小写) 。
- 列类型无法根据 JSON 值进行强制转换 (例如,
DateTime列中出现非 ISO-8601 格式的字符串) 。 - 自管道启动以来,目标表的 schema 已发生变化——schema 只会在启动时拉取一次。应用 schema 变更后,请重启该作业。
检查 ClickHouse 死信表中的 error_message 和 stack_trace 列 (或 Pub/Sub 死信消息中的 errorMessage attribute) ,以确定根本原因。
管道已启动,但没有行写入 ClickHouse
- 确认订阅正在接收消息——查看 Dataflow 作业页面上的
messages-received指标。 - 在基于时间的模式下 (仅使用
windowSeconds) ,只有到达窗口边界时才会刷写行。可适当调低windowSeconds,以确认是否发生了刷写。 - 验证 Dataflow 工作线程与 ClickHouse 端点之间的网络连通性 (防火墙、VPC 对等互连或 Private Service Connect) 。
模板源代码
该模板的源代码位于:
GoogleCloudPlatform/DataflowTemplates— 上游的 Google Cloud Platform 代码仓库。ClickHouse/DataflowTemplates— ClickHouse 的 fork。