跳到主要内容
版本:Next

BigQuery

BigQuery 数据接收器连接器

支持的引擎

Spark
Flink
Seatunnel Zeta

主要特性

描述

用于 Google Cloud BigQuery 的数据接收器连接器,使用 Storage Write API 实现高性能数据摄取。

支持的数据源信息

数据源支持的版本Maven
BigQueryBOM 26.72.0下载

配置选项

名称类型是否必须默认值描述
project_idstring-GCP 项目 ID
dataset_idstring-BigQuery 数据集 ID
table_idstring-BigQuery 表 ID
service_account_key_pathstring-GCP 服务账号 JSON 密钥文件路径
service_account_key_jsonstring-内联 GCP 服务账号 JSON 密钥内容
write_modestringbatch写入模式。支持的值:batchstreaming
sequence_number_columnstring-用于 CDC 去重的序列号列名。仅在 write_modestreaming 时适用
schema_evolution_enabledbooleanfalse是否将 ADD COLUMN Schema 变更事件应用到目标 BigQuery 表
schema_evolution_relax_not_nullbooleanfalseSchema 演进时是否将源端非空列创建为 BigQuery NULLABLE 字段
batch_sizeint1000发送到 BigQuery 之前批量处理的行数
emulator_hoststring-BigQuery emulator REST 地址,例如 localhost:9050。该参数仅用于测试。
emulator_grpc_hoststring-BigQuery emulator Storage Write API 地址,例如 localhost:9060;默认回退到 emulator_host。仅用于测试。
universe_domainstring-Google Cloud 宇宙域/环境域名,例如主权云 S3NS 环境配置为 s3nsapis.fr
schema_save_modeenumCREATE_SCHEMA_WHEN_NOT_EXISTSchema 保存模式。详见下文。
data_save_modeenumAPPEND_DATAData 保存模式。详见下文。
custom_sqlstring-data_save_mode 选择 CUSTOM_PROCESSING 时,需要填写的自定义 SQL 语句。
multi_table_sink_replicaint-Sink 通用参数,用于控制多表运行时每张表的 sink 副本数。
common-options-Sink 通用参数,详见 Sink Common Options

认证参数

生产 BigQuery 任务必须使用下面任意一种认证方式。只有配置 emulator_host 做测试时才会跳过认证。

  1. service_account_key_path:服务账号 JSON 密钥文件路径。
  2. service_account_key_json:直接填写服务账号 JSON 密钥内容。
  3. 默认凭据:如果前两项都不配置,则使用 Google Application Default Credentials。

表选项

目标 BigQuery 表可以通过 SeaTunnel SaveMode 自动创建。 通过将 schema_save_mode 配置为 CREATE_SCHEMA_WHEN_NOT_EXIST(默认值)或 RECREATE_SCHEMA,连接器在初始化时可以基于上游 schema 信息自动创建 BigQuery 数据集和数据表。

连接器写入的目标表由 project_id.dataset_id.table_id 决定。 在多表同步场景下,您可以将 table_id 配置为包含 ${table_name} 的表达式(例如 table_id = "${table_name}"table_id = "prefix_${table_name}"),从而将数据动态路由到不同的 BigQuery 表中。在这种多表设置下,连接器将根据上游表信息自动在 BigQuery 中创建相应的目标表。

schema_save_mode [Enum]

在同步任务启动之前,控制如何处理目标表的结构。

  • RECREATE_SCHEMA :如果目标表存在,则先删除该表然后重新创建;如果不存在则直接创建。
  • CREATE_SCHEMA_WHEN_NOT_EXIST :如果目标表不存在则创建它;如果已存在则跳过创建。
  • ERROR_WHEN_SCHEMA_NOT_EXIST :如果目标表不存在,则抛出异常并报错。
  • IGNORE :忽略目标表结构的处理,不执行任何与结构相关的检查或 DDL 动作。

data_save_mode [Enum]

在同步任务启动之前,控制如何处理目标表中的已有数据。

  • DROP_DATA :删除目标表中的已有数据。
  • APPEND_DATA :保留目标表中的已有数据,并将新数据追加写入。
  • CUSTOM_PROCESSING :执行用户自定义的处理。此选项需要配合 custom_sql 参数。
  • ERROR_WHEN_DATA_EXISTS :如果目标表已包含数据,则抛出异常并报错。

custom_sql [String]

data_save_mode 被设置为 CUSTOM_PROCESSING 时,此参数中填写的自定义 SQL 语句将在数据开始写入前被执行。

