BigQuery 接收器连接器会将数据从 Kafka 流式传输到 BigQuery 表,从而在 BigQuery 中实现实时数据注入和分析。
BigQuery Sink 连接器的使用场景包括:
数据仓储。将流式数据加载到 BigQuery 中,以进行分析和报告。
填充支持实时信息中心的 BigQuery 表。
准备工作
在创建 BigQuery 接收器连接器之前,请确保您具备以下条件:
为您的 Connect 集群创建 Managed Service for Apache Kafka 集群。此集群是与 Connect 集群关联的 Kafka 主集群。此集群也是 BigQuery 接收器连接器流水线的一端所连接的源集群。
创建 Connect 集群以托管 BigQuery 接收器连接器。
创建一个 BigQuery 数据集来存储从 Kafka 流式传输的数据。
在源集群中创建并配置 Kafka 主题。数据从相应 Kafka 主题转移到目标 BigQuery 数据集。
所需的角色和权限
如需获得创建连接器所需的权限,请让您的管理员为您授予项目的 Managed Kafka Connector Editor (roles/managedkafka.connectorEditor) IAM 角色。如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限。
此预定义角色包含创建连接器所需的权限。如需查看所需的确切权限,请展开所需权限部分:
所需权限
创建连接器需要以下权限:
-
创建连接器:
managedkafka.connectors.create
授予向 BigQuery 表写入数据的权限
Managed Kafka 服务账号必须有权将消息写入 BigQuery 表。向包含相应表的项目中的服务账号授予 BigQuery Data Editor (roles/bigquery.dataEditor) 角色。
受管 Kafka 服务账号采用以下格式:service-PROJECT_NUMBER@gcp-sa-managedkafka.,其中 PROJECT_NUMBER 是 Connect 集群的项目编号。
如果您的 Connect 集群与 Managed Service for Apache Kafka 集群位于不同的项目中,请参阅 在其他项目中创建 Connect 集群。
关键配置参数
以下是有关 BigQuery 接收器连接器配置的建议。
队列大小
queueSize 设置用于控制向 BigQuery 发送的写入请求数上限。默认情况下,此值未配置,这可能会导致 Connect 集群出现内存不足错误。建议您配置此值,以减少此类错误的发生几率。队列中的项是写入请求,不直接对应于消息数量或字节数。因此,很难提前知道建议值。
从 Kafka 轮询消息时,连接器会根据消息的目标表拆分工作。如果轮询中的所有消息都要写入单个表,则只会添加一个工作项。但如果轮询中的消息要写入三个表,则会添加三个工作项。工作项的大小可以根据消息轮询配置、Kafka 消息的大小和写入的表数量大致估算。
BigQuery 接收器连接器的架构
BigQuery 接收器连接器使用配置的值转换器 (value.converter) 将 Kafka 记录值解析为字段。然后,它会将这些字段写入 BigQuery 表中具有相同名称的列。
连接器需要架构才能运行。您可以通过以下方式提供架构:
- 基于消息的架构:架构包含在每条消息中。
- 基于表的架构:连接器会从 BigQuery 表架构推断消息架构。
- 架构注册表:连接器从架构注册表(例如 Managed Service for Apache Kafka 架构注册表 [预览版])读取架构。
接下来的部分将介绍这些选项。
基于消息的架构
在此模式下,每个 Kafka 记录都包含一个 JSON 架构。连接器使用架构将记录数据写入为 BigQuery 表行。
如需使用基于消息的架构,请在连接器上设置以下属性:
value.converter=org.apache.kafka.connect.json.JsonConvertervalue.converter.schemas.enable=true
Kafka 记录值示例:
{
"schema": {
"type": "struct",
"fields": [
{
"field": "user",
"type": "string",
"optional": false
},
{
"field": "age",
"type": "int64",
"optional": false
}
]
},
"payload": {
"user": "userId",
"age": 30
}
}
如果目标表已存在,则 BigQuery 表架构必须与嵌入式消息架构兼容。如果值为 autoCreateTables=true,连接器会在需要时自动创建目标表。如需了解详情,请参阅创建表。
如果您希望连接器在消息架构发生变化时更新 BigQuery 表架构,请将 allowNewBigQueryFields、allowSchemaUnionization 或 allowBigQueryRequiredFieldRelaxation 设置为 true。
基于表格的架构
在此模式下,Kafka 记录包含没有明确架构的纯 JSON 数据。连接器会根据目标表推断架构。
要求:
- BigQuery 表必须已存在。
- Kafka 记录数据必须与表架构兼容。
- 此模式不支持根据传入的消息动态更新架构。
如需使用基于表的架构,请在连接器上设置以下属性:
value.converter=org.apache.kafka.connect.json.JsonConvertervalue.converter.schemas.enable=falsebigQueryPartitionDecorator=false
如果 BigQuery 表使用基于时间的分区(每日分区),则 bigQueryPartitionDecorator 可以是 true。否则,请将此属性设置为 false。
Kafka 记录值示例:
{
"user": "userId",
"age": 30
}
架构注册表
在此模式下,每个 Kafka 记录都包含 Apache Avro 数据,并且消息架构存储在架构注册表中。
如需将 BigQuery 接收器连接器与架构注册表搭配使用,请在连接器上设置以下属性:
value.converter=io.confluent.connect.avro.AvroConvertervalue.converter.schema.registry.url=SCHEMA_REGISTRY_URL
将 SCHEMA_REGISTRY_URL 替换为架构注册表的网址。
如需将连接器与 Managed Service for Apache Kafka 架构注册表搭配使用,请设置以下属性:
value.converter.bearer.auth.credentials.source=GCP
如需了解详情,请参阅将 Kafka Connect 与架构注册表搭配使用。
Apache Iceberg 受管理表
BigQuery Sink 连接器支持将 Apache Iceberg 受管表(以下简称“Iceberg 受管表”)作为接收器目标。
Iceberg 受管表为在 Google Cloud上构建开放格式湖仓一体提供了基础。Iceberg 受管理表提供与 BigQuery 表相同的全代管式体验,但使用 Parquet 将数据存储在客户拥有的存储分区中,以便与 Apache Iceberg 开放表格式进行互操作。
如需了解如何创建 Apache Iceberg 表,请参阅创建 Apache Iceberg 表。
创建 BigQuery 接收器连接器
控制台
在 Google Cloud 控制台中,前往连接集群页面。
点击要创建连接器的 Connect 集群。
点击创建连接器。
对于连接器名称,请输入一个字符串。
有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。
对于连接器插件,选择 BigQuery Sink。
在主题部分中,指定要从中读取数据的 Kafka 主题。您可以指定主题列表或正则表达式,以匹配主题名称。
方法 1:选择选择一个 Kafka 主题列表。在 Kafka 主题列表中,选择一个或多个主题。点击确定。
方法 2:选择使用主题正则表达式。在主题正则表达式字段中,输入正则表达式。
点击数据集,然后指定一个 BigQuery 数据集。您可以选择现有数据集,也可以创建新数据集。
可选:在配置框中,添加配置属性或修改默认属性。如需了解详情,请参阅配置连接器。
选择任务重启政策。如需了解详情,请参阅任务重启政策。
点击创建。
gcloud
-
在 Google Cloud 控制台中,激活 Cloud Shell。
Cloud Shell 会话随即会在 Google Cloud 控制台的底部启动,并显示命令行提示符。Cloud Shell 是一个已安装 Google Cloud CLI 且已为当前项目设置值的 Shell 环境。该会话可能需要几秒钟时间来完成初始化。
运行
gcloud managed-kafka connectors create命令:gcloud managed-kafka connectors create CONNECTOR_ID \ --location=LOCATION \ --connect-cluster=CONNECT_CLUSTER_ID \ --config-file=CONFIG_FILE替换以下内容:
CONNECTOR_ID:连接器的 ID 或名称。 有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。连接器的名称是不可变的。
LOCATION:您创建连接器的位置。此位置必须与您创建 Connect 集群的位置相同。
CONNECT_CLUSTER_ID:创建连接器的 Connect 集群的 ID。
CONFIG_FILE:BigQuery Sink 连接器的 YAML 配置文件路径。
以下是 BigQuery Sink 连接器的配置文件示例:
name: "BQ_SINK_CONNECTOR_ID" project: "GCP_PROJECT_ID" topics: "GMK_TOPIC_ID" tasks.max: 3 connector.class: "com.wepay.kafka.connect.bigquery.BigQuerySinkConnector" key.converter: "org.apache.kafka.connect.storage.StringConverter" value.converter: "org.apache.kafka.connect.json.JsonConverter" value.converter.schemas.enable: "false" defaultDataset: "BQ_DATASET_ID"替换以下内容:
BQ_SINK_CONNECTOR_ID:BigQuery Sink 连接器的 ID 或名称。有关如何命名连接器的指南,请参阅 Managed Service for Apache Kafka 资源命名指南。连接器的名称是不可变的。
GCP_PROJECT_ID:BigQuery 数据集所在的 Google Cloud项目的 ID。
GMK_TOPIC_ID:Managed Service for Apache Kafka 主题的 ID,数据从该主题流向 BigQuery Sink 连接器。
BQ_DATASET_ID:作为流水线接收器的 BigQuery 数据集的 ID。
Terraform
您可以使用 Terraform 资源创建连接器。
如需了解如何应用或移除 Terraform 配置,请参阅基本 Terraform 命令。
Go
在尝试此示例之前,请按照 安装客户端库中的 Go 设置说明进行操作。如需了解详情,请参阅 Managed Service for Apache Kafka Go API 参考文档。
如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据(ADC)。如需了解详情,请参阅为本地开发环境设置 ADC。
Java
在尝试此示例之前,请按照 安装客户端库中的 Java 设置说明进行操作。如需了解详情,请参阅 Managed Service for Apache Kafka Java API 参考文档。
如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据。如需了解详情,请参阅 为本地开发环境设置 ADC。
Python
在尝试此示例之前,请按照 安装客户端库中的 Python 设置说明进行操作。如需了解详情,请参阅 Managed Service for Apache Kafka Python API 参考文档。
如需向 Managed Service for Apache Kafka 进行身份验证,请设置应用默认凭据。如需了解详情,请参阅为本地开发环境设置 ADC。
创建连接器后,您可以修改、删除、暂停、停止或重启连接器。
配置连接器
本部分介绍了一些可在连接器上设置的配置属性。如需查看此连接器特有的属性的完整列表,请参阅 BigQuery Sink 连接器配置。
表名称
默认情况下,连接器使用主题名称作为 BigQuery 表名称。如需使用其他表名称,请按以下格式设置 topic2TableMap 属性:
topic2TableMap=TOPIC_1:TABLE_1,TOPIC_2:TABLE_2,...
创建表格
如果目标表不存在,BigQuery 接收器连接器可以创建这些表。
如果值为
autoCreateTables=true,连接器会尝试创建任何不存在的 BigQuery 表。此设置是默认行为。如果为
autoCreateTables=false,连接器不会创建任何表。如果目标表不存在,则会发生错误。
当 autoCreateTables 为 true 时,您可以使用以下配置属性来更精细地控制连接器创建和配置新表的方式:
allBQFieldsNullableclusteringPartitionFieldNamesconvertDoubleSpecialValuespartitionExpirationMssanitizeFieldNamessanitizeTopicstimestampPartitionFieldName
如需了解这些属性,请参阅 BigQuery Sink 连接器配置。
Kafka 元数据
您可以通过分别配置 kafkaDataFieldName 和 kafkaKeyFieldName 字段,将 Kafka 中的其他数据(例如元数据信息和键信息)映射到 BigQuery 表中。元数据信息的示例包括 Kafka 主题、分区、偏移量和插入时间。