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

# Google BigQuery Sink Connector

> 介绍如何在 AutoMQ Connect 中配置和运行 Google BigQuery Sink Connector，包括前置条件、配置、监控和故障排查。

## 概述

Google BigQuery Sink Connector 将 Kafka Topic 中的 Kafka Connect 记录写入 Google BigQuery。默认情况下，Topic 名称映射为目标表名，`defaultDataset` 作为目标数据集；也可以使用显式 Topic 到表映射，或将 Topic 名称改写为 `dataset:table` 来选择目标。Connector 支持自动建表、Schema 更新、按时间分区、聚类、并发写入和可选的 GCS 批量加载。

## 前置条件

* Google Cloud 项目已启用 BigQuery API，并为 Connector 服务账号授予目标项目、数据集、表以及可选 GCS Bucket 所需权限。

## 授权许可

使用 Apache License 2.0。

## 快速开始

准备好 Connect Cluster、Kafka、Google Cloud 项目和 BigQuery 数据集，并确认 Connect Worker 可以访问 Google Cloud。Connector 的创建、更新和状态查看请参考 AutoMQ 的[管理 Connector](../manage-connectors)。

```properties theme={null}
connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector
topics=orders
project=<gcp-project>
defaultDataset=<dataset>
keyfile=<service-account-json-path>
keySource=FILE
autoCreateTables=true
schemaRetriever=com.wepay.kafka.connect.bigquery.retrieve.IdentitySchemaRetriever
```

将占位符替换为实际资源。`keySource=FILE` 表示 `keyfile` 是服务账号 JSON 文件路径；也可以使用 `keySource=JSON` 传入 JSON 内容，但不要把凭证提交到代码仓库或写入日志。

## 配置

### Kafka Connect 框架

#### `connector.class`

指定要加载的 BigQuery Sink Connector 实现类。

* **类型**：`class`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填，使用 `com.wepay.kafka.connect.bigquery.BigQuerySinkConnector`。

### 消费范围

#### `topics`

要消费的 Topic 列表，使用逗号分隔。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：高
* **有效值 / 注意事项**：必须与 `topics.regex` 二选一。

#### `topics.regex`

使用 Java 正则表达式选择要消费的 Topic。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：必须与 `topics` 二选一。

### 目标与认证

#### `project`

要写入的 Google Cloud 项目。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填。

#### `defaultDataset`

默认 BigQuery 数据集。

* **类型**：`string`
* **默认值**：无
* **重要级别**：高
* **有效值 / 注意事项**：必填；Topic 路由为 `dataset:table` 时可覆盖默认数据集。

#### `keyfile`

服务账号 JSON 文件路径或 JSON 内容。

* **类型**：`password`
* **默认值**：`null`
* **重要级别**：中
* **有效值 / 注意事项**：`keySource=FILE` 时填写文件路径，`keySource=JSON` 时填写 JSON 内容；使用应用默认凭证时不要配置。不要暴露凭证。

#### `keySource`

指定 `keyfile` 是文件路径还是 JSON 内容。

* **类型**：`string`
* **默认值**：`FILE`
* **重要级别**：中
* **有效值 / 注意事项**：`FILE`、`JSON` 或 `APPLICATION_DEFAULT`。使用 `APPLICATION_DEFAULT` 时，Connector 从 Worker 环境获取 Google Cloud 应用默认凭证，并且不应配置 `keyfile`。

### 表与 Schema

#### `autoCreateTables`

目标表不存在时自动创建表。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：高
* **有效值 / 注意事项**：需要 `schemaRetriever`。

#### `schemaRetriever`

用于自动创建表或更新 Schema 的 `SchemaRetriever` 实现类。

* **类型**：`class`
* **默认值**：`com.wepay.kafka.connect.bigquery.retrieve.IdentitySchemaRetriever`
* **重要级别**：中
* **有效值 / 注意事项**：实现类必须有无参构造函数；自定义实现可替换默认实现。

#### `allowNewBigQueryFields`

允许在后续 Schema 更新中添加新字段。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：需要 `schemaRetriever`。

#### `allowBigQueryRequiredFieldRelaxation`

允许将 BigQuery `REQUIRED` 字段改为 `NULLABLE`。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：需要 `schemaRetriever`。

#### `allowSchemaUnionization`

Schema 更新时，将现有表 Schema 与当前批次中的记录 Schema 合并。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：需要 `schemaRetriever`。启用前确认上游 Schema 质量；BigQuery 不支持删除列或任意修改已有列类型。

### 更新与删除

#### `upsertEnabled`

根据 Kafka 记录 Key，通过中间表和周期性 MERGE 更新目标行。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：启用时必须配置 `kafkaKeyFieldName`，并至少启用一种 MERGE 触发方式。

#### `deleteEnabled`

根据 Kafka 记录 Key，将 Tombstone 记录合并为目标表删除操作。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：启用时必须配置 `kafkaKeyFieldName`，并至少启用一种 MERGE 触发方式。

#### `intermediateTableSuffix`

为 upsert 或 delete 中间表添加的名称后缀。

* **类型**：`string`
* **默认值**：`tmp`
* **重要级别**：低
* **有效值 / 注意事项**：必须为非空字符串；中间表名称以目标表名和该后缀开头。

#### `mergeIntervalMs`

启用 upsert 或 delete 时，按时间间隔触发 MERGE，单位为毫秒。

* **类型**：`long`
* **默认值**：`60000`
* **重要级别**：低
* **有效值 / 注意事项**：使用正整数，或使用 `-1` 关闭按时间触发；`mergeIntervalMs` 与 `mergeRecordsThreshold` 不能同时为 `-1`。

#### `mergeRecordsThreshold`

中间表累计到指定记录数时触发 MERGE。

* **类型**：`long`
* **默认值**：`-1`
* **重要级别**：低
* **有效值 / 注意事项**：使用正整数，或使用 `-1` 关闭按记录数触发；`mergeIntervalMs` 与 `mergeRecordsThreshold` 不能同时为 `-1`。

### 表名、字段与分区

#### `sanitizeTopics`

清理 Topic 派生的表名。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：启用后表名可能与原始 Topic 不同。

#### `topic2TableMap`

显式指定 Topic 到 BigQuery 表的映射。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：低
* **有效值 / 注意事项**：格式为 `topic-1:table-1,topic-2:table-2`。启用后忽略 `sanitizeTopics`；未命中的 Topic 仍使用 Topic 名作为表名。不要同时使用会改写 Topic 名称的正则 SMT。

#### `sanitizeFieldNames`

将非法字段名字符替换为下划线，并为数字开头的字段名前置下划线。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：中
* **有效值 / 注意事项**：不同源字段可能清理为同一个名称。

#### `kafkaKeyFieldName`

保存 Kafka 记录 Key 的字段名。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：为空时不写入 Kafka Key。

#### `kafkaDataFieldName`

保存 Kafka 数据和元数据结构的字段名。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：为空时不添加该结构。

#### `bigQueryPartitionDecorator`

使用 BigQuery 分区装饰符写入按日期分区的表。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：高
* **有效值 / 注意事项**：不能与 `timestampPartitionFieldName` 同时使用。

#### `bigQueryMessageTimePartitioning`

使用 Kafka 消息时间生成分区日期，而不是 Connector 处理时间。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：高
* **有效值 / 注意事项**：消息必须带有有效时间戳。

#### `timestampPartitionFieldName`

值中用于时间戳分区的字段名。

* **类型**：`string`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：不能与 `bigQueryPartitionDecorator=true` 同时使用。

#### `clusteringPartitionFieldNames`

用于 BigQuery 聚类的字段列表，使用逗号分隔。

* **类型**：`list`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：最多四个字段，且需要分区表。

#### `timePartitioningType`

自动创建表时使用的时间分区粒度。

* **类型**：`string`
* **默认值**：`DAY`
* **重要级别**：低
* **有效值 / 注意事项**：`HOUR`、`DAY`、`MONTH`、`YEAR` 或 `NONE`；不会修改已存在表的分区类型。

#### `partitionExpirationMs`

自动创建表时设置分区过期时间，单位为毫秒。

* **类型**：`long`
* **默认值**：`null`
* **重要级别**：低
* **有效值 / 注意事项**：仅影响 Connector 新建的表，不会修改已有表；过期分区中的数据会被永久删除。

### 吞吐、重试与批量加载

#### `threadPoolSize`

每个 Task 并发写入 BigQuery 的最大线程数。

* **类型**：`int`
* **默认值**：`10`
* **重要级别**：中
* **有效值 / 注意事项**：至少为 `1`。

#### `queueSize`

写请求队列达到该软上限后暂停分区。

* **类型**：`long`
* **默认值**：`-1`
* **重要级别**：高
* **有效值 / 注意事项**：`-1` 表示不限制；队列下降后恢复分区。

#### `bigQueryRetry`

BigQuery 后端错误或配额超限错误的每请求重试次数。

* **类型**：`int`
* **默认值**：`0`
* **重要级别**：中
* **有效值 / 注意事项**：至少为 `0`。

#### `bigQueryRetryWait`

重试之间的最小等待时间，单位为毫秒。

* **类型**：`long`
* **默认值**：`1000`
* **重要级别**：中
* **有效值 / 注意事项**：至少为 `0`。

#### `max.retries`

Task 遇到可重试写入异常后，重新处理当前记录批次的最大次数。

* **类型**：`int`
* **默认值**：`10`
* **重要级别**：中
* **有效值 / 注意事项**：至少为 `1`。该配置控制 Task 批次级重试，与单次 BigQuery 请求的 `bigQueryRetry` 不同。

#### `enableBatchLoad`

指定通过 GCS 批量加载到 BigQuery 的 Topic 列表。

* **类型**：`list`
* **默认值**：空列表
* **重要级别**：低
* **有效值 / 注意事项**：Beta 功能；启用后必须配置 `gcsBucketName`。

#### `gcsBucketName`

GCS 批量加载对象所在的 Bucket。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：高
* **有效值 / 注意事项**：仅在启用批量加载时必填。

#### `gcsFolderName`

GCS 批量加载对象使用的文件夹前缀。

* **类型**：`string`
* **默认值**：空字符串
* **重要级别**：中
* **有效值 / 注意事项**：仅在启用批量加载时使用。

#### `batchLoadIntervalSec`

执行 GCS 到 BigQuery 加载任务的间隔，单位为秒。

* **类型**：`int`
* **默认值**：`120`
* **重要级别**：低
* **有效值 / 注意事项**：仅在启用批量加载时使用。

#### `autoCreateBucket`

批量加载时自动创建不存在的 GCS Bucket。

* **类型**：`boolean`
* **默认值**：`true`
* **重要级别**：中
* **有效值 / 注意事项**：仅在启用批量加载时使用。

#### `convertDoubleSpecialValues`

将特殊浮点值转换为 BigQuery 可接受的有限值。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：启用前确认下游对转换后数值的解释。

#### `avroDataCacheSize`

Avro Schema 转换缓存大小。

* **类型**：`int`
* **默认值**：`100`
* **重要级别**：低
* **有效值 / 注意事项**：至少为 `0`。

#### `allBQFieldsNullable`

将生成的 BigQuery 字段转换为 `NULLABLE`，数组字段仍为 `REPEATED`。

* **类型**：`boolean`
* **默认值**：`false`
* **重要级别**：低
* **有效值 / 注意事项**：启用时必须同时允许 REQUIRED 字段放宽。

## 最佳实践

### 控制高吞吐写入的并发和背压

**适用业务场景**：写入速率可能超过 BigQuery 处理能力，需要限制积压并对临时后端错误重试。

**配置示例**：

```properties theme={null}
connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector
topics=orders
project=<gcp-project>
defaultDataset=<dataset>
keyfile=<service-account-json-path>
autoCreateTables=false
threadPoolSize=10
queueSize=1000
bigQueryRetry=3
bigQueryRetryWait=1000
```

**关键说明**：`threadPoolSize` 控制并发写入，`queueSize` 达到软上限后会暂停分区。并发度和重试次数应结合 BigQuery 配额与可接受延迟调整。

### 按事件时间写入分区表

**适用业务场景**：记录带有业务事件时间，需要按该字段分区而不是按 Connector 接收时间分区。

**配置示例**：

```properties theme={null}
connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector
topics=orders
project=<gcp-project>
defaultDataset=<dataset>
keyfile=<service-account-json-path>
autoCreateTables=true
schemaRetriever=com.wepay.kafka.connect.bigquery.retrieve.IdentitySchemaRetriever
bigQueryPartitionDecorator=false
timestampPartitionFieldName=event_time
clusteringPartitionFieldNames=customer_id,region
```

**关键说明**：字段分区模式与默认的 `bigQueryPartitionDecorator=true` 互斥。聚类字段最多四个，且必须用于分区表。

## 监控

### 监控内容

监控 Worker、Connector 和 Task 状态，记录吞吐、延迟、Offset 提交、错误、重试和失败任务；同时关注 Worker JVM 堆、GC、线程数和 CPU。启用 `queueSize` 后观察积压以及分区暂停/恢复；只有部署了错误处理和 DLQ 时，才监控 DLQ 活动。

### 导入 Grafana 大盘

下载 [AutoMQ Connect Cluster Dashboard](https://automq-download-center.oss-cn-hangzhou.aliyuncs.com/connect-dashboard/automq-connect-cluster-dashboard.json)，配置包含 Kafka Connect 指标的 Prometheus 数据源和标签，然后在 Grafana 中导入 JSON 大盘。

## 限制条件

* `topics` 与 `topics.regex` 互斥，且至少配置其中一个。
* `timestampPartitionFieldName` 不能与 `bigQueryPartitionDecorator=true` 同时使用。
* `clusteringPartitionFieldNames` 最多四个字段，并且只能用于分区表。
* 字段名清理可能让不同源字段映射为同一个 BigQuery 字段名并产生冲突。
* Connector 不声明 BigQuery 写入的 exactly-once 或绝对不重复保障。

## 常见问题

### Connector 启动时报 Topic 选择冲突，怎么办？

只保留 `topics` 或 `topics.regex` 其中一个，并检查配置中心、环境变量或 Config Provider 是否注入了另一个属性。

### 启动时报缺少 `defaultDataset` 或 `project`，怎么办？

这两个属性没有默认值。填写目标 Google Cloud 项目和 BigQuery 数据集，并确认服务账号有访问权限。

### 目标表不存在时写入失败，怎么办？

保持 `autoCreateTables=true` 并配置 `SchemaRetriever` 实现；如果关闭自动建表，则预先创建目标表并检查其 Schema。

### 为什么记录进入了不同的数据集或表？

检查 Topic 是否被 SMT 改写为 `dataset:table`，以及是否启用了 `sanitizeTopics`。未使用路由时，Connector 使用 `defaultDataset` 和 Topic 名称。

### 写入延迟持续升高时应该检查什么？

检查 BigQuery 配额、错误和重试次数，再检查 `threadPoolSize`、`queueSize` 以及 Task 状态。队列长期不下降时，应降低输入速率或谨慎调整并发和重试。