Schema 演进

Schema 演进默认关闭。需要在 BigQuery sink 中设置 schema_evolution_enabled = true,并在支持的 CDC source 中设置 schema-changes.enabled = true,才能将源表的 ADD COLUMN 事件同步到配置的目标表。

仅支持物理列的 ADD COLUMN 事件。默认情况下,新增的标量列或 struct 列必须允许为空。设置 schema_evolution_relax_not_null = true 后,源端的非空标量列或 struct 列会在 BigQuery 中创建为 NULLABLE 字段;这是因为目标表中的历史数据没有新列对应的值。

源端 array 列必须是非空列,并会创建为 BigQuery REPEATED 字段。nullable array 会被拒绝,因为 BigQuery array 不能为 NULL;静默映射会丢失 NULL 与空数组之间的区别。不支持 DROP COLUMNRENAME COLUMNMODIFY COLUMN。BigQuery 会把新字段追加到目标 Schema 末尾,因此源事件中的 FIRSTAFTER 位置提示不会改变 BigQuery 的物理字段顺序。数据行按字段名编码,sink 会在接收使用新字段的数据前刷新 writer Schema。

遇到不支持的 Schema 变更时,任务会失败而不会静默跳过,因为在源端和目标端 Schema 不一致的情况下继续运行可能导致后续数据错位或损坏。从同一个 checkpoint 恢复会再次回放该事件并重复失败。优先在 BigQuery 上手动执行与兼容性校验相符的等价 DDL,然后从同一个 checkpoint 恢复,这样回放的事件就会成为一次幂等的空操作。仅在别无选择时才从更晚的 source 位置重新启动,因为跳过会丢失从上一个 checkpoint 到新起始位置之间的所有源端记录;如果数据链路本身可能产生不支持的 DDL,请考虑关闭 schema-changes.enabled,并在 SeaTunnel 外部管理这些 Schema 变更。

Schema 更新使用 ALTER TABLE ... ADD COLUMN IF NOT EXISTS。如果目标表已经存在同名字段,其类型和模式必须兼容,否则任务会失败。除 Storage Write API 写入所需权限外,凭据还必须能够执行 DDL job 并读取更新后的表元数据。

当 sink 并行度大于 1 时,多个 subtask 可能收到同一个 Schema 变更事件并并发提交相同的 ALTER TABLE 语句。sink 已经对此做了容错处理:竞争失败的 subtask 会重新读取表结构并将该列视为已存在;如果 BigQuery 因单表元数据更新配额返回 rateLimitExceeded,sink 会按有界指数退避策略自动重试。并行度非常高时,建议将启用 schema_evolution_enabled 的任务保持在适中的并行度,以便 Schema 变更更快收敛。

写入模式

  • batch:使用 BigQuery buffered write stream,并在 SeaTunnel checkpoint/commit 阶段提交数据。主要特性中的精确一次能力指的是该模式。
  • streaming:使用默认 stream,并携带 BigQuery change 字段写入 CDC 记录。该模式适合 CDC 的 upsert/delete 数据,但该连接器没有将它标记为精确一次。

使用 streaming 模式写入 CDC 数据时,请先在 BigQuery 中创建好带 Primary Key 的目标表。连接器会把 SeaTunnel 的行类型转换为 BigQuery change 记录:INSERTUPDATE_AFTER 会写成 UPSERTDELETEUPDATE_BEFORE 会写成 DELETE

sequence_number_column

sequence_number_column 是可选配置。

当配置了 sequence_number_column 时,该列的值会作为 _CHANGE_SEQUENCE_NUMBER 发送到 BigQuery,用于启用 BigQuery 侧的去重。在 source 重新发送数据时,具有相同 primary key 和相同 sequence number 的行可以由 BigQuery 进行去重。 如果没有配置 sequence_number_column,则不会发送 _CHANGE_SEQUENCE_NUMBER,BigQuery 也不会执行基于 sequence number 的去重。

注意

  • BigQuery 要求 _CHANGE_SEQUENCE_NUMBER 是十六进制 STRING。对于整数列以及精确的整数 decimal 值(例如映射为 DECIMAL(20, 0) 的 MySQL BIGINT UNSIGNED),connector 会将 unsigned 64-bit 范围内的非负值转换为十六进制字符串;对于字符串列,connector 会将值视为已编码的十六进制 sequence number,仅进行校验而不转换。
  • sequence number 最多可以包含 4 个以 / 分隔的 section,每个 section 最多包含 16 个十六进制字符。Null、负数、空值或格式错误的值会被拒绝。
  • sequence_number_column 应该引用 source 表中单调递增的列,例如以 epoch millis 表示的 updated_atversionseq_id
  • 如果要在 streaming 模式下启用 BigQuery 侧的去重,目标 BigQuery 表必须定义 Primary Key。否则,无论是否配置 sequence number,BigQuery 都会将每次写入视为 append 操作。

emulator_host

emulator_host 只用于本地测试或 CI 测试,用于配置 emulator 的 REST 地址。配置后,SeaTunnel 会无凭据连接 BigQuery emulator。当 emulator 的 Storage Write API 使用不同地址时,需要设置 emulator_grpc_host;例如 goccy BigQuery emulator 默认使用 9060 端口。未配置时,gRPC 地址会回退到 emulator_host。生产任务不要使用这些参数。

任务示例

简单批处理示例

env {
parallelism = 1
job.mode = "BATCH"
}

source {
FakeSource {
row.num = 10
string.fake.mode = "template"
string.template = ["key", "value"]
schema = {
fields {
c_map = "map<string, string>"
c_array = "array<int>"
c_string = string
c_boolean = boolean
c_tinyint = tinyint
c_smallint = smallint
c_int = int
c_bigint = bigint
c_float = float
c_double = double
c_decimal = "decimal(30, 8)"
c_bytes = bytes
c_date = date
c_timestamp = timestamp
c_time = time
}
}
}
}

sink {
BigQuery {
project_id = "test-project"
dataset_id = "test_dataset"
table_id = "test_table"
batch_size = 2
emulator_host = "localhost:9050"
emulator_grpc_host = "localhost:9060"
}
}

CDC 流式模式(MySQL 到 BigQuery)

目标 BigQuery 表需要提前创建,并且应定义 CDC 源表使用的主键。例如:

CREATE TABLE `my-gcp-project.cdc_dataset.orders` (
uuid INT64 NOT NULL,
name STRING,
score INT64,
PRIMARY KEY (uuid) NOT ENFORCED
)
OPTIONS (max_staleness = INTERVAL 0 MINUTE);
env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 10000
}

source {
MySQL-CDC {
parallelism = 1
server-id = 5652
username = "st_user_source"
password = "mysqlpw"
table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
schema-changes.enabled = true
}
}

sink {
BigQuery {
project_id = "my-gcp-project"
dataset_id = "cdc_dataset"
table_id = "orders"
service_account_key_path = "/path/to/key.json"
write_mode = "streaming"
schema_evolution_enabled = true
batch_size = 500
}
}

如果上游 CDC 源能产生单调递增的列(例如 updated_at 毫秒时间戳或行版本号),可以把它配到 sequence_number_column,让 BigQuery 端对重试批次做去重。目标表必须定义主键(上例使用 PRIMARY KEY (uuid) NOT ENFORCED),否则 BigQuery 会把每次写入都当作 append,跳过去重。

sink {
BigQuery {
project_id = "my-gcp-project"
dataset_id = "cdc_dataset"
table_id = "orders"
service_account_key_path = "/path/to/key.json"
write_mode = "streaming"
sequence_number_column = "updated_at"
batch_size = 500
}
}

复杂数据类型示例

source {
FakeSource {
row.num = 100
schema = {
fields {
order_id = "bigint"
customer = {
name = "string"
email = "string"
}
items = "array<string>"
metadata = "map<string, string>"
order_date = "date"
}
}
}
}

sink {
BigQuery {
project_id = "my-gcp-project"
dataset_id = "orders"
table_id = "customer_orders"
service_account_key_path = "/path/to/key.json"
batch_size = 500
}
}

内联服务账号密钥

如果不便挂载密钥文件(例如 CI runner、把密钥放在 Kubernetes Secret 中以环境变量形式注入),可以直接把 JSON 内容放到 service_account_key_json 中。

sink {
BigQuery {
project_id = "my-gcp-project"
dataset_id = "orders"
table_id = "customer_orders"
service_account_key_json = "${GCP_SA_KEY_JSON}"
batch_size = 500
}
}

测试

该连接器同时使用 BigQuery REST API 和 Storage Write API。使用 goccy BigQuery emulator 时,请将 emulator_host 配置为 REST 端口(默认 9050),并将 emulator_grpc_host 配置为 gRPC 端口(默认 9060)。 Emulator 适合用于本地和 CI 覆盖,但生产可用性仍应在真实 BigQuery 环境中验证。

更新日志

Change Log
ChangeCommitVersion