概述
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。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 处理能力,需要限制积压并对临时后端错误重试。 配置示例:threadPoolSize 控制并发写入,queueSize 达到软上限后会暂停分区。并发度和重试次数应结合 BigQuery 配额与可接受延迟调整。
按事件时间写入分区表
适用业务场景:记录带有业务事件时间,需要按该字段分区而不是按 Connector 接收时间分区。 配置示例:bigQueryPartitionDecorator=true 互斥。聚类字段最多四个,且必须用于分区表。
监控
监控内容
监控 Worker、Connector 和 Task 状态,记录吞吐、延迟、Offset 提交、错误、重试和失败任务;同时关注 Worker JVM 堆、GC、线程数和 CPU。启用queueSize 后观察积压以及分区暂停/恢复;只有部署了错误处理和 DLQ 时,才监控 DLQ 活动。
导入 Grafana 大盘
下载 AutoMQ Connect Cluster Dashboard,配置包含 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 状态。队列长期不下降时,应降低输入速率或谨慎调整并发和重试。