Skip to main content

概述

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
  • 重要级别:中
  • 有效值 / 注意事项FILEJSONAPPLICATION_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 关闭按时间触发;mergeIntervalMsmergeRecordsThreshold 不能同时为 -1

mergeRecordsThreshold

中间表累计到指定记录数时触发 MERGE。
  • 类型long
  • 默认值-1
  • 重要级别:低
  • 有效值 / 注意事项:使用正整数,或使用 -1 关闭按记录数触发;mergeIntervalMsmergeRecordsThreshold 不能同时为 -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
  • 重要级别:低
  • 有效值 / 注意事项HOURDAYMONTHYEARNONE;不会修改已存在表的分区类型。

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 大盘。

限制条件

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

常见问题

Connector 启动时报 Topic 选择冲突,怎么办?

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

启动时报缺少 defaultDatasetproject,怎么办?

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

目标表不存在时写入失败,怎么办?

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

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

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

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

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