Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Dataflow Pub/Sub 到 ClickHouse 的模板

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 必须已存在。

模板参数



参数名称 参数说明 必填 说明
inputSubscription 要从中读取消息的 Pub/Sub 订阅。示例:projects/<PROJECT_ID>/subscriptions/<SUBSCRIPTION_NAME> 消息必须采用 JSON 编码。
clickHouseUrl ClickHouse 端点 URL。SSL 连接 (ClickHouse Cloud) 使用 https://,非 SSL 连接使用 http://。示例:https://<HOST>:8443http://<HOST>:8123 对于 ClickHouse Cloud,请使用端口 8443 上的 HTTPS 端点。
clickHouseDatabase 目标表所在的 ClickHouse 数据库名称。示例:default
clickHouseTable 要写入数据的 ClickHouse 表名称。 运行该管道前,该表必须已存在。
clickHouseUsername 用于向 ClickHouse 进行身份验证的用户名。
clickHousePassword 用于向 ClickHouse 进行身份验证的密码。
clickHouseDeadLetterTable 用于写入失败消息的 ClickHouse 表。示例:my_table_dead_letter 必须提供 clickHouseDeadLetterTabledeadLetterTopic 中的至少一个。该表必须已存在,并且具有死信处理中所示的死信 schema。
deadLetterTopic 用于发布失败消息的 Pub/Sub topic。示例:projects/<PROJECT_ID>/topics/<TOPIC_NAME> 必须提供 clickHouseDeadLetterTabledeadLetterTopic 中的至少一个。失败载荷会发布到该 topic,并将 errorMessagefailedAt 设置为消息 attribute。
windowSeconds 基于时间的批处理窗口时长 (秒) 。 有关它与 batchRowCount 的相互作用,请参见批处理与窗口。如果两者都未设置,则组合模式默认使用 30s1000 行。
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 表的列名完全一致。

为将传入消息映射到目标表,管道会在启动时执行以下操作:

  1. 拉取目标 ClickHouse 表的 schema。
  2. 根据该 ClickHouse schema 构建 Beam Row schema。
  3. 对每条传入的 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 映射或类型强制转换过程中失败的消息,会被路由到已配置的死信目标端。必须至少提供 clickHouseDeadLetterTabledeadLetterTopic 之一;如果两者都已设置,则失败的消息会同时发送到这两者。

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。

  1. 点击 CREATE JOB FROM TEMPLATE 按钮。

    Dataflow 控制台
  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 选项卡 以监控作业状态。你可以在其中查看作业详情,包括进度和错误信息:

显示正在运行的从 Pub/Sub 到 ClickHouse 作业的 Dataflow 控制台

该模板还会在 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_messagestack_trace 列 (或 Pub/Sub 死信消息中的 errorMessage attribute) ,以确定根本原因。

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

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

模板源代码

该模板的源代码位于:

Navigation