Skip to main content
Pub/Sub 到 ClickHouse 模板是一个流式管道,它从 Pub/Sub 订阅中读取 JSON 编码的消息,并将其写入 ClickHouse 表。 无法解析或无法映射到目标 schema 的消息会被路由到死信目标端:ClickHouse 表、Pub/Sub topic,或两者兼有。

管道要求

  • 源 Pub/Sub 订阅必须已存在。
  • 发布到该订阅的消息必须是有效的 JSON。
  • 目标 ClickHouse 表必须已存在,且其列名必须与 JSON 载荷中的字段名匹配。
  • Dataflow 工作线程所在的机器必须能够访问 ClickHouse 主机。
  • 必须至少提供一个死信目标端 (clickHouseDeadLetterTabledeadLetterTopic) 。如果两者都提供,失败的消息将同时路由到这两个目标端。
  • 设置 clickHouseDeadLetterTable 时,死信表必须已在 ClickHouse 中存在,并且其 schema 必须与死信处理中所示一致。
  • 设置 deadLetterTopic 时,Pub/Sub topic 必须已存在。

模板参数



所有 ClickHouseIO 参数的默认值可在 ClickHouseIO Apache Beam Connector 中找到。

消息格式与 schema 映射

Pub/Sub 消息必须是 JSON 对象,且其顶层字段名必须与目标 ClickHouse 表的列名完全一致。 为将传入消息映射到目标表,管道会在启动时执行以下操作:
  1. 拉取目标 ClickHouse 表的 schema。
  2. 根据该 ClickHouse schema 构建 Beam Row schema。
  3. 对每条传入的 Pub/Sub 消息,解析 JSON 载荷,并根据 ClickHouse schema 中定义的字段组装出一行数据。

JSON 字段名必须与 ClickHouse 列名完全一致 (区分大小写) 。消息中与 ClickHouse 列不对应的字段会被忽略。如果某个 ClickHouse 列在 JSON 载荷中没有对应字段,管道会尝试向该列写入 NULL——这仅在该列被声明为 Nullable 时才会成功。无法解析的消息、其值无法强制转换为列类型的消息,或会向非 Nullable 列写入 NULL 的消息,都会被路由到死信目标端。

类型转换

JSON 值会被转换为对应的 ClickHouse 列类型:

批处理和窗口化

由于该管道是流式的,传入的行会先在窗口中累积,然后再刷写到 ClickHouse。窗口策略取决于你提供的参数: 调整这些值可以在延迟和插入效率之间进行权衡。较小的窗口可降低端到端延迟;较大的窗口则会生成更少但更大的 INSERT 批次。

死信处理

在 JSON 解析、schema 映射或类型强制转换过程中失败的消息,会被路由到已配置的死信目标端。必须至少提供 clickHouseDeadLetterTabledeadLetterTopic 之一;如果两者都已设置,则失败的消息会同时发送到这两者。

ClickHouse 死信表

设置 clickHouseDeadLetterTable 后,死信表必须已存在,并且采用以下固定 schema: 单节点部署的最小定义:
请根据你的部署情况调整引擎和 ORDER BY 子句:复制表请使用 ReplicatedMergeTree,分布式部署请添加 ON CLUSTER,并按需调整分区策略或生存时间 (TTL)。

Pub/Sub 死信 topic

设置 deadLetterTopic 后,每条处理失败的消息都会重新发布到该 topic,并附带:
  • 载荷:原始消息字节。
  • 属性 errorMessage:失败时捕获的异常消息。
  • 属性 failedAt:该行失败时的处理时间戳。
这样一来,在底层 schema 或生产者问题解决后,就可以方便地重放失败消息。

运行模板

可在 Google Cloud Console 中使用 Pub/Sub 到 ClickHouse 模板。
请务必阅读本文档,尤其是上述各节,以充分了解该模板的配置要求和前置条件。
登录 Google Cloud Console 并搜索 Dataflow。
  1. 点击 CREATE JOB FROM TEMPLATE 按钮。
  2. 打开模板表单后,输入作业名称并选择所需的区域。
  3. Dataflow Template 输入框中,输入 ClickHousePub/Sub,然后选择 Pub/Sub to ClickHouse 模板。
  4. 选中后,表单会展开。请填写:
    • Pub/Sub 输入订阅,格式为 projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME>
    • ClickHouse 端点 URL——对于 ClickHouse Cloud,请使用 https://<HOST>:8443
    • ClickHouse 数据库、目标表、用户名和密码。
    • 至少一个死信目标端:ClickHouse 表或 Pub/Sub topic (或两者都填) 。
  5. 你也可以按需自定义批处理 (windowSecondsbatchRowCount) 以及 ClickHouseIO 调优参数,详见模板参数一节。

监控作业

前往 Google Cloud Console 中的 Dataflow Jobs 选项卡 以监控作业状态。你可以在其中查看作业详情,包括进度和错误信息: 该模板还会在 PubSubToClickHouse 命名空间下导出以下自定义指标,可在 Dataflow 作业页面查看:

故障排查

超出内存限制 (总量) 错误 (代码 241)

当 ClickHouse 在处理大批次数据时内存耗尽,就会出现此错误。要解决此问题:
  • 增加实例资源:将 ClickHouse server 升级到内存更大的实例,以承载数据处理负载。
  • 减小批次大小:在 Dataflow job 配置中调小 batchRowCount (和/或 maxInsertBlockSize) ,以向 ClickHouse 发送更小的数据块,从而降低每个批次的内存消耗。

所有消息都被发送到死信目标端

最常见的原因是:
  • JSON 字段名与 ClickHouse 列名不完全一致 (匹配区分大小写) 。
  • 列类型无法根据 JSON 值进行强制转换 (例如,DateTime 列中出现非 ISO-8601 格式的字符串) 。
  • 自管道启动以来,目标表的 schema 已发生变化——schema 只会在启动时拉取一次。应用 schema 变更后,请重启该作业。
检查 ClickHouse 死信表中的 error_messagestack_trace 列 (或 Pub/Sub 死信消息中的 errorMessage attribute) ,以确定根本原因。

管道已启动,但没有行写入 ClickHouse

  • 确认订阅正在接收消息——查看 Dataflow 作业页面上的 messages-received 指标。
  • 在基于时间的模式下 (仅使用 windowSeconds) ,只有到达窗口边界时才会刷写行。可适当调低 windowSeconds,以确认是否发生了刷写。
  • 验证 Dataflow 工作线程与 ClickHouse 端点之间的网络连通性 (防火墙、VPC 对等互连或 Private Service Connect) 。

模板源代码

该模板的源代码位于:
最后修改于 2026年7月23日