CockroachDB
Debezium CockroachDB 连接器
Debezium CockroachDB 连接器从 CockroachDB 增强型 changefeed 中捕获行级变更,并将其流式传输到 Kafka 主题中。
CockroachDB 的 changefeed 能够近乎实时地检测数据库中的数据变更——插入、更新和删除。该连接器会创建一个 changefeed,将事件发布到中间 Kafka 集群,然后消费这些事件,将其转换为标准的 Debezium 信封格式,并发送到下游消费者所订阅的输出 Kafka 主题中。
| CockroachDB 连接器是一个孵化中的连接器。孵化中的连接器是指为预览目的而发布的连接器,其后续变更可能并非始终向后兼容。 |
|---|
概述
CockroachDB 是一个分布式 SQL 数据库,提供强一致性、水平扩展能力和高存活能力。与依赖预写日志(WAL)进行变更数据捕获的传统数据库不同,CockroachDB 提供了原生的 changefeed 机制,可直接从存储层流式传输行级变更。
Debezium CockroachDB 连接器利用了这一原生 changefeed 能力。连接器为其捕获的每一个行级插入、更新和删除操作生成一个变更事件,并将每个表的变更事件记录发送到各自独立的 Kafka 主题中。客户端应用程序读取与所关注的数据库表相对应的 Kafka 主题,并可对从这些主题接收到的每个行级事件作出响应。
连接器具备容错能力。在读取变更和生成事件的过程中,连接器会记录其处理的最后一个已解析时间戳。如果连接器因任何原因(包括通信故障、网络问题或崩溃)而停止,连接器重启后会从中断处的最后一条记录继续流式传输。
连接器的工作原理
要对 Debezium CockroachDB 连接器进行最佳配置和运行,了解该连接器如何执行快照、流式传输变更事件、确定 Kafka 主题名称以及使用元数据是很有帮助的。
架构
CockroachDB 连接器采用两阶段 Kafka 架构:
- CockroachDB 变更流到中间 Kafka:连接器创建单个 CockroachDB 变更流,对所有已配置的表使用
enriched信封格式(CREATE CHANGEFEED FOR table1, table2, …)。CockroachDB 会自动将事件路由到中间集群中按表划分的 Kafka 主题。enriched 格式同时包含模式元数据以及变更前后的完整行状态。 - 中间 Kafka 到 Debezium 再到输出 Kafka:连接器在单个 KafkaConsumer 中订阅所有按表划分的 Kafka 主题;根据主题名称将每个事件路由到对应的表;将 enriched 变更流事件转换为标准的 Debezium 信封格式(包含
before、after、source和op字段);并将其生产到最终的输出 Kafka 主题。
通过采用这种架构,连接器能够利用 CockroachDB 原生的高可用变更流基础设施来可靠地捕获数据,并以标准 Debezium 事件格式提供给下游消费者。默认情况下,连接器会将所有已配置的表合并到单个变更流中。通过将数据整合到单个变更流作业中,连接器符合建议的每个 CockroachDB 集群变更流数量上限(撰写本文时约为 80 个)。为了提高吞吐量并避免性能耦合,CockroachDB 的最佳实践建议不要使用单个变更流来监控大量表。在表数量很多的环境中,可以考虑设置 cockroachdb.changefeed.max.tables.per.changefeed 将表拆分到多个变更流中,或者运行多个连接器实例,使每个连接器捕获相关的一组表。
快照
CockroachDB 变更流原生支持 initial_scan 选项,该选项会在流式传输持续变更之前回填所有现有行。Debezium CockroachDB 连接器将其 snapshot.mode 配置映射到 CockroachDB 的 initial_scan 变更流选项,因此连接器不需要单独的基于 JDBC 的快照阶段。
在初始扫描阶段,捕获的事件被标记为 op=r,表示读操作,以此将快照事件与持续变更事件区分开来。初始扫描完成后,变更流转入流式传输阶段,事件会被标记为相应的操作类型(c 表示创建、u 表示更新、d 表示删除)。
下表展示了每种 snapshot.mode 与 CockroachDB initial_scan 选项的对应关系:
| 快照模式 | initial_scan 值 | 说明 |
|---|---|---|
initial(默认) | 首次启动为 yes,重启时为 no | 首次启动(无历史偏移量)时,快照会回填所有现有行。在存在偏移量的情况下重启时,流式传输将从存储的游标位置继续。 |
always | yes | 快照始终回填所有现有行,即使在重启时也是如此。 |
initial_only | only | 快照回填所有现有行,然后连接器停止。适用于一次性数据迁移。 |
no_data / never | no | 完全跳过初始扫描。仅捕获变更数据流(changefeed)创建之后产生的持续变更。 |
when_needed | 首次启动为 yes,重启时为 no | 与 initial 相同,但如果存储的偏移量不再有效,连接器将执行一次新的初始快照。 |
表 1. 快照模式与 initial_scan 的映射关系
流式传输变更
初始扫描完成后,CockroachDB 连接器将持续流式传输变更。当 CockroachDB 中发生行级变更时,变更数据流会将相应的事件写入中间 Kafka 主题。连接器消费这些事件,将其转换为 Debezium 变更事件,并转发到输出 Kafka 主题。
CockroachDB 变更数据流会以可配置的时间间隔发出已解析时间戳消息。这些消息表示该时间戳之前的所有变更均已发出。连接器使用已解析时间戳进行偏移量跟踪,以便在重启时能够创建一个 cursor 指向最后一个已解析时间戳的新变更数据流,确保不会遗漏任何事件。
当连接器接收到变更时,它会将事件转换为 Debezium 的 read、create、update 或 delete 事件记录。连接器将这些变更记录转发给运行在同一进程中的 Kafka Connect 框架。Kafka Connect 进程会按照事件生成的相同顺序,异步地将变更事件记录写入相应的 Kafka 主题。
Kafka Connect 会定期将最新的 offset 写入另一个 Kafka 主题。offset 表示 Debezium 随每个事件一同携带的、与源相关的位点信息。对于 CockroachDB 连接器,最后解析的时间戳就是 offset。
当 Kafka Connect 正常关闭时,它会停止所有连接器,将所有事件记录刷写到 Kafka,并记录从每个连接器收到的最后一个 offset。当 Kafka Connect 重启时,它会读取每个连接器最后记录的 offset,并从各自最后记录的 offset 处启动每个连接器。
复用现有的 changefeed
默认情况下,连接器会自行创建并管理其 CockroachDB changefeed。在创建 changefeed 之前,连接器会检查是否已有正在运行的 changefeed 覆盖了所配置的表。如果找到匹配的 changefeed,连接器会跳过创建,改为消费现有的 changefeed。这种行为使连接器的重启具有幂等性,避免在集群上创建重复的 changefeed 作业,同时也允许运维人员预先配置一个 changefeed,随后由连接器接入。
只有当正在运行的 changefeed 作业同时满足以下两个条件时,该 changefeed 才被视为匹配:
- 其描述中包含每个已配置表的完全限定名称。
- 其描述中包含连接器所期望的
topic_prefix,即cockroachdb.changefeed.sink.topic.prefix的值;若未设置 sink 主题前缀,则为topic.prefix的值。
为了使连接器能够正确消费现有的 changefeed,该 changefeed 必须发布到名称与连接器订阅的主题相匹配的 Kafka 主题,格式为 topicPrefix.database.schema.table。为生成这种主题布局,并以连接器所期望的结构发出事件,请使用以下选项创建 changefeed:
full_table_name
必填,使 CockroachDB 以 database.schema.table 形式命名主题。
topic_prefix='prefix.'
必填,且必须等于上文列表中描述的连接器主题前缀。指定的前缀必须以点号(.)结尾。
envelope='enriched'
必填。连接器消费的是 enriched 信封。
enriched_properties
可选。连接器读取的是基础 enriched 字段(op、ts_ns、after,以及设置了 diff 时的 before),因此消费该 feed 不需要任何 enriched 属性。默认值 source 可确保中间主题包含来源和提交时间元数据。
format='json'
必填。连接器消费的是经 JSON 编码的 changefeed 消息。
diff、updated、resolved
将这些选项设置为与连接器配置一致,以便前像、更新标记和解析时间戳可用。
如果其中任何一个条件不满足,连接器就不会识别现有的 changefeed,而是创建自己的 changefeed。使用不同的 topic_prefix 创建的 changefeed,或者未使用 full_table_name 的 changefeed,都不会被识别为匹配项,因此连接器会创建自己的 changefeed。如果正在运行的 changefeed 与连接器的主题前缀和表相匹配,但它使用的封套(envelope)不是 enriched,连接器就无法消费它。此时连接器将启动失败,并返回一个错误,提示你必须使用 envelope='enriched' 重新创建该 changefeed,或者使用其他主题前缀。
| 对于大多数部署,最简单的方式是让连接器创建并管理 changefeed。只有在有明确理由时才复用现有的 changefeed,例如保留现有的游标位置或自定义分区方案。 |
|---|
心跳机制
CockroachDB changefeed 中的已解析时间戳可作为天然的心跳信号。当连接器收到已解析时间戳时,会执行以下操作:
- 将存储的偏移量游标更新为该已解析时间戳的值。
- 分发一个 Debezium 心跳事件,以推进 Kafka Connect 的偏移量。
通过发送心跳消息,连接器有助于确保即使在没有数据变更的空闲时段,偏移量也能持续推进。如果偏移量无法推进,监控系统可能会报告连接器无响应,这会导致 CockroachDB 的受保护时间戳不断累积,从而阻碍垃圾回收。
要在 __debezium-heartbeat.<topic.prefix> Kafka 主题上启用 Debezium 心跳记录,请将 heartbeat.interval.ms 设置为大于 0 的值(单位为毫秒)。cockroachdb.changefeed.resolved.interval 属性用于控制 CockroachDB 发送已解析时间戳的频率(默认值:10s)。
模式演进
连接器会自动检测 ALTER TABLE ADD COLUMN、DROP COLUMN 和 RENAME COLUMN 等 DDL 变更,无需重启。当传入的 changefeed 事件包含与已注册表模式不匹配的字段时,连接器会通过执行以下任务来调整模式:
- 通过将事件字段名与已注册的列名进行比较来检测不匹配。
- 重新查询
information_schema,以获取更新后的表定义。 - 刷新内部模式,并根据调整后的列布局继续处理该表。
CockroachDB 变更流通过执行 backfill(回填)操作原生处理模式变更,期间 Schema Change Manager 会使用更新后的模式重新发布所有现有行。回填可确保在模式转换期间不会丢失任何事件。回填完成后,刷新后的列会立即对下游消费者可用。
连接器会调整其发布的变更事件的模式,使输出主题上的记录始终反映当前的表结构。它不维护数据库历史,也不会将 DDL 变更事件发布到单独的模式变更主题中。连接器所使用的模式来自它自身通过 JDBC 对表的发现,而不是来自变更流消息中的 schema 块,因此即使 cockroachdb.changefeed.enriched.properties 仅设置为 source,模式变更也能得到处理。 |
|---|
增量快照
CockroachDB 连接器支持 Debezium 的基于信号的增量快照。增量快照会从一个或多个表中重新读取现有行,并将其作为 op=r(读取)事件发布,与此同时连接器持续不间断地捕获正在进行的变更。
你可以使用增量快照执行以下任务:
- 重新填充丢失了数据的下游消费者。
- 为新添加的接收端或主题进行数据回填。
- 通过重新快照并比较来验证源与目标的一致性。
准备使用增量快照
要使连接器能够使用增量快照,你必须创建一个信号表,然后更新连接器配置以识别该信号表。
操作步骤
通过执行以下 SQL 命令在 CockroachDB 中创建信号表:
CREATE TABLE debezium_signal ( id STRING PRIMARY KEY, type STRING NOT NULL, data STRING );在连接器配置中添加以下属性,以引用该信号表:
{ "signal.data.collection": "mydb.public.debezium_signal", "table.include.list": "public.my_table,public.debezium_signal" }
你必须将信号表包含在 table.include.list 中,这样变更馈送才能向连接器传递信号事件。
触发增量快照
通过在连接器的信号表中插入一行来触发增量快照。
操作步骤
使用以下格式运行 SQL 命令,向信号表中添加一行以触发增量快照:
INSERT INTO debezium_signal (id, type, data) VALUES ('snap-1', 'execute-snapshot', '{"data-collections": ["mydb.public.my_table"]}');
连接器会从指定的表中重新读取所有行,并将它们作为 op=r 事件发出。
有关信号传递的更多信息,请参阅 向 Debezium 连接器发送信号。
主题名称
CockroachDB 连接器会将某张表的所有数据变更事件(插入、更新和删除)发送到专门用于该表的 Kafka 主题中。默认情况下,Kafka 主题名称为 topicPrefix.schemaName.tableName,其中:
- topicPrefix 是由
topic.prefix连接器配置属性指定的主题前缀。 - schemaName 是数据库模式的名称(默认值:
public)。 - tableName 是发生该操作的数据库表的名称。
例如,假设某个连接器捕获包含两张表 orders 和 customers 的数据库的变更,其主题前缀为 cockroachdb。该连接器将把记录流式传输到以下两个 Kafka 主题:
cockroachdb.public.orderscockroachdb.public.customers
该连接器采用了与其他 Debezium 连接器相似的命名约定,并支持标准的主题命名策略。
数据变更事件
Debezium CockroachDB 连接器会为每个行级的 INSERT、UPDATE 和 DELETE 操作生成一个数据变更事件。每个事件都包含一个键和一个值。键和值是相互独立的文档。键和值的结构取决于发生变更的表。
Debezium 和 Kafka Connect 是围绕事件消息的持续流设计的。然而,这些事件的结构可能会随时间发生变化,这对消费者来说可能难以处理。为了帮助消费者适应结构上的多变性,连接器发出的是自包含的事件。也就是说,每个事件消息要么包含其内容的模式(schema),要么在使用模式注册中心的环境中,每条消息都包含一个模式 ID,消费者可用它从注册中心获取模式。
以下骨架 JSON 文档展示了事件消息中键字段和值字段的基本结构。不过,键和值文档的具体表示形式取决于你在应用中所配置的 Kafka Connect 转换器。schema 字段只有在你配置转换器生成它时,才会出现在变更事件键或变更事件值中。同样地,事件键和事件载荷也只有在你配置转换器生成它们时才会出现。如果你使用 JSON 转换器并配置它生成模式,变更事件将具有如下结构:
// Key
{
"schema": { (1)
...
},
"payload": { (2)
...
}
}
// Value
{
"schema": { (3)
...
},
"payload": { (4)
...
}
}| 项目 | 字段名 | 描述 |
|---|---|---|
| 1 | schema | 第一个 schema 字段是事件键的一部分。它指定了一个 Kafka Connect 模式,用于描述事件键的 payload 部分包含的内容。换句话说,第一个 schema 字段描述的是主键的结构。 |
| 2 | payload | 第一个 payload 字段是事件键的一部分。它具有前一个 schema 字段所描述的结构,并且包含被更改行的键。 |
| 3 | schema | 第二个 schema 字段是事件值的一部分。它指定了一个 Kafka Connect 模式,用于描述事件值的 payload 部分包含的内容。换句话说,第二个 schema 描述的是被更改行的结构。通常,该模式中包含嵌套模式。 |
| 4 | payload | 第二个 payload 字段是事件值的一部分。它具有前一个 schema 字段所描述的结构,并且包含被更改行的实际数据。 |
表 2. 变更事件基本内容概览
默认的 xref:主题命名行为 会使连接器将变更事件记录流式传输到这样一个主题,其名称与该事件所描述的表名一致。
变更事件键
表的变更事件键包含若干字段,每个字段对应事件创建时该表主键中存在的一个列。
例如,考虑 defaultdb 数据库中 orders 表的如下定义:
示例表
CREATE TABLE orders (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
customer_name STRING NOT NULL,
amount DECIMAL NOT NULL,
status STRING DEFAULT 'pending',
created_at TIMESTAMP DEFAULT now()
);变更事件键示例
根据前面 orders 表的定义,该表的每个变更事件都使用相同的键结构。以下 JSON 表示此键结构:
{
"schema": {
"type": "struct",
"name": "cockroachdb.public.orders.Key",
"optional": false,
"fields": [
{
"type": "string",
"optional": false,
"field": "id"
}
]
},
"payload": {
"id": "5f8a1c2e-3b4d-4e6f-8a9b-1c2d3e4f5a6b"
}
}更改事件值
考虑与更改事件键示例中相同的示例 orders 表:
CREATE TABLE orders (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
customer_name STRING NOT NULL,
amount DECIMAL NOT NULL,
status STRING DEFAULT 'pending',
created_at TIMESTAMP DEFAULT now()
);基于前面的 orders 表定义,以下示例展示了下列类型变更事件值的典型结构:
create 事件
以下示例展示了连接器为 orders 表中插入数据的操作生成的变更事件的值部分:
{
"schema": { ... },
"payload": {
"before": null, (1)
"after": { (2)
"id": "5f8a1c2e-3b4d-4e6f-8a9b-1c2d3e4f5a6b",
"customer_name": "Alice",
"amount": "99.99",
"status": "pending",
"created_at": "2026-01-15T10:30:00"
},
"source": { (3)
"version": "3.6.3.Final",
"connector": "cockroachdb",
"name": "cockroachdb_connector",
"ts_ms": 1705312200000,
"ts_us": 1705312200000000,
"ts_ns": 1705312200000000000,
"snapshot": "false",
"db": "defaultdb",
"sequence": "[]",
"schema": "public",
"table": "orders",
"cluster": "cockroachdb_connector",
"resolved_ts": null,
"ts_hlc": null
},
"op": "c", (4)
"ts_ms": 1705312200123, (5)
"ts_us": 1705312200123456,
"ts_ns": 1705312200123456789
}
}表 3. create 事件值字段说明
| 序号 | 字段名 | 说明 |
|---|---|---|
| 1 | before |
可选字段,指定事件发生前行的状态。当 op 字段为 c(create,创建)时,before 字段为 null,因为该变更事件针对的是新内容。 |
| 2 | after |
可选字段,指定事件发生后行的状态。在本例中,after 字段包含新行的 id、customer_name、amount、status 和 created_at 列的值。 |
| 3 | source |
必填字段,描述事件的源元数据。该字段包含的信息可用于将此事件与其他事件进行比较,涉及事件的来源、事件发生的顺序,以及事件是否属于同一事务。源元数据包括: - Debezium 版本 - 连接器类型和名称 - 包含新行的数据库和表 - 该事件是否属于快照(初始扫描)的一部分 - 模式(schema)名称 - 逻辑集群名称 - 已解析的时间戳(用于一致性跟踪) - 混合逻辑时钟(HLC)时间戳。HLC 是 CockroachDB 的内部时间戳格式。 |
| 4 | op |
必填字符串,描述导致连接器生成该事件的操作类型。在本例中,c 表示该操作创建了一行。有效值为: - r = 读取(初始扫描 / 快照) - c = 创建 - u = 更新 - d = 删除 |
| 5 | ts_ms、ts_us、ts_ns |
可选字段,显示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。 在 source 对象中,ts_ms 表示数据库中进行更改的时间。通过比较 payload.source.ts_ms 的值与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。 |
update 事件
示例 orders 表中更新事件的事件记录值部分,与该表的 create 事件具有相同的模式。同样,事件值的负载(payload)也具有相同的结构。不过,update 事件的事件值负载中包含的值有所不同。以下示例展示了连接器为 orders 表中的一次更新所生成的变更事件值部分的 JSON:
{
"schema": { ... },
"payload": {
"before": null, (1)
"after": { (2)
"id": "5f8a1c2e-3b4d-4e6f-8a9b-1c2d3e4f5a6b",
"customer_name": "Alice",
"amount": "99.99",
"status": "shipped",
"created_at": "2026-01-15T10:30:00"
},
"source": { (3)
"version": "3.6.3.Final",
"connector": "cockroachdb",
"name": "cockroachdb_connector",
"ts_ms": 1705312500000,
"ts_us": 1705312500000000,
"ts_ns": 1705312500000000000,
"snapshot": "false",
"db": "defaultdb",
"sequence": "[]",
"schema": "public",
"table": "orders",
"cluster": "cockroachdb_connector",
"resolved_ts": null,
"ts_hlc": null
},
"op": "u", (4)
"ts_ms": 1705312500123,
"ts_us": 1705312500123456,
"ts_ns": 1705312500123456789
}
}| 项目 | 字段名称 | 描述 |
|---|---|---|
| 1 | before | 一个可选字段,包含更新前该行的状态。默认情况下,除非启用了 diff 变更馈送选项(cockroachdb.changefeed.include.diff=true),否则 CockroachDB 变更馈送不会包含 before 状态。当未启用 diff 时,该字段为 null。 |
| 2 | after | 一个可选字段,用于指定事件发生后该行的状态。在此示例中,status 的值已被更改为 shipped。 |
| 3 | source | 必填字段,用于描述事件的来源元数据。source 字段的结构与 create 事件中的字段相同,但部分取值有所不同。 |
| 4 | op | 一个必填字符串,用于描述操作的类型。在 update 事件的值中,op 字段的值为 u,表示该行因更新而发生变更。 |
表 4. update 事件值字段说明
若要在更新事件中包含 before 状态,请在连接器配置中将 cockroachdb.changefeed.include.diff 设置为 true。若要在更新事件中包含某一行的 before 状态,请在连接器配置中将 CockroachDB 变更馈送的 diff 选项(即 cockroachdb.changefeed.include.diff)设置为 true。 |
|---|
delete 事件
delete 变更事件中的值,其 schema 部分与同一表的 create 事件和 update 事件相同。delete 变更事件记录为使用者提供处理行删除所需的信息。
Debezium CockroachDB 连接器生成的事件在设计上可与 Kafka 日志压缩协同工作。日志压缩允许删除部分较旧的消息,只要每个键至少保留最新的那条消息即可。删除较旧的消息可以让 Kafka 回收存储空间,同时确保主题中包含完整的数据集,可用于重新加载基于键的状态。
下面的示例展示了连接器针对 orders 表中的删除事件所生成的变更事件 payload 部分的 JSON:
{
"schema": { ... },
"payload": {
"before": null,
"after": null, (1)
"source": { (2)
"version": "3.6.3.Final",
"connector": "cockroachdb",
"name": "cockroachdb_connector",
"ts_ms": 1705312800000,
"ts_us": 1705312800000000,
"ts_ns": 1705312800000000000,
"snapshot": "false",
"db": "defaultdb",
"sequence": "[]",
"schema": "public",
"table": "orders",
"cluster": "cockroachdb_connector",
"resolved_ts": null,
"ts_hlc": null
},
"op": "d", (3)
"ts_ms": 1705312800123,
"ts_us": 1705312800123456,
"ts_ns": 1705312800123456789
}
}| 项 | 字段名称 | 描述 |
|---|---|---|
| 1 | after | after 字段为 null,表示该行已不复存在。 |
| 2 | source | 必填字段,用于描述事件的源元数据。在 delete 事件值中,source 字段的结构与同一表的 create 和 update 事件相同。 |
| 3 | op | 描述操作类型的必填字符串。op 字段的值为 d,表示该行已被删除。 |
表 5. delete 事件值字段说明
delete 变更事件记录为使用者提供了处理该行移除所需的信息。
墓碑事件
当你删除一行时,delete 事件值仍然可以与日志压缩配合工作,因为 Kafka 可以删除所有具有相同键的较早消息。但是,要让 Kafka 删除所有具有相同键的消息,消息的值必须为 null。为了实现这一点,连接器会在 delete 事件之后发送一个特殊的 tombstone(墓碑)事件,该事件具有相同的键,但值为 null。
数据类型映射
CockroachDB 连接器使用与行所在表结构相同的事件来表示行的变更。事件中的每个列值都有一个对应字段。Debezium 在事件记录中表示值的方式取决于 CockroachDB 列的数据类型。CockroachDB 数据应用以下数据类型映射:
CockroachDB 类型
| CockroachDB 数据类型 | Kafka Connect schema 类型 | 说明 |
|---|---|---|
BOOL、BOOLEAN | BOOLEAN | |
INT2、SMALLINT | INT16 | |
INT4、INT、INTEGER | INT32 | |
INT8、BIGINT、SERIAL | INT64 | |
FLOAT4、REAL | FLOAT32 | |
FLOAT8、DOUBLE PRECISION、FLOAT | FLOAT64 | |
NUMERIC、DECIMAL、DEC | STRING | 表示为字符串以保留任意精度。 |
STRING、VARCHAR、CHAR、CHARACTER VARYING、TEXT | STRING | |
UUID | STRING | 表示为标准 UUID 字符串格式。 |
BYTES、BYTEA、BLOB | BYTES | |
DATE | INT32(io.debezium.time.Date 逻辑类型) | 自 Unix 纪元以来的天数。 |
TIME | INT64(io.debezium.time.MicroTime 逻辑类型) | 自午夜以来的微秒数。 |
TIMETZ | STRING(io.debezium.time.ZonedTime 逻辑类型) | 带时区偏移量的 ISO-8601 时间(例如 14:36:34.873+02:00)。 |
TIMESTAMP | INT64(io.debezium.time.MicroTimestamp 逻辑类型) | 自 Unix 纪元以来的微秒数,按 UTC 解释。 |
TIMESTAMPTZ | STRING(io.debezium.time.ZonedTimestamp 逻辑类型) | 带时区偏移量的 ISO-8601 时间戳(例如 2026-06-15T14:36:34.873Z)。 |
INTERVAL | STRING | CockroachDB 间隔格式。 |
JSONB、JSON | STRING(Json 逻辑类型) | 使用 Debezium Json 逻辑类型表示为 JSON 字符串。 |
INET | STRING | IP 地址字符串。 |
BIT、VARBIT | STRING | 位串表示形式。 |
ARRAY | STRING | JSON 数组格式。 |
ENUM | STRING | 枚举标签字符串。 |
GEOMETRY、GEOGRAPHY | STRING | GeoJSON 或 WKT 格式。 |
VECTOR | FLOAT64 的 ARRAY(DoubleVector) | 与 pgvector 兼容的类型(CockroachDB 24.2+)。使用 Debezium DoubleVector 逻辑类型。 |
表 6. CockroachDB 数据类型的映射关系
默认值
当 CockroachDB 的列指定 DEFAULT 子句时,连接器会将 JDBC 报告的默认表达式解析为与该列的 Kafka Connect 模式相匹配的 Java 类型,并将其设置为该字段的模式默认值。CockroachDB 会在默认表达式后附加 CockroachDB 特有的 :::TYPE 后缀(例如 0:::INT8、'PENDING':::STRING、'[1.0,2.0,3.0]':::VECTOR)。连接器会去除该注解,解包带引号的字面量,并将值转换为该列的 Java 类型。
由函数生成的默认值(例如 current_timestamp()、gen_random_uuid()、now() 和 unique_rowid())不会由连接器计算。对于这些列,连接器会从模式中省略默认值,从而让 CockroachDB 在插入时计算该值。
设置 CockroachDB
在部署 Debezium CockroachDB 连接器之前,请准备好 CockroachDB 环境,以便连接器能够创建和读取 changefeed。总体而言,需要确保所需的组件和网络连通性已经就绪,在集群上启用 rangefeed,并为连接器所使用的数据库用户授予 changefeed 所需的权限。有关启用 rangefeed 和授予权限的信息,请参阅授予权限。
前置条件
CockroachDB 集群
一个支持基于 sink 的 changefeed 的 CockroachDB 集群。无需企业许可证即可在所有 CockroachDB 版本中使用全部 changefeed 功能。
中间 Kafka 集群
CockroachDB 用于发布 changefeed 事件的 Kafka 集群。它可以是 Kafka Connect 所使用的同一个 Kafka 集群,也可以是独立的集群。
网络连通性
CockroachDB 集群必须能够连接到中间 Kafka 集群,以发布 changefeed 事件。Kafka Connect 工作进程必须能够读取 CockroachDB 集群(用于 JDBC)和中间 Kafka 集群(用于消费 changefeed 事件)。
Changefeed 用户账户
一个在目标表上拥有以下权限之一的 CockroachDB 用户:
CHANGEFEED权限(CockroachDB v22.2+)ALL权限admin角色的成员身份
授予 changefeed 权限
启用全集群的 rangefeed,并为连接器所使用的数据库用户授予必要的权限。CockroachDB 会在执行 CREATE CHANGEFEED 时校验 changefeed 权限,如果用户权限不足,将返回准确的、与版本相关的错误消息。
操作步骤
以 CockroachDB 管理员用户身份,运行以下 SQL 命令:
-- Enable rangefeeds (required cluster-wide setting) SET CLUSTER SETTING kv.rangefeed.enabled = true; -- Grant changefeed privilege to a specific user GRANT CHANGEFEED ON TABLE orders TO myuser; GRANT VIEWCLUSTERSETTING TO myuser; -- Or grant on all tables in a database GRANT CHANGEFEED ON ALL TABLES IN DATABASE defaultdb TO myuser;
授予 VIEWCLUSTERSETTING 权限后,连接器便可以验证 kv.rangefeed.enabled 是否设置为 true,这是变更馈送(changefeed)的先决条件。
保护变更馈送接收器的安全
当变更馈送接收器使用 TLS 或双向 TLS(mTLS)时,连接器可以从磁盘加载 PEM 文件,并在启动时自动将其 base64 编码的内容添加到接收器 URI 中。这种自动注入方式免去了手动将大段编码后的证书字符串复制到连接器配置中的麻烦。
设置以下一个或多个属性,以指定连接器用于加密接收器连接的机制:
cockroachdb.changefeed.sink.tls.ca.cert.file
用于验证接收器代理证书的 PEM 编码 CA 证书的文件路径。
cockroachdb.changefeed.sink.tls.client.cert.file
用于与接收器建立双向 TLS 连接的 PEM 编码客户端证书的文件路径。
cockroachdb.changefeed.sink.tls.client.key.file
用于与接收器建立双向 TLS 连接的 PEM 编码客户端私钥的文件路径。
对于每个已设置的属性,连接器会读取该文件,验证其可读且非空,对内容进行 base64 编码,再对结果进行 URL 编码,然后按照所配置的接收器类型所要求的格式,将相应的查询参数追加到 cockroachdb.changefeed.sink.uri 中。对于 Kafka 接收器,追加的参数为 ca_cert=…、client_cert=…、client_key=… 和 tls_enabled=true(只要三个 TLS 文件选项中的任何一个被设置,就会自动添加)。基于文件的值会覆盖接收器 URI 中已内联存在的同名查询参数。
使用 mTLS 保护的 Kafka 接收器的配置示例
{
"cockroachdb.changefeed.sink.type": "kafka",
"cockroachdb.changefeed.sink.uri": "kafka://kafka.example.com:9093",
"cockroachdb.changefeed.sink.tls.ca.cert.file": "/etc/kafka/secrets/ca.pem",
"cockroachdb.changefeed.sink.tls.client.cert.file": "/etc/kafka/secrets/client.pem",
"cockroachdb.changefeed.sink.tls.client.key.file": "/etc/kafka/secrets/client-key.pem"
}连接器在启动时会校验每个已配置的文件路径,如果文件不存在、不可读或为空,会立即失败并退出。目前只有 Kafka 接收器(sink)会使用这些 TLS 参数;对于其他类型的接收器,接收器 URI 会原样透传。
部署
要部署 Debezium CockroachDB 连接器,需要获取连接器插件文件归档包,将 JAR 文件添加到 Kafka Connect 环境中,然后编辑 plugin.path 配置,使其指向这些文件所在的位置。
如果你使用的是不可变容器,请参阅 Kafka 和 Kafka Connect 的 Debezium 容器镜像。
你从 quay.io 获取的 Debezium 容器镜像未经过严格测试或安全分析,仅供测试和评估之用。这些镜像不适用于生产环境。为降低生产部署中的风险,请仅部署由受信任的供应商积极维护并经过潜在漏洞全面测试的容器。 |
|---|
你也可以在 Kubernetes 和 OpenShift 上运行 Debezium。
前置条件
- Kafka 和 Kafka Connect 已部署并正在运行。
操作步骤
- 将归档文件下载到临时目录。
- 将下载的文件复制到 Kafka Connect 环境中的某个目录,例如
/usr/local/share/kafka/plugins。 - 将 JAR 文件解压到该目录中,然后删除下载的归档文件。
- 将包含 JAR 文件的目录添加到 Kafka Connect 的
plugin.path中。 - 配置连接器,并将配置内容添加到你的 Kafka Connect 集群中。
- 重启 Kafka Connect 进程,使其加载新的 JAR 文件。
连接器配置示例
以下示例展示了 CockroachDB 连接器的一种可能配置,该连接器连接到 CockroachDB 集群,并捕获 defaultdb 数据库中 public 模式下所有表的变更。
{
"name": "cockroachdb-connector", (1)
"config": {
"connector.class": "io.debezium.connector.cockroachdb.CockroachDBConnector", (2)
"database.hostname": "cockroachdb-host", (3)
"database.port": "26257", (4)
"database.user": "myuser", (5)
"database.password": "mypassword", (6)
"database.dbname": "defaultdb", (7)
"topic.prefix": "cockroachdb", (8)
"cockroachdb.changefeed.sink.uri": "kafka://kafka-broker:9092", (9)
"snapshot.mode": "initial", (10)
"tasks.max": "1" (11)
}
}| 1 | 连接器注册到 Kafka Connect 服务时使用的名称。 |
|---|---|
| 2 | 此 CockroachDB 连接器类的名称。 |
| 3 | CockroachDB 主机的地址。 |
| 4 | CockroachDB 服务器的端口号(默认:26257)。 |
| 5 | CockroachDB 用户的名称。 |
| 6 | CockroachDB 用户的密码。 |
| 7 | 要连接的 CockroachDB 数据库名称。 |
| 8 | 连接器写入的 Kafka 主题的主题前缀。 |
| 9 | CockroachDB 发布变更流(changefeed)事件所在的 Kafka 集群 URI。 |
| 10 | 快照模式;initial 表示在首次启动时回填现有行。 |
| 11 | 最大任务数。 |
有关可以在连接器配置中指定的属性的更多信息,请参阅 CockroachDB 连接器属性完整列表 。
添加连接器配置
编辑 CockroachDB 连接器配置,指定连接器的运行方式。完成配置后,使用 Kafka Connect API 将该配置应用到集群。
前提条件
- CockroachDB 正在运行且可访问。
- 中间 Kafka 集群正在运行,并且可以从 CockroachDB 集群访问。
- CockroachDB 连接器已安装。
操作步骤
使用 Kafka Connect REST API 提交
POST请求,将连接器配置的 JSON 应用到 Kafka Connect 集群。例如:
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" \ http://localhost:8083/connectors/ -d @connector-config.json
结果
当连接器启动时,它会连接 CockroachDB 数据库,创建变更订阅(changefeed),并开始为行级操作生成数据变更事件,将变更事件记录流式传输到 Kafka 主题中。
监控
Debezium CockroachDB 连接器内置支持 Kafka 和 Kafka Connect 提供的 JMX 指标,以及快照指标和流式指标,这些指标有助于了解连接器在初始扫描期间以及变更事件记录捕获和流式传输期间的行为。
其他资源
自定义 MBean 名称
Debezium 连接器通过连接器的 MBean 名称暴露指标。这些指标针对每个连接器实例,提供有关连接器快照、流式传输和架构历史处理过程行为的数据。
默认情况下,当你部署一个配置正确的连接器时,Debezium 会为每种不同的连接器指标生成一个唯一的 MBean 名称。要查看连接器处理过程的指标,你需要配置可观测性栈来监控其 MBean。但这些默认 MBean 名称取决于连接器配置;配置更改可能导致 MBean 名称发生变化。MBean 名称的变更会破坏连接器实例与 MBean 之间的关联,从而干扰监控活动。在这种情况下,如果你想恢复监控,必须重新配置可观测性栈以使用新的 MBean 名称。
为防止 MBean 名称变更导致的监控中断,你可以配置自定义指标标签。通过在连接器配置中添加 custom.metric.tags 属性来配置自定义指标。该属性接受键值对,其中每个键表示 MBean 对象名称的一个标签,对应的值表示该标签的值。例如:k1=v1,k2=v2。Debezium 会将指定的标签附加到连接器的 MBean 名称中。
为连接器配置 custom.metric.tags 属性后,你就可以配置可观测性栈以检索与指定标签关联的指标。随后,可观测性栈将使用指定的标签(而非易变的 MBean 名称)来唯一标识连接器。之后,如果 Debezium 重新定义了构建 MBean 名称的方式,或者连接器配置中的 topic.prefix 发生变化,指标收集也不会中断,因为指标抓取任务使用指定的标签模式来标识连接器。
使用自定义标签的另一个好处是,你可以使用能够反映数据管道架构的标签,从而按照符合你运维需求的方式来组织指标。例如,你可以指定标签值,用于声明连接器活动的类型、应用程序上下文或数据源,比如 db1-streaming-for-application-abc。如果你指定多个键值对,那么所有指定的键值对都会被追加到连接器的 MBean 名称中。
下面的示例说明了标签如何修改默认的 MBean 名称。
示例 1. 自定义标签如何修改连接器 MBean 名称
默认情况下,CockroachDB 连接器为流式指标使用以下 MBean 名称:
debezium.CockroachDB:type=connector-metrics,context=streaming,server=<topic.prefix>如果你将 custom.metric.tags 的值设置为 database=salesdb-streaming,table=inventory,Debezium 将生成以下自定义 MBean 名称:
debezium.CockroachDB:type=connector-metrics,context=streaming,server=<topic.prefix>,database=salesdb-streaming,table=inventory快照指标
该连接器通过 CockroachDB 变更流的 initial_scan 回填已有行,并围绕该过程驱动标准的快照生命周期,因此快照指标会在初始扫描期间填充。
MBean 为 debezium.CockroachDB:type=connector-metrics,context=snapshot,server=<topic.prefix>。
下表列出了可用于监控 Debezium 快照操作的 JMX 指标,包括行数、表进度、持续时间和队列容量。除非快照操作正在进行中,或者自上次连接器启动以来发生过快照,否则不会暴露快照指标。
| 指标名称 | 描述 |
|---|---|
Snapshotting |
快照操作是否正在进行。 |
TotalTablesToSnapshot |
要进行快照的表的总数。 |
SnapshotCompleted |
快照是否已完成。 |
SnapshotDurationInSeconds |
快照所花费的秒数。 |
RowsScanned |
已扫描的行数。 |
TableId |
要进行快照的表的标识符。 |
TableName |
要进行快照的表的名称。 |
TotalRowsScanned |
已扫描的总行数。 |
SnapshotLockingTime |
快照期间用于获取锁的时间。 |
QueueCapacity |
快照队列的容量。 |
QueueRemainingCapacity |
快照队列的剩余容量。 |
QueueTotalCapacity |
快照队列的总容量。 |
| 属性 | 类型 | 描述 |
|---|---|---|
LastEvent | string | 连接器读取的最后一个快照事件。 |
MilliSecondsSinceLastEvent | long | 自连接器读取并处理最近一个事件以来经过的毫秒数。 |
NumberOfErroneousEvents | long | 记录连接器在快照操作期间识别为错误的变更事件数量。每当连接器在初始快照、增量快照或临时快照过程中遇到无法处理的事件时,该指标就会递增。事件处理失败的原因可能包括:事件格式错误、与模式不兼容,或在转换过程中发生失败。该指标值在连接器任务的整个生命周期内保持有效。如果快照被中断且连接器任务重启,该指标计数将重置为 0。 |
TotalNumberOfEventsSeen | long | 自上次启动或重置以来,此连接器看到的事件总数。 |
NumberOfEventsFiltered | long | 被连接器上配置的包含/排除列表过滤规则所过滤的事件数量。 |
CapturedTables | string[] | 连接器捕获的表的列表。 |
QueueTotalCapacity | int | 用于在快照器和主 Kafka Connect 循环之间传递事件的队列长度。 |
QueueRemainingCapacity | int | 用于在快照器和主 Kafka Connect 循环之间传递事件的队列的剩余容量。 |
TotalTableCount | int | 被纳入快照的表的总数。 |
RemainingTableCount | int | 快照尚未复制的表的数量。 |
SnapshotRunning | boolean | 快照是否已启动。 |
SnapshotPaused | boolean | 快照是否已暂停。 |
SnapshotAborted | boolean | 快照是否已中止。 |
SnapshotCompleted | boolean | 快照是否已完成。 |
SnapshotSkipped | boolean | 快照是否被跳过。 |
SnapshotDurationInSeconds | long | 快照到目前为止所花费的总秒数,即使尚未完成也计入。同时包含快照处于暂停状态的时间。 |
SnapshotPausedDurationInSeconds | long | 快照被暂停的总秒数。如果快照被暂停多次,暂停时间将累计。 |
RowsScanned | Map<String, Long> | 包含快照中每张表已扫描行数的映射。在处理过程中,表会被逐步添加到该映射中。每扫描 10,000 行以及完成一张表时都会更新。 |
TableChunkCounts | Map<String, Long> | 在使用基于分块的多线程快照时,包含快照中每张表的分块数量的映射。 |
TableChunksCompletedCounts | Map<String, Long> | 在使用基于分块的多线程快照时,包含快照中每张表已完成分块数量的映射。 |
MaxQueueSizeInBytes | long | 队列的最大缓冲区大小(以字节为单位)。如果 max.queue.size.in.bytes 被设置为正的 long 值,则可使用此指标。 |
CurrentQueueSizeInBytes | long | 队列中当前记录的容量(以字节为单位)。 |
下表列出了连接器执行增量快照时可用的其他 JMX 指标,其中包括可用于跟踪快照进度的分块和表边界标识符。
| 属性 | 类型 | 描述 |
|---|---|---|
ChunkId | string | 当前快照分块的标识符。 |
ChunkFrom | string | 定义当前分块的主键集下界。 |
ChunkTo | string | 定义当前分块的主键集上界。 |
TableFrom | string | 当前正在快照的表的主键集下界。 |
TableTo | string | 当前正在快照的表的主键集上界。 |
流处理指标
MBean 为 debezium.CockroachDB:type=connector-metrics,context=streaming,server=<topic.prefix>。
下表列出了可用于监控 Debezium 流处理操作的 JMX 指标,包括按类型统计的事件数量、相对于源的延迟、队列容量以及连接状态。
| 属性 | 类型 | 描述 |
|---|---|---|
LastEvent | string | 连接器读取到的最后一个流式事件。 |
MilliSecondsSinceLastEvent | long | 自连接器读取并处理最近一个事件以来经过的毫秒数。 |
NumberOfErroneousEvents | long | 记录连接器在流式传输过程中识别为错误的变更事件数量。在流式会话的生命周期内,每当连接器遇到无法处理的事件时,该指标就会递增。事件处理失败的原因可能包括格式错误、与模式不兼容,或在转换过程中出现失败。该指标的值在连接器任务的生命周期内保持不变。连接器重启后,该指标计数会重置为 0。 |
TotalNumberOfEventsSeen | long | 自上次启动连接器或重置指标以来,源数据库报告的数据变更事件总数。表示 Debezium 需要处理的数据变更工作负载。 |
TotalNumberOfCreateEventsSeen | long | 自上次启动或重置指标以来,连接器处理的创建事件总数。 |
TotalNumberOfUpdateEventsSeen | long | 自上次启动或重置指标以来,连接器处理的更新事件总数。 |
TotalNumberOfDeleteEventsSeen | long | 自上次启动或重置指标以来,连接器处理的删除事件总数。 |
NumberOfEventsFiltered | long | 被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。 |
NumberOfUnchangedEventsSkipped | long | 自上次启动连接器或重置指标以来,因被监视列未发生变化而跳过的更新事件数量。如果 skip.messages.without.change 为 false,则默认值为 -1。对于 CockroachDB 以及其他不支持跳过未变更事件的连接器,即使将该属性设置为 true,该值也可能始终保持为 0。 |
CapturedTables | string[] | 连接器捕获的表列表。 |
QueueTotalCapacity | int | 用于在流式程序与主 Kafka Connect 循环之间传递事件的队列长度。 |
QueueRemainingCapacity | int | 用于在流式程序与主 Kafka Connect 循环之间传递事件的队列的剩余可用容量。 |
Connected | boolean | 表示连接器当前是否已连接到数据库服务器的标志。 |
MilliSecondsBehindSource | long | 最后一个变更事件的时间戳与连接器处理该事件之间相差的毫秒数。该值会包含运行数据库服务器和连接器的机器之间的时钟差异。 |
MilliSecondsBehindSourceMinValue | long | 连接器运行期间观测到的相对于源的最小滞后毫秒数。 |
MilliSecondsBehindSourceMaxValue | long | 连接器运行期间观测到的相对于源的最大滞后毫秒数。 |
MilliSecondsBehindSourceAverageValue | double | 连接器运行期间所有观测值计算得出的相对于源的平均滞后毫秒数。 |
MilliSecondsBehindSourceP50 | double | 相对于源的滞后毫秒数的第 50 百分位数(中位数)。由于该指标受异常值影响较小,因此比平均值更能稳健地衡量典型滞后。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
MilliSecondsBehindSourceP95 | double | 相对于源的滞后毫秒数的第 95 百分位数。该指标表示 95% 的滞后测量值低于此值,可用于识别尾部延迟并设定 SLA 阈值。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
MilliSecondsBehindSourceP99 | double | 相对于源的滞后毫秒数的第 99 百分位数。该指标表示 99% 的滞后测量值低于此值,可用于了解最坏情况下的性能表现。当 statistics.metrics.enabled 设置为 true(默认值)时可用。 |
连接器配置属性
Debezium CockroachDB 连接器提供了许多配置属性,你可以利用它们来实现应用所需的正确连接器行为。许多属性都有默认值。属性信息按如下方式组织:
除非提供了默认值,否则以下配置属性均为必填。
| 属性 | 默认值 | 描述 |
|---|---|---|
name | 无默认值 | 连接器的唯一名称。注册连接器时使用的名称必须唯一,如果重复使用连接器名称,注册将会失败。所有 Kafka Connect 连接器都要求设置此属性。 |
connector.class | 无默认值 | 连接器对应的 Java 类名。对于 CockroachDB 连接器,始终使用 io.debezium.connector.cockroachdb.CockroachDBConnector。 |
tasks.max | 1 | 连接器可以创建的最大任务数。 |
database.hostname | 无默认值 | CockroachDB 数据库服务器的 IP 地址或主机名。 |
database.port | 26257 | CockroachDB 数据库服务器的端口号(整数)。CockroachDB 默认的 SQL 端口为 26257。 |
database.user | 无默认值 | 用于连接数据库的 CockroachDB 数据库用户名。 |
database.password | 无默认值 | 用于连接数据库的 CockroachDB 数据库用户密码。 |
database.dbname | 无默认值 | 要从中流式传输变更的 CockroachDB 数据库名称。 |
topic.prefix | 无默认值 | 主题前缀,为特定的 CockroachDB 数据库服务器或集群提供命名空间。主题前缀在所有其他连接器中应当是唯一的,因为它将作为所有接收该连接器发出事件的 Kafka 主题名称的前缀。主题前缀中只能使用字母、数字、连字符、点和下划线。 |
cockroachdb.changefeed.sink.uri | 无默认值 | CockroachDB 发布变更流(changefeed)事件的中间 Kafka 集群的 URI。格式:kafka://host:port。此属性为必填项。 |
表 7. 连接器必填配置属性
以下属性用于配置到 CockroachDB 数据库的 JDBC 连接。
| 属性 | 默认值 | 描述 |
|---|---|---|
database.sslmode | prefer | 指定 Debezium 是否与 CockroachDB 建立加密连接。可设置为以下选项之一:disable、allow、prefer、require、verify-ca、verify-full。详情请参阅 CockroachDB 认证文档。 |
database.sslrootcert | 无默认值 | 包含用于验证数据库服务器的受信任根证书的文件。 |
database.sslcert | 无默认值 | 包含客户端 SSL 证书的文件。 |
database.sslkey | 无默认值 | 包含客户端 SSL 私钥的文件。 |
database.sslpassword | 无默认值 | 用于从 database.sslkey 指定的文件中访问客户端私钥的密码。 |
database.tcpKeepAlive | true | 指定是否启用 TCP keep-alive 探测以避免 TCP 连接断开。 |
database.on.connect.statements | 无默认值 | 指定连接器在建立 JDBC 连接时执行的、以分号分隔的 SQL 语句列表。如需表示字面分号而非分隔符,请使用双分号(;;)。 |
connection.timeout.ms | 30000 | 指定连接器等待与数据库建立连接的时间(以毫秒为单位)。 |
connection.retry.delay.ms | 1000 | 连接重试尝试之间的基础延迟(以毫秒为单位)。实际延迟会乘以尝试次数(线性退避)。 |
connection.max.retries | 3 | 连接器在放弃连接尝试之前重试连接的最大次数。 |
connection.validation.timeout.seconds | 5 | 验证现有 JDBC 连接是否仍然可用的超时时间(以秒为单位)。 |
表 8. 连接配置属性
以下属性用于配置 CockroachDB 变更 feed 的行为。
| 属性 | 默认值 | 描述 |
|---|---|---|
cockroachdb.changefeed.resolved.interval | 10s | 已解析时间戳消息的发送间隔。格式示例:10s、1m 等。已解析时间戳用于偏移量跟踪。 |
cockroachdb.changefeed.include.updated | false | 指定 UPDATE 事件是否包含哪些列被更新的信息。 |
cockroachdb.changefeed.include.diff | false | 指定是否在变更流事件中包含更新前后的差异信息。启用后,UPDATE 事件会在 before 字段中包含该行的先前状态。 |
cockroachdb.changefeed.enriched.properties | source | 以逗号分隔的增强封装(enriched envelope)属性列表,将原样传递给 CockroachDB 的 CREATE CHANGEFEED 语句。连接器始终使用增强封装,因此该设置对其创建的每个变更流都生效。有效值集合由 CockroachDB 定义,目前为 source 和 schema,并在创建变更流时进行校验。连接器通过 JDBC 发现的表元数据来推导每个事件的模式,并根据事件到达的主题来识别该事件所属的表。连接器不会读取 schema 块,因为它所包含的信息与连接器已有的信息重复;schema 会被接受并透传,供中间主题的其他消费者使用,但连接器本身并不使用它。默认值为 source,因为该块标识了源数据库、模式、表以及提交时间戳。当你需要检查中间主题,或有其他工具同时消费这些主题时,这些信息很有用。 |
cockroachdb.changefeed.cursor | now | 变更流的起始游标位置。使用 now 表示从当前时间开始,也可以指定一个绝对时间戳。当连接器从存储的偏移量重启时,它会使用存储的已解析时间戳,而非该值。 |
cockroachdb.changefeed.batch.size | 1000 | 变更流处理的批处理大小。 |
cockroachdb.changefeed.poll.interval.ms | 100 | 变更流处理的轮询间隔(毫秒)。 |
cockroachdb.changefeed.sink.type | kafka | 变更流事件的接收器(sink)类型。当前支持:kafka。 |
cockroachdb.changefeed.sink.topic.prefix | 空 | 连接器添加到中间变更流主题名称前的前缀字符串。连接器会原样使用你指定的前缀字符串,因此生成的主题名称为 prefix__database.schema.table。如需自定义分隔符,请自行包含在前缀中:例如,crdb. 会生成 crdb.mydb.public.orders,env-prod- 会生成 env-prod-mydb.public.orders。如果未为该属性提供值,连接器将默认使用 topic.prefix 的值并在其后加上一个点。 |
cockroachdb.changefeed.max.tables.per.changefeed | 0 | 单个 CockroachDB 变更流中包含的最大表数。默认值 0 表示将所有已配置的表放在同一个变更流中。将该值设置为正数,可指示连接器把捕获的表拆分到多个变更流中。指定的值即为每个生成的变更流所包含的最大表数。如果 CockroachDB 返回关于单个变更流监听表数过多的警告,则应将捕获的表拆分到多个变更流中。较小的值可降低耦合度,但会创建更多变更流作业。请指定一个既能优化性能、又符合 CockroachDB 关于限制每个集群变更流作业数量建议的值。 |
cockroachdb.changefeed.sink.options | 无默认值 | 以逗号分隔的接收器附加选项列表,格式为 key=value。 |
cockroachdb.changefeed.sink.tls.ca.cert.file | 无默认值 | 连接器用于验证变更流接收器服务器证书的 PEM 编码 CA 证书文件路径。设置该属性后,连接器会读取文件、对内容进行 base64 编码,并按所配置的接收器类型所要求的格式将 CA 证书追加到接收器 URI 中(例如,Kafka 接收器使用 ca_cert=…)。指定的值会覆盖 cockroachdb.changefeed.sink.uri 中存在的任何等效查询参数。有关更多信息,参见保护变更流接收器的安全。 |
cockroachdb.changefeed.sink.tls.client.cert.file | 无默认值 | 用于与变更流接收器建立双向 TLS 的 PEM 编码客户端证书文件路径。设置该属性后,连接器会读取文件、对内容进行 base64 编码,并按所配置的接收器类型所要求的格式将客户端证书追加到接收器 URI 中(例如,Kafka 接收器使用 client_cert=…)。指定的值会覆盖 cockroachdb.changefeed.sink.uri 中存在的任何等效查询参数。有关更多信息,参见保护变更流接收器的安全。 |
| cockroachdb.changefeed.sink.tls.client.key.file | 无默认值 | 用于与变更流接收器建立双向 TLS 的 PEM 编码客户端私钥文件路径。设置后,连接器会读取文件、对内容进行 base64 编码,并按所配置的接收器类型所要求的格式将客户端私钥追加到接收器 URI 中(例如,Kafka 接收器使用 client_key=…)。指定的值会覆盖 cockroachdb.changefeed.sink.uri 中存在的任何等效查询参数。设置这三个接收器 TLS 文件选项中的任何一个,也会隐含地在接收器 URI 上设置 `tls
表 9. Changefeed 配置属性
设置以下属性,可配置连接器的快照与 schema 行为。
| 属性 | 默认值 | 描述 |
|---|---|---|
snapshot.mode | initial | 指定连接器启动时执行快照的条件。CockroachDB 使用原生 changefeed 的 initial_scan 来回填现有行。有关每种快照模式如何映射到 CockroachDB 的 initial_scan 选项,请参见快照。可设置以下选项之一:always、initial、initial_only、no_data、never、when_needed、configuration_based、custom。 |
snapshot.isolation.mode | serializable | 指定连接器使用的事务隔离级别。可设置以下选项之一:serializable(CockroachDB 默认值)、read_committed。 |
snapshot.locking.mode | none | 指定连接器在执行模式快照时如何对表加锁。可设置以下选项之一:shared、none、custom。由于 CockroachDB 中基于 changefeed 的快照不需要表锁,因此默认且推荐的设置是 none。 |
schema.include.list | 无默认值 | 可选的以逗号分隔的正则表达式列表,用于匹配你想要捕获其更改的模式(schema)名称。未包含在 schema.include.list 中的任何模式名称,其更改都不会被捕获。默认情况下,连接器会捕获所有非系统模式中的更改。为匹配模式名称,Debezium 会将你指定的正则表达式作为锚定正则表达式应用,即该表达式会与模式的完整名称字符串进行匹配,而不会匹配模式名称中可能出现的子字符串。如果你在配置中包含此属性,请勿同时设置 schema.exclude.list 属性。 |
schema.exclude.list | 无默认值 | 可选的以逗号分隔的正则表达式列表,用于匹配你不想捕获其更改的模式(schema)名称。名称未包含在 schema.exclude.list 中的任何模式,其更改都会被捕获,系统模式除外。为匹配模式名称,Debezium 会将你指定的正则表达式作为锚定正则表达式应用,即该表达式会与模式的完整名称字符串进行匹配,而不会匹配模式名称中可能出现的子字符串。如果你在配置中包含此属性,请勿同时设置 schema.include.list 属性。 |
read.only | false | 指定连接器是否使用替代方法向 Debezium 传递信号,而不是写入信号表。 |
status.update.interval.ms | 10000 | 指定连接器向服务器发送状态更新的时间间隔,单位为毫秒。 |
表 10. 连接器配置属性
以下 高级 配置属性的默认值适用于大多数场景,因此很少需要在连接器配置中进行指定。
| 属性 | 默认值 | 描述 |
|---|---|---|
schema.name.adjustment.mode | none | 指定如何调整架构(schema)名称,以兼容连接器所使用的消息转换器。可设置为以下选项之一:none、avro、avro_unicode。 |
field.name.adjustment.mode | none | 指定如何调整字段名称,以兼容连接器所使用的消息转换器。可设置为以下选项之一:none、avro、avro_unicode。详情参见 Avro 命名。 |
unavailable.value.placeholder | __debezium_unavailable_value | 不可用值的占位符(例如,用于被外存(toasted)的列)。 |
heartbeat.interval.ms | 0 | 指定连接器向 __debezium-heartbeat.<topic.prefix> Kafka 主题发送心跳记录的时间间隔(以毫秒为单位)。设置为 0(默认值)可禁用心跳记录。即使禁用,已解析的时间戳仍会推进连接器的内部偏移量。更多信息参见心跳。 |
表 11. 高级连接器配置属性
出现问题时的行为
Debezium 是一个分布式系统,用于捕获多个上游数据库中的所有变更;它绝不会遗漏或丢失任何事件。当系统正常运行并得到妥善管理时,Debezium 会为每条变更事件记录提供 一次且仅一次(exactly once)的投递。
如果发生故障,系统的设计会防止任何事件丢失。但是,在从故障中恢复的过程中,可能会重复投递某些变更事件。在这种异常情况下,Debezium 与 Kafka 一样,会为变更事件提供 至少一次(at least once)的投递。
本节余下的内容将介绍 Debezium 如何处理各类故障和问题。
配置和启动错误
在以下情况下,连接器在尝试启动时会失败,在日志中报告错误或异常,然后停止运行:
- 连接器的配置无效。
- 连接器无法使用指定的连接参数成功连接到 CockroachDB。
- 连接器无法创建 changefeed,因为用户缺少足够的权限(
CHANGEFEED、ALL或admin角色成员资格)。
在这些情况下,错误消息中会包含问题的详细信息,并可能提供建议的解决方法。修正配置或解决 CockroachDB 的问题后,请重新启动连接器。
CockroachDB 变得不可用
连接器运行期间,CockroachDB 集群可能因各种原因而变得不可用。由于 CockroachDB 是分布式数据库,集群会自动处理单个节点的故障。如果整个集群变得不可用,changefeed 将停止产生事件。集群恢复后,连接器会从最后存储的偏移量(offset)处继续处理。
连接器包含重试逻辑,并提供可配置的参数(connection.max.retries、connection.retry.delay.ms)来处理瞬时连接故障,包括 CockroachDB 特有的错误,如序列化失败(SQL 状态 40001)和连接错误(SQL 状态 08xxx)。
Kafka Connect 进程正常停止
假设 Kafka Connect 以分布式模式运行,并且一个 Kafka Connect 进程正常停止。在关闭该进程之前,Kafka Connect 会将该进程的连接器任务迁移到该组中的另一个 Kafka Connect 进程。新的连接器任务会从先前任务停止的位置继续处理。在连接器任务正常停止并在新进程上重新启动期间,处理会有短暂的延迟。
Kafka Connect 进程崩溃
如果 Kafka 连接器进程意外停止,它所运行的任何连接器任务都会随之终止,且不会记录最近已处理的偏移量。当 Kafka Connect 以分布式模式运行时,Kafka Connect 会在其他进程中重新启动这些连接器任务。但是,CockroachDB 连接器会从先前进程已记录的最后一个偏移量处恢复。这意味着新替换的任务可能会重新生成崩溃之前刚刚处理过的部分变更事件。重复事件的数量取决于偏移量刷新周期以及崩溃前数据变更的量。
由于故障恢复期间有可能出现重复事件,因此消费者应始终对部分重复事件有所预期。
Kafka 变得不可用
在连接器生成变更事件的过程中,Kafka Connect 框架会使用 Kafka producer API 将这些事件记录到 Kafka 中。Kafka Connect 会按照你在配置中指定的频率,定期记录这些变更事件中出现的最新偏移量。如果 Kafka broker 变得不可用,运行连接器的 Kafka Connect 进程会反复尝试重新连接到 Kafka broker。换句话说,连接器任务会暂停,直到重新建立连接为止,届时连接器将从上次中断的位置继续运行。
连接器被长时间停止
如果连接器被正常停止,数据库仍可继续使用。连接器重新启动后,会从中断处继续流式传输变更。也就是说,它会为连接器停止期间发生的所有数据库变更生成变更事件记录。
CockroachDB 的 垃圾回收 TTL(默认:4 小时)限制了 changefeed 可以回溯启动的时间范围。如果连接器停止的时间超过了配置的 TTL 间隔,存储的偏移量可能已不再有效。在这种情况下,请使用 snapshot.mode=when_needed,以便在偏移量过期时自动重新执行快照。
启用调试日志
要诊断连接器的行为,请在 Kafka Connect worker 的日志配置中,为连接器的一个或多个 logger 启用 DEBUG(或 TRACE)日志级别。连接器使用 SLF4J,因此标准的 connect-log4j.properties 机制(或运行时的 /admin/loggers REST 端点)可以控制日志级别。
下表列出了用于排查问题时最有用的 logger 名称。
| 日志记录器 | 追踪内容 |
|---|---|
io.debezium.connector.cockroachdb | 顶层连接器生命周期:任务启动/停止、配置、信号处理。 |
io.debezium.connector.cockroachdb.CockroachDBSchema | 架构发现(发现了哪些表和 schema、哪些被 schema.include.list / table.include.list 过滤掉了),以及架构变更时的逐表刷新。 |
io.debezium.connector.cockroachdb.CockroachDBStreamingChangeEventSource | Changefeed 查询构建、逐事件分发、已解析时间戳的推进,以及 Kafka 消费者的接线。 |
io.debezium.connector.cockroachdb.CockroachDBSnapshotChangeEventSource | 快照决策以及向 CockroachDB initial_scan 模式的委托。 |
io.debezium.connector.cockroachdb.CockroachDBDefaultValueConverter | 列 DEFAULT 表达式的解析;当默认表达式无法映射到该列的 Java 类型时,会以 DEBUG 级别记录日志。 |
io.debezium.connector.cockroachdb.connection.CockroachDBConnection | JDBC 连接尝试、重试以及成功连接。 |
例如,要启用架构发现和流式传输的跟踪日志,请将以下内容添加到 connect-log4j.properties 中:
log4j.logger.io.debezium.connector.cockroachdb.CockroachDBSchema=DEBUG
log4j.logger.io.debezium.connector.cockroachdb.CockroachDBStreamingChangeEventSource=DEBUG或者在运行时针对正在运行的 Kafka Connect 工作进程设置级别:
curl -X PUT -H "Content-Type: application/json" \
--data '{"level":"DEBUG"}' \
http://localhost:8083/admin/loggers/io.debezium.connector.cockroachdb限制
仅支持 Kafka sink
该连接器目前仅支持将 Kafka 作为中间 changefeed sink。对 webhook、Pub/Sub 以及云存储 sink 的支持已在计划中(参见 debezium/dbz#1632)。
更新事件中的前值
默认情况下,CockroachDB 的 changefeed 不包含行的前一状态。若要在更新事件中包含 before 字段,请设置 cockroachdb.changefeed.include.diff=true 以启用 diff changefeed 选项。
评论
登录后参与评论
KnowForge