管道要求
- 源 Pub/Sub 订阅必须已存在。
- 发布到该订阅的消息必须是有效的 JSON。
- 目标 ClickHouse 表必须已存在,且其列名必须与 JSON 载荷中的字段名匹配。
- Dataflow 工作线程所在的机器必须能够访问 ClickHouse 主机。
- 必须至少提供一个死信目标端 (
clickHouseDeadLetterTable或deadLetterTopic) 。如果两者都提供,失败的消息将同时路由到这两个目标端。 - 设置
clickHouseDeadLetterTable时,死信表必须已在 ClickHouse 中存在,并且其 schema 必须与死信处理中所示一致。 - 设置
deadLetterTopic时,Pub/Sub topic 必须已存在。
模板参数
所有
ClickHouseIO 参数的默认值可在 ClickHouseIO Apache Beam Connector 中找到。消息格式与 schema 映射
- 拉取目标 ClickHouse 表的 schema。
- 根据该 ClickHouse schema 构建 Beam
Rowschema。 - 对每条传入的 Pub/Sub 消息,解析 JSON 载荷,并根据 ClickHouse schema 中定义的字段组装出一行数据。
类型转换
批处理和窗口化
调整这些值可以在延迟和插入效率之间进行权衡。较小的窗口可降低端到端延迟;较大的窗口则会生成更少但更大的
INSERT 批次。
死信处理
clickHouseDeadLetterTable 或 deadLetterTopic 之一;如果两者都已设置,则失败的消息会同时发送到这两者。
ClickHouse 死信表
clickHouseDeadLetterTable 后,死信表必须已存在,并且采用以下固定 schema:
单节点部署的最小定义:
请根据你的部署情况调整引擎和
ORDER BY 子句:复制表请使用 ReplicatedMergeTree,分布式部署请添加 ON CLUSTER,并按需调整分区策略或生存时间 (TTL)。Pub/Sub 死信 topic
deadLetterTopic 后,每条处理失败的消息都会重新发布到该 topic,并附带:
- 载荷:原始消息字节。
- 属性
errorMessage:失败时捕获的异常消息。 - 属性
failedAt:该行失败时的处理时间戳。
运行模板
请务必阅读本文档,尤其是上述各节,以充分了解该模板的配置要求和前置条件。
-
点击
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调优参数,详见模板参数一节。
监控作业
PubSubToClickHouse 命名空间下导出以下自定义指标,可在 Dataflow 作业页面查看:
故障排查
超出内存限制 (总量) 错误 (代码 241)
- 增加实例资源:将 ClickHouse server 升级到内存更大的实例,以承载数据处理负载。
- 减小批次大小:在 Dataflow job 配置中调小
batchRowCount(和/或maxInsertBlockSize) ,以向 ClickHouse 发送更小的数据块,从而降低每个批次的内存消耗。
所有消息都被发送到死信目标端
- JSON 字段名与 ClickHouse 列名不完全一致 (匹配区分大小写) 。
- 列类型无法根据 JSON 值进行强制转换 (例如,
DateTime列中出现非 ISO-8601 格式的字符串) 。 - 自管道启动以来,目标表的 schema 已发生变化——schema 只会在启动时拉取一次。应用 schema 变更后,请重启该作业。
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。