> ## Documentation Index
> Fetch the complete documentation index at: https://clickhouse.com/docs/llms.txt
> Use this file to discover all available pages before exploring further.

> 您可以使用 Google Dataflow 模板将 JSON 消息从 Pub/Sub 流式导入 ClickHouse

# Dataflow Pub/Sub 到 ClickHouse 的模板

export const Image = ({img, alt, size = "lg"}) => {
  const normalizedSize = ["sm", "md", "lg"].includes(size) ? size : "lg";
  return <div className={`ch-image-${normalizedSize}`}>
      <Frame>
        <img src={img} alt={alt} />
      </Frame>
    </div>;
};

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

<div id="pipeline-requirements">
  ## 管道要求
</div>

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

<div id="template-parameters">
  ## 模板参数
</div>

<br />

<br />

| 参数名称                        | 参数说明                                                                                                                            | 必填 | 说明                                                                                                                          |
| --------------------------- | ------------------------------------------------------------------------------------------------------------------------------- | -- | --------------------------------------------------------------------------------------------------------------------------- |
| `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` 中的至少一个。该表必须已存在，并且具有[死信处理](#dead-letter-handling)中所示的死信 schema。         |
| `deadLetterTopic`           | 用于发布失败消息的 Pub/Sub topic。示例：`projects/<PROJECT_ID>/topics/<TOPIC_NAME>`。                                                         |    | 必须提供 `clickHouseDeadLetterTable` 或 `deadLetterTopic` 中的至少一个。失败载荷会发布到该 topic，并将 `errorMessage` 和 `failedAt` 设置为消息 attribute。 |
| `windowSeconds`             | 基于时间的批处理窗口时长 (秒) 。                                                                                                              |    | 有关它与 `batchRowCount` 的相互作用，请参见[批处理与窗口](#batching-and-windowing)。如果两者都未设置，则组合模式默认使用 `30s` 和 `1000` 行。                        |
| `batchRowCount`             | 在刷写到 ClickHouse 前要累积的行数。                                                                                                        |    | 有关它与 `windowSeconds` 的相互作用，请参见[批处理与窗口](#batching-and-windowing)。                                                            |
| `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` 选项。                                                                                                       |

<Note>
  所有 `ClickHouseIO` 参数的默认值可在 [`ClickHouseIO` Apache Beam Connector](/docs/zh/integrations/connectors/data-ingestion/etl-tools/apache-beam#clickhouseiowrite-parameters) 中找到。
</Note>

<div id="message-format-and-schema-mapping">
  ## 消息格式与 schema 映射
</div>

Pub/Sub 消息必须是 JSON 对象，且其顶层字段名必须与目标 ClickHouse 表的列名完全一致。

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

1. 拉取目标 ClickHouse 表的 schema。
2. 根据该 ClickHouse schema 构建 Beam `Row` schema。
3. 对每条传入的 Pub/Sub 消息，解析 JSON 载荷，并根据 ClickHouse schema 中定义的字段组装出一行数据。

<br />

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

<div id="type-conversion">
  ### 类型转换
</div>

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

| ClickHouse 类型                                                                      | 说明                                                 |
| ---------------------------------------------------------------------------------- | -------------------------------------------------- |
| [`Float32`](/docs/zh/reference/data-types/float)                                        | 通过 `Float.valueOf` 进行解析。                           |
| [`Float64`](/docs/zh/reference/data-types/float)                                        | 通过 `Double.valueOf` 进行解析。                          |
| [`Date`](/docs/zh/reference/data-types/date)                                            | 解析为 ISO-8601 日期字符串。                                |
| [`DateTime`](/docs/zh/reference/data-types/datetime)                                    | 解析为 ISO-8601 日期时间字符串 (例如 `2026-01-15T12:34:56Z`) 。 |
| [`Array(T)`](/docs/zh/reference/data-types/array)                                       | JSON 数组；每个元素都会转换为元素类型 `T`。空数组或缺失的数组都会生成空数组。        |
| Integer types (`Int8`/`Int16`/`Int32`/`Int64`, `UInt8`/`UInt16`/`UInt32`/`UInt64`) | 从 JSON 数值或其字符串表示形式中解析。                             |
| [`String`](/docs/zh/reference/data-types/string)                                        | 对于文本字段，直接按原样使用；非文本 JSON 节点会被序列化为其 JSON 字符串形式。      |

<div id="batching-and-windowing">
  ## 批处理和窗口化
</div>

由于该管道是流式的，传入的行会先在窗口中累积，然后再刷写到 ClickHouse。窗口策略取决于你提供的参数：

| `windowSeconds` | `batchRowCount` | 行为                                       |
| --------------- | --------------- | ---------------------------------------- |
| 已设置             | 未设置             | 基于时间的固定窗口，窗口长度为 `windowSeconds`。         |
| 未设置             | 已设置             | 全局窗口，使用计数触发器；每达到 `batchRowCount` 行时触发一次。 |
| 均已设置            | 均已设置            | 全局窗口，使用组合触发器；哪个条件先满足 (时间 **或** 行数) 就先触发。 |
| 均未设置            | 均未设置            | 使用默认值的组合模式：`30` 秒或 `1000` 行，以先达到者为准。     |

调整这些值可以在延迟和插入效率之间进行权衡。较小的窗口可降低端到端延迟；较大的窗口则会生成更少但更大的 `INSERT` 批次。

<div id="dead-letter-handling">
  ## 死信处理
</div>

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

<div id="clickhouse-dead-letter-table">
  ### ClickHouse 死信表
</div>

设置 `clickHouseDeadLetterTable` 后，死信表必须已存在，并且采用以下固定 schema：

| 列               | 类型         | 描述                             |
| --------------- | ---------- | ------------------------------ |
| `raw_message`   | `String`   | 原始 Pub/Sub 消息载荷，采用 UTF-8 文本格式。 |
| `error_message` | `String`   | 说明该行失败原因的异常消息。                 |
| `stack_trace`   | `String`   | 失败时捕获的完整 Java 堆栈跟踪。            |
| `failed_at`     | `DateTime` | 该行进入失败状态时的处理时间戳。               |

单节点部署的最小定义：

```sql theme={null}
CREATE TABLE my_table_dead_letter (
    raw_message   String,
    error_message String,
    stack_trace   String,
    failed_at     DateTime
) ENGINE = MergeTree()
ORDER BY failed_at;
```

<Note>
  请根据你的部署情况调整引擎和 `ORDER BY` 子句：复制表请使用 `ReplicatedMergeTree`，分布式部署请添加 `ON CLUSTER`，并按需调整分区策略或生存时间 (TTL)。
</Note>

<div id="pubsub-dead-letter-topic">
  ### Pub/Sub 死信 topic
</div>

设置 `deadLetterTopic` 后，每条处理失败的消息都会重新发布到该 topic，并附带：

* **载荷**：原始消息字节。
* **属性** `errorMessage`：失败时捕获的异常消息。
* **属性** `failedAt`：该行失败时的处理时间戳。

这样一来，在底层 schema 或生产者问题解决后，就可以方便地重放失败消息。

<div id="running-the-template">
  ## 运行模板
</div>

可在 Google Cloud Console 中使用 Pub/Sub 到 ClickHouse 模板。

<Note>
  请务必阅读本文档，尤其是上述各节，以充分了解该模板的配置要求和前置条件。
</Note>

登录 Google Cloud Console 并搜索 Dataflow。

1. 点击 `CREATE JOB FROM TEMPLATE` 按钮。
   <Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/google-dataflow/create_job_from_template_button.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=ca429a13d8a9e99c43ae477bf14ad1a9" border alt="Dataflow 控制台" width="1872" height="886" data-path="images/integrations/data-ingestion/google-dataflow/create_job_from_template_button.webp" />

2. 打开模板表单后，输入作业名称并选择所需的区域。

3. 在 `Dataflow Template` 输入框中，输入 `ClickHouse` 或 `Pub/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. 你也可以按需自定义批处理 (`windowSeconds`、`batchRowCount`) 以及 `ClickHouseIO` 调优参数，详见[模板参数](#template-parameters)一节。

<div id="monitor-the-job">
  ### 监控作业
</div>

前往 Google Cloud Console 中的 [Dataflow Jobs 选项卡](https://console.cloud.google.com/dataflow/jobs) 以监控作业状态。你可以在其中查看作业详情，包括进度和错误信息：

<Image img="https://mintcdn.com/private-7c7dfe99/pIetLsS_hOGHqoPJ/images/integrations/data-ingestion/google-dataflow/pubsub-inqueue-job.webp?fit=max&auto=format&n=pIetLsS_hOGHqoPJ&q=85&s=c5922b3ad406648be710f93d856f5fe8" size="lg" border alt="显示正在运行的从 Pub/Sub 到 ClickHouse 作业的 Dataflow 控制台" width="3016" height="470" data-path="images/integrations/data-ingestion/google-dataflow/pubsub-inqueue-job.webp" />

该模板还会在 `PubSubToClickHouse` 命名空间下导出以下自定义指标，可在 Dataflow 作业页面查看：

| 指标                      | 类型           | 描述                           |
| ----------------------- | ------------ | ---------------------------- |
| `messages-received`     | Counter      | 解析步骤接收到的 Pub/Sub 消息总数。       |
| `rows-parsed-ok`        | Counter      | 成功转换为一行并路由到主输出的消息数。          |
| `rows-parse-failed`     | Counter      | 解析或 schema 映射失败，并被路由到死信的消息数。 |
| `message-payload-bytes` | Distribution | 传入 Pub/Sub 消息载荷大小的分布，单位为字节。  |

<div id="troubleshooting">
  ## 故障排查
</div>

<div id="code-241-dbexception-memory-limit-total-exceeded">
  ### 超出内存限制 (总量) 错误 (代码 241)
</div>

当 ClickHouse 在处理大批次数据时内存耗尽，就会出现此错误。要解决此问题：

* 增加实例资源：将 ClickHouse server 升级到内存更大的实例，以承载数据处理负载。
* 减小批次大小：在 Dataflow job 配置中调小 `batchRowCount` (和/或 `maxInsertBlockSize`) ，以向 ClickHouse 发送更小的数据块，从而降低每个批次的内存消耗。

<div id="all-messages-going-to-dlq">
  ### 所有消息都被发送到死信目标端
</div>

最常见的原因是：

* JSON 字段名与 ClickHouse 列名不完全一致 (匹配区分大小写) 。
* 列类型无法根据 JSON 值进行强制转换 (例如，`DateTime` 列中出现非 ISO-8601 格式的字符串) 。
* 自管道启动以来，目标表的 schema 已发生变化——schema 只会在启动时拉取一次。应用 schema 变更后，请重启该作业。

检查 ClickHouse 死信表中的 `error_message` 和 `stack_trace` 列 (或 Pub/Sub 死信消息中的 `errorMessage` attribute) ，以确定根本原因。

<div id="no-rows-arriving">
  ### 管道已启动，但没有行写入 ClickHouse
</div>

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

<div id="template-source-code">
  ## 模板源代码
</div>

该模板的源代码位于：

* [`GoogleCloudPlatform/DataflowTemplates`](https://github.com/GoogleCloudPlatform/DataflowTemplates/tree/main/v2/googlecloud-to-clickhouse) — 上游的 Google Cloud Platform 代码仓库。
* [`ClickHouse/DataflowTemplates`](https://github.com/ClickHouse/DataflowTemplates) — ClickHouse 的 fork。
