卡珊德拉
Debezium Cassandra 连接器
Cassandra 连接器可以监控 Cassandra 集群并记录所有行级别的变更。该连接器必须部署在 Cassandra 集群中每个节点的本地。连接器首次连接到 Cassandra 节点时,会对所有键空间中启用了 CDC 的表执行快照。连接器还会读取写入 Cassandra 提交日志的变更,并生成相应的插入、更新和删除事件。每张表的所有事件都记录在单独的 Kafka 主题中,应用程序和服务可以方便地消费这些事件。
有关与此连接器兼容的 Cassandra 版本信息,请参阅 Debezium 版本概览。
概述
Cassandra 是一款开源的 NoSQL 数据库。与大多数数据库类似,Cassandra 的写入路径首先是立即将变更记录到其提交日志中。提交日志存储在每个节点的本地,记录对该节点进行的每一次写入。
自 Cassandra 3.0 起,引入了变更数据捕获(CDC)功能。通过将表属性设置为 cdc=true,可以在表级别启用 CDC 功能。启用后,任何包含启用了 CDC 的表的数据的提交日志,在丢弃时都会被移动到 cassandra.yaml 中指定的 CDC 目录。
Cassandra 连接器驻留在每个 Cassandra 节点上,并监控 cdc_raw 目录中的变更。它在检测到变更后会处理所有本地提交日志段,为提交日志中的每一行插入、更新和删除操作生成变更事件,将每张表的所有变更事件发布到单独的 Kafka 主题中,最后从 cdc_raw 目录中删除提交日志。最后这一步非常重要,因为一旦启用 CDC,Cassandra 本身将无法清除提交日志。如果 cdc_free_space_in_mb 被填满,对启用了 CDC 的表的写入将会被拒绝。
连接器具备容错能力。当连接器读取提交日志并生成事件时,它会将每个提交日志段的文件名和位置与每个事件一起记录下来。如果连接器因任何原因停止运行(包括通信故障、网络问题或崩溃),重启后它会从上次停止的地方继续读取提交日志。这同样适用于快照:如果连接器停止时快照尚未完成,重启后它将开始新的快照。稍后我们将讨论出现问题时连接器的行为。
| Cassandra 与其他 Debezium 连接器不同,它并非基于 Kafka Connect 框架实现。相反,它是一个单独的 JVM 进程,旨在部署在每个 Cassandra 节点上,并通过 Kafka 生产者将事件发布到 Kafka。 |
|---|
Cassandra 连接器目前不支持以下特性。由这些特性引起的变更将被忽略:
- 集合类型列的 TTL
- 范围删除
- 静态列
- 触发器
- 物化视图
- 二级索引
- 轻量级事务
设置 Cassandra
在使用 Debezium Cassandra 连接器监控 Cassandra 集群中的变更之前,必须在节点级别和表级别启用 CDC。
在节点上启用 CDC
要启用 CDC,请在 cassandra.yaml 中更新以下 CDC 配置:
cdc_enabled: true其他 CDC 配置的默认值如下:
cdc_raw_directory: $CASSANDRA_HOME/data/cdc_raw
cdc_free_space_in_mb: 4096
cdc_free_space_check_interval_ms: 250cdc_enabled在节点级别启用或禁用 CDC 操作cdc_raw_directory决定所有对应 memtable 刷新完成后,提交日志段(commit log segments)被移动到的目标位置cdc_free_space_in_mb是用于存储提交日志段的最大容量,默认值取 4096 MB 与卷空间的 1/8 中的较小者cdc_free_space_check_interval_ms是重新计算cdc_raw_directory所占空间的频率,以避免在已达容量上限时白白浪费 CPU 周期
在表上启用 CDC
在 Cassandra 节点上启用 CDC 之后,还必须通过 CREATE TABLE 或 ALTER TABLE 命令显式地在每张表上启用 CDC。例如:
CREATE TABLE foo (a int, b text, PRIMARY KEY(a)) WITH cdc=true;
ALTER TABLE foo WITH cdc=true;Cassandra 连接器的工作原理
本节详细介绍 Cassandra 连接器如何执行快照、如何将提交日志事件转换为 Debezium 更改事件、如何处理提交日志的生命周期、如何将事件记录到 Kafka、如何管理模式演变,以及在出现问题时的表现。
快照
当 Cassandra 连接器首次在某个 Cassandra 节点上启动时,默认会对集群执行初始快照。这是默认模式,因为大多数情况下 CDC 是在非空表上启用的,而且提交日志并不包含完整的历史记录。
快照读取器会发出 SELECT 语句来查询表中的所有列。Cassandra 允许在全局或语句级别设置一致性级别。对于快照,默认在语句级别将一致性级别设置为 ALL,以提供最高的一致性。这意味着如果在快照期间有一个节点宕机,快照将无法继续,需要在该节点恢复上线后重新执行快照。你可以将快照的一致性级别调整为较低的级别以提高可用性,前提是你要了解这与一致性之间的权衡。
在 Cassandra 3.X 中,无法严格从本地 Cassandra 节点读取。从 Cassandra 4.0 开始,将新增 NODE_LOCAL 一致性级别。这将使 Cassandra 连接器仅从其所在的节点读取(这与提交日志的处理方式一致)。 |
|---|
与关系型数据库不同,快照期间不会施加读锁,因此在快照过程中对 Cassandra 的写入不会被阻塞。如果在快照期间被查询的数据被其他客户端修改,这些更改可能会反映在快照结果集中。
如果连接器在快照完成之前失败或停止,连接器在重启时会开始新的快照。在默认的快照模式(initial)下,一旦连接器完成初始快照,就不会再执行任何额外的快照。唯一的例外是在连接器重启期间:如果某张表启用了 CDC,随后连接器被重启,那么该表将被再次快照。
第二种快照模式(always)允许连接器在必要时执行快照。它会定期检查新启用 CDC 的表,并在检测到后立即对这些表进行快照。
第三种快照模式(never)确保连接器从不执行快照。当新连接器以这种方式配置时,它只会读取 CDC 目录中的提交日志(commit log)。这不是默认行为,因为以这种模式(不进行快照)启动新连接器要求提交日志包含所有启用 CDC 的表的完整历史记录,而实际情况往往并非如此。该模式的另一个使用场景是:如果已经有一个连接器在执行快照,你可以禁用其他连接器上的快照功能,以避免重复工作。
读取提交日志
Cassandra 连接器绝大部分时间通常都花在读取 Cassandra 节点上的本地提交日志上。在 Cassandra 4.0 中,每次段(segment)执行 fsync 时,都会更新索引文件以反映最新的偏移量。这消除了 Cassandra 3.X 中 CDC 功能的处理延迟,可以通过将配置 commit.log.real.time.processing.enabled 设置为 true 在 Cassandra 4 Debezium 连接器中启用该功能。索引文件的轮询频率由 commit.log.marked.complete.poll.interval.ms 决定。
提交日志的二进制数据由 Cassandra 的 CommitLogReader 和 CommitLogReadHandler 反序列化。每个反序列化后的对象在 Cassandra 中称为一个 mutation(变更)。一个 mutation 包含一个或多个变更事件。
随着 Cassandra 连接器读取提交日志,它会将日志事件转换为 Debezium 的 创建、更新 或 删除 事件,这些事件包含在提交日志中找到该事件的位置。Cassandra 连接器使用 Kafka Connect 转换器对这些变更事件进行编码,并将它们发布到相应的 Kafka 主题中。
提交日志的局限性
Cassandra 的提交日志存在一系列局限性,这些局限性对于正确解读 CDC 事件至关重要:
只有在提交日志写满时,才会被写入
cdc_raw目录,此时该提交日志会被刷新/丢弃。这意味着事件被记录与被捕获之间存在延迟。单个 Cassandra 节点上的提交日志并不反映整个集群的所有写入,它们只反映存储在该节点上的写入。这就是为什么有必要监控 Cassandra 集群中所有节点上的变更。然而,由于复制因子的存在,这也意味着这些事件的下游消费者必须处理去重问题。
对单个 Cassandra 节点的写入会在到达时被记录下来。但是,这些事件的到达顺序可能与它们的发出顺序不一致。这些事件的下游消费者必须理解并实现类似于 Cassandra 读取路径的逻辑,才能获得正确的输出。
表的结构变更不会被记录在提交日志中,只有数据变更才会被记录。因此,结构变更是以尽力而为的方式被检测的。为避免数据丢失,建议在进行结构变更期间暂停对该表的写入。
Cassandra 不执行写前读取,因此提交日志不会记录变更行中每一列的值,它只记录被修改列的值(分区键列除外,因为它们在 Cassandra DML 命令中是必需的,所以始终会被记录)。
由于 CQL 的特性,insert DML 可能导致行的插入或更新;update DML 可能导致行的插入、更新或删除;delete DML 可能导致行的更新或删除。由于查询不会被记录在提交日志中,CDC 事件类型是基于其在关系数据库意义上对行产生的影响来分类的。
TODO:是否有办法确定与实际 Cassandra DML 语句相对应的事件类型?如果有,这种方式是否优于这些事件的语义分类?
管理提交日志的生命周期
默认情况下,Cassandra 连接器会删除已处理的提交日志。不建议在禁用提交日志删除的情况下启动连接器,因为这可能导致磁盘存储膨胀,并阻止对 Cassandra 集群的进一步写入。若要以自定义方式管理提交日志(例如将其上传到云服务商),可以实现 CommitLogTransfer 接口。
主题名称
Cassandra 连接器会将单个表上所有插入、更新和删除操作的事件写入单个 Kafka 主题。Kafka 主题的名称始终采用以下形式:
clusterName.keyspaceName.tableName
其中 clusterName 是通过 topic.prefix 配置属性指定的连接器逻辑名称,keyspaceName 是发生操作所在键空间的名称,tableName 是发生操作所在表的名称。
例如,考虑一个 Cassandra 安装,其中包含一个 inventory 键空间,该键空间中有四张表:products、products_on_hand、customers 和 orders。如果监控此数据库的连接器被赋予逻辑服务器名称 fulfillment,那么该连接器将在这四个 Kafka 主题上产生事件:
fulfillment.inventory.productsfulfillment.inventory.products_on_handfulfillment.inventory.customersfulfillment.inventory.orders
TODO:关于主题名称,使用 clusterName.keyspaceName.tableName 是否合适?还是应该使用 connectorName.keyspaceName.tableName 或 connectorName.clusterName.keyspaceName.tableName?
模式演变
DDL 不会被记录到提交日志中。当表的模式发生更改时,该更改由某个 Cassandra 节点发出,并通过 Gossip 协议传播到其他节点。
Cassandra 中的模式更改将由一个已实现的 SchemaChangeListener 检测到,延迟小于 1 秒,随后它会更新从 Cassandra 加载的模式实例,以及为每张表缓存的 Kafka 键值模式。
请注意,按照当前的模式演变方式,在以下情况下,Cassandra 连接器在一小段时间内将无法提供准确的数据变更信息:
- 如果某张表的 CDC 被禁用,则在 CDC 被禁用之前发生的数据变更将被跳过。
- 如果某列从表中被移除,则在该列被移除之前涉及该列的数据变更将无法被正确反序列化,并会被跳过。
事件
Cassandra 连接器产生的所有数据变更事件都包含键和值,尽管键和值的结构取决于产生该变更事件的表(参见主题名称)。
变更事件的键
对于给定的表,变更事件的键将具有这样的结构:包含事件创建时该表主键中每个列对应的一个字段。考虑一个 inventory 数据库,其中 customers 表定义如下:
CREATE TABLE customers (
id bigint,
registration_date timestamp,
first_name text,
last_name text,
email text,
PRIMARY KEY (id, registration_date)
);customers 表在保持当前定义期间,其每个变更事件都会使用相同的键 schema,用 JSON 表示如下:
{
"type": "record",
"name": "cassandra-cluster-1.inventory.customers.Key",
"namespace": "io.debezium.connector.cassandra",
"fields": [
{
"name": "id",
"type": "long"
},
{
"name": "registration_date",
"type": "long",
"logicalType": "timestamp-millis"
}
]
}对于 id = 1001 和 registration_date = 1562202942545,键负载(key payload)的 JSON 表示形式如下:
{
"id": 1001,
"registration_date": 1562202942545
}虽然 field.exclude.list 配置属性允许你从事件值中移除某些列,但主键中的所有列始终会包含在事件的键中。 |
|---|
变更事件的值
变更事件消息的值稍微复杂一些。Cassandra 连接器生成的每个变更事件值都具有信封结构,包含以下字段:
op
一个必填字段,包含描述操作类型的字符串值。Cassandra 连接器的取值为:i 表示插入,u 表示更新,d 表示删除。
after
一个可选字段,如果存在,则包含事件发生之后行的状态。该结构由 cassandra-cluster-1.inventory.customers.Value Kafka Connect 模式描述,该模式代表事件所属的集群、键空间和表。
source
一个必填字段,包含描述事件源元数据的结构。对于 Cassandra,它包含以下字段:
- Debezium 版本。
- 连接器名称。
- Cassandra 集群名称。
- 记录该事件的提交日志文件的名称、该事件在提交日志文件中出现的位置、该事件是否属于快照的一部分、受影响的键空间和表的名称,以及分区更新的微秒级最大时间戳。
ts_ms
(可选)如果存在,则包含连接器处理该事件的时间,该时间基于运行 Cassandra 连接器的 JVM 的系统时钟。
由于 Cassandra 不执行写前读取,Cassandra 提交日志不会记录变更应用前的行值。因此,Cassandra 变更事件记录不包含 before 字段。 |
|---|
以下是我们 customers 表中某个 create 事件的值模式的 JSON 表示:
{
"type": "record",
"name": "cassandra-cluster-1.inventory.customers.Envelope",
"namespace": "io.debezium.connector.cassandra",
"fields": [
{
"name": "op",
"type": "string"
},
{
"name": "ts_ms",
"type": "long",
"logicalType": "timestamp-millis"
},
{
"name": "after",
"type": "record",
"fields": [
{
"name": "id",
"type": [
"null",
{
"name": "id",
"type": "record",
"fields": [
{
"name":"value",
"type": "string"
},
{
"name":"deletion_ts",
"type": ["null", "long"],
"default" : "null"
},
{
"name":"set",
"type": "boolean"
}
]
}
]
},
{
"name": "registration_date",
"type": [
"null",
{
"name": "registration_date",
"type": "record",
"fields": [
{
"name":"value",
"type": "long",
"logical_type": "timestamp-millis"
},
{
"name":"deletion_ts",
"type": ["null", "long"],
"default" : "null"
},
{
"name":"set",
"type": "boolean"
}
]
}
]
},
{
"name": "first_name",
"type": [
"null",
{
"name": "first_name",
"type": "record",
"fields": [
{
"name":"value",
"type": "string"
},
{
"name":"deletion_ts",
"type": ["null", "long"],
"default" : "null"
},
{
"name":"set",
"type": "boolean"
}
]
}
]
},
{
"name": "last_name",
"type": [
"null",
{
"name": "last_name",
"type": "record",
"fields": [
{
"name":"value",
"type": "string"
},
{
"name":"deletion_ts",
"type": ["null", "long"],
"default" : "null"
},
{
"name":"set",
"type": "boolean"
}
]
}
]
},
{
"name": "last_name",
"type": [
"null",
{
"name": "email",
"type": "record",
"fields": [
{
"name":"value",
"type": "string"
},
{
"name":"deletion_ts",
"type": ["null", "long"],
"default" : "null"
},
{
"name":"set",
"type": "boolean"
}
]
}
]
}
]
},
{
"name": "source",
"type": "record",
"fields": [
{
"name": "version",
"type": "string"
},
{
"name": "connector",
"type": "string"
},
{
"name": "cluster",
"type": "string"
},
{
"name": "snapshot",
"type": "boolean"
},
{
"name": "keyspace",
"type": "string"
},
{
"name": "table",
"type": "string"
},
{
"name": "file",
"type": "string"
},
{
"name": "position",
"type": "int"
},
{
"name": "ts_ms",
"type": "long",
"logicalType": "timestamp-micros"
}
]
}
]
}待办:验证在删除类 DDL 中,最大时间戳 != 删除时间戳
对于以下 insert DML:
INSERT INTO customers (
id,
registration_date,
first_name,
last_name,
email)
VALUES (
1001,
now(),
"Anne",
"Kretchmar",
"[email protected]"
);以 JSON 表示时,值负载(value payload)如下所示:
{
"op": "c",
"ts_ms": 1562202942832,
"ts_us": 1562202942832014,
"ts_ns": 1562202942832014962,
"after": {
"id": {
"value": 1001,
"deletion_ts": null,
"set": true
},
"registration_date": {
"value": 1562202942545,
"deletion_ts": null,
"set": true
},
"first_name": {
"value": "Anne",
"deletion_ts": null,
"set": true
},
"last_name": {
"value": "Kretchmar",
"deletion_ts": null,
"set": true
},
"email": {
"value": "[email protected]",
"deletion_ts": null,
"set": true
}
},
"source": {
"version": "3.6.3.Final",
"connector": "cassandra",
"cluster": "cassandra-cluster-1",
"snapshot": false,
"keyspace": "inventory",
"table": "customers",
"file": "commitlog-6-123456.log",
"pos": 54,
"ts_ms": 1562202942666382,
"ts_us": 1562202942666382000,
"ts_ns": 1562202942666382000000
}
}给定以下 update DML:
UPDATE customers
SET email = "[email protected]"
WHERE id = 1001 AND registration_date = 1562202942545以 JSON 表示时,值负载(value payload)如下所示:
{
"op": "u",
"ts_ms": 1562202942912,
"ts_us": 1562202942912014,
"ts_ns": 1562202942912014982,
"after": {
"id": {
"value": 1001,
"deletion_ts": null,
"set": true
},
"registration_date": {
"value": 1562202942545,
"deletion_ts": null,
"set": true
},
"first_name": null,
"last_name": null,
"email": {
"value": "[email protected]",
"deletion_ts": null,
"set": true
}
},
"source": {
"version": "3.6.3.Final",
"connector": "cassandra",
"cluster": "cassandra-cluster-1",
"snapshot": false,
"keyspace": "inventory",
"table": "customers",
"file": "commitlog-6-123456.log",
"pos": 102,
"ts_ms": 1562202942666490,
"ts_us": 1562202942666490000,
"ts_ns": 1562202942666490000000
}
}将其与 insert 事件中的值进行比较,我们可以看到几处不同:
op字段的值现在是u,表示该行因为更新操作而发生了变化。after字段现在包含了该行更新后的状态,这里可以看到 email 的值变成了[email protected]。注意first_name和last_name为 null,这是因为这些字段在本次更新中没有发生变化。然而id和registration_date仍然被包含在内,因为它们是该表的主键。source字段的结构与之前相同,但由于此事件来自提交日志(commit log)中的不同位置,其值有所不同。ts_ms显示了连接器处理此事件时的毫秒时间戳。
最后,给定如下的 delete DML:
DELETE FROM customers
WHERE id = 1001 AND registration_date = 1562202942545;JSON 表示形式中的值负载如下所示:
{
"op": "d",
"ts_ms": 1562202942912,
"ts_us": 1562202942912047,
"ts_ns": 1562202942912047921,
"after": {
"id": {
"value": 1001,
"deletion_ts": 1562202972545,
"set": true
},
"registration_date": {
"value": 1562202942545,
"deletion_ts": 1562202972545,
"set": true
},
"first_name": null,
"last_name": null,
"email": null
},
"source": {
"version": "3.6.3.Final",
"connector": "cassandra",
"cluster": "cassandra-cluster-1",
"snapshot": false,
"keyspace": "inventory",
"table": "customers",
"file": "commitlog-6-123456.log",
"pos": 102,
"ts_ms": 1562202942666490,
"ts_us": 1562202942666490000,
"ts_ns": 1562202942666490000000
}
}将其与 insert 和 update 事件中的值进行比较,我们会发现几处不同:
op字段的值现在为d,表示该行的更改是由删除操作引起的。after字段只包含id和registration_date的值,因为这是按主键进行的删除。source字段的结构与之前相同,但值有所不同,因为该事件在提交日志中处于不同的位置。ts_ms显示连接器处理该事件的时间戳(毫秒)。
TODO:鉴于目前尚不支持 TTL,是否最好移除 delete_ts?通过查看每一列是否为 null 来判断某个字段是否被设置,这样是否也可以?
TODO:讨论 Cassandra 连接器中的墓碑事件
数据类型
如上所述,Cassandra 连接器使用结构与该行所在表相同的事件来表示行的更改。事件针对每个列值包含一个字段,该值在事件中的表示方式取决于该列的 Cassandra 数据类型。本节介绍这种映射关系。
下表描述了连接器如何将每种 Cassandra 数据类型映射到 Kafka Connect 数据类型。
| Cassandra 数据类型 | 字面量类型(Schema 类型) | 语义类型(Schema 名称) |
|---|---|---|
ascii | string | n/a |
bigint | int64 | n/a |
blob | bytes | n/a |
boolean | boolean | n/a |
counter | int64 | n/a |
date | int32 | io.debezium.time.Date |
decimal | float64 | n/a |
double | float64 | n/a |
float | float32 | n/a |
frozen | bytes | n/a |
inet | string | n/a |
int | int32 | n/a |
list | array | n/a |
map | map | n/a |
set | array | n/a |
smallint | int16 | n/a |
text | string | n/a |
time | int64 | n/a |
timestamp | int64 | io.debezium.time.Timestamp |
timeuuid | string | io.debezium.data.Uuid |
tinyint | int8 | n/a |
tuple | map | n/a |
uuid | string | io.debezium.data.Uuid |
varchar | string | n/a |
varint | int64 | n/a |
duration | int64 | io.debezium.time.NanoDuration(时长值以纳秒表示的近似形式) |
TODO:添加逻辑类型
任意精度整数类型
Cassandra 连接器根据 varint.handling.mode 连接器配置属性的设置来处理 varint 值。
varint.handling.mode=long
| Cassandra 类型 | 字面类型 | 语义类型 |
|---|---|---|
varint | INT64 | 不适用 |
表 1. varint.handling.mode=long 时的映射
varint.handling.mode=precise
表 2. decimal.handling.mode=precise 时的映射
Cassandra 类型 字面类型 语义类型
varint
BYTES
org.apache.kafka.connect.data.Decimalscale 架构参数被设置为零。
varint.handling.mode=string
| Cassandra 类型 | 字面类型 | 语义类型 |
|---|---|---|
varint | STRING | 不适用 |
表 3. varint.handling.mode=string 时的映射
十进制类型
Cassandra 连接器根据 decimal.handling.mode 连接器配置属性的设置来处理 decimal 值。
decimal.handling.mode=double
| Cassandra 类型 | 字面类型 | 语义类型 |
|---|---|---|
decimal | FLOAT64 | 不适用 |
表 4. decimal.handling.mode=double 时的映射
decimal.handling.mode=precise
表 5. decimal.handling.mode=precise 时的映射
Cassandra 类型 字面类型 语义类型
decimal
STRUCT
io.debezium.data.VariableScaleDecimal
包含一个具有两个字段的结构:类型为 INT32 的 scale,其中保存所传输值的标度;类型为 BYTES 的 value,其中以非缩放形式保存原始值。
decimal.handling.mode=string
| Cassandra 类型 | 字面类型 | 语义类型 |
|---|---|---|
decimal | STRING | 不适用 |
表 6. decimal.handling.mode=string 时的映射
如果默认的数据类型转换无法满足你的需求,可以为该连接器创建自定义转换器。
出现问题时
配置与启动错误
如果配置无效,或者连接器无法使用指定的连接参数成功连接到 Cassandra,Cassandra 连接器将在启动时失败,在日志中报告错误或异常并停止运行。在这种情况下,错误信息会包含有关该问题的更多细节,并可能给出解决建议。修正配置后即可重新启动连接器。
Cassandra 变为不可用
连接器运行后,如果 Cassandra 节点因任何原因变得不可用,连接器将会失败并停止。在这种情况下,请在服务器恢复可用后重新启动连接器。如果这一情况发生在快照期间,连接器将从表的开头重新引导整个表。
Cassandra 连接器正常停止
如果 Cassandra 连接器正常关闭,在停止进程之前,它会确保将 ChangeEventQueue 中的所有事件刷新到 Kafka。Cassandra 连接器每次将流式记录发送到 Kafka 时,都会记录文件名和偏移量。因此,当连接器重新启动时,它会从中断的位置恢复。具体做法是:搜索目录中最旧的提交日志(commit log),开始处理该提交日志,跳过已读取的记录,直到找到最新一条尚未处理的记录。如果 Cassandra 连接器在快照过程中停止,它将从该表恢复,但会重新引导整个表。
Cassandra 连接器崩溃
如果 Cassandra 连接器意外崩溃,它很可能在未记录最新处理偏移量的情况下终止。在这种情况下,当连接器重新启动时,它将从最近记录的偏移量恢复。这意味着很可能会产生重复数据(这并不重要,因为我们已经会从复制因子(RF)中得到重复数据)。请注意,偏移量仅在记录成功发送到 Kafka 后才更新,因此在崩溃期间丢失 ChangeEventQueue 中尚未发布的数据是可以接受的,因为这些事件会被重新生成。
Kafka 变得不可用
当连接器生成变更事件时,它会使用 Kafka 生产者 API 将这些事件发布到 Kafka。如果 Kafka Broker 变得不可用(生产者遇到 TimeoutException),Cassandra 连接器将每秒尝试重新连接一次 Broker,直到重试成功为止。
Cassandra 连接器停止一段时间
根据表的写入负载,当 Cassandra 连接器长时间停止时,可能会触及 cdc_total_space_in_mb 的容量上限。一旦达到这个上限,Cassandra 将停止接受该表的写入操作;这意味着在运行 Cassandra 连接器期间,监控此空间非常重要。如果最坏的情况发生,请完成以下步骤:
- 关闭 Cassandra 连接器。
- 禁用该表的 CDC 功能,使其停止生成额外的写入。由于提交日志未被过滤,同一节点上其他启用 CDC 的表的写入仍可能影响提交日志文件的生成。
- 从偏移量文件中删除已记录的偏移量
- 在容量增加或目录使用空间得到控制后,重新启动连接器,使其重新引导该表。
Cassandra 表的 CDC 启用、临时禁用后再次启用
如果 Cassandra 表暂时禁用了 CDC,然后在一段时间后重新启用它,则必须重新引导。要重新引导单个表,可以手动从 snapshot_offset.properties 文件中删除与该表对应的已记录偏移量行。
部署连接器
Cassandra 连接器应部署到 Cassandra 集群中的每个节点。Cassandra 连接器 Jar 文件需要接收一个 CDC 配置(.properties)文件。参考示例配置。
示例配置
以下是一个本地运行和测试 Cassandra 连接器的 .properties 配置文件示例:
connector.name=test_connector
commit.log.relocation.dir=/Users/test_user/debezium-connector-cassandra/test_dir/relocation/
http.port=8000
cassandra.config=/usr/local/etc/cassandra/cassandra.yaml
cassandra.hosts=127.0.0.1
cassandra.port=9042
kafka.producer.bootstrap.servers=127.0.0.1:9092
kafka.producer.retries=3
kafka.producer.retry.backoff.ms=1000
topic.prefix=test_prefix
key.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url: http://localhost:8081
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url: http://localhost:8081
offset.backing.store.dir=/Users/test_user/debezium-connector-cassandra/test_dir/
snapshot.consistency=ONE
snapshot.mode=ALWAYS
latest.commit.log.only=true连接配置
Cassandra 连接器使用 Cassandra 驱动来配置与 Cassandra 的连接,这些配置必须通过单独的 application.conf 文件提供。你可以在此处找到该驱动配置的完整参考文档,下面是一个示例:
datastax-java-driver {
basic {
request.timeout = 20 seconds
contact-points = [ "spark-master-1:9042" ]
load-balancing-policy {
local-datacenter = "dc1"
}
}
advanced {
auth-provider {
class = PlainTextAuthProvider
username = user
password = pass
}
ssl-engine-factory {
...
}
}
}为了让 Debezium 连接器读取/使用此应用程序配置文件,必须在连接器属性文件中按如下方式设置:
cassandra.driver.config.file=/path/to/application/configuration.conf监控
Cassandra 连接器内置了对 JMX 指标的支持。Cassandra 驱动程序还会发布大量关于驱动活动的指标,这些指标也可以通过 JMX 进行监控。连接器有两类指标。
快照指标
MBean 为 debezium.cassandra:type=connector-metrics,context=snapshot,server=<topic.prefix>。
除非快照操作正在进行中,或自上次启动连接器以来已执行过快照,否则不会暴露快照指标。
下表列出了可用的快照指标。
| 属性名称 | 类型 | 描述 |
|---|---|---|
TotalTableCount | int | 快照中包含的表的总数。 |
RemainingTableCount | int | 快照尚未复制的表的数量。 |
SnapshotRunning | boolean | 快照是否已启动。 |
SnapshotAborted | boolean | 快照是否已中止。 |
SnapshotCompleted | boolean | 快照是否已完成。 |
SnapshotDurationInSeconds | long | 快照到目前为止所花费的总秒数,即使尚未完成。 |
RowsScanned | Map<String, Long> | 一个映射,包含快照中每张表已扫描的行数。在处理过程中,表会被逐步添加到该映射中。每扫描 10,000 行以及完成一张表时都会更新。 |
流处理指标
MBean 为 debezium.cassandra:type=connector-metrics,context=streaming,server=<topic.prefix>。
下表列出了可用的流处理指标。
| 属性名称 | 类型 | 说明 |
|---|---|---|
CommitLogFilename | string | 连接器最近读取的提交日志文件的文件名。 |
CommitLogPosition | long | 连接器已读取的提交日志中的最近位置(以字节计)。 |
NumberOfProcessedMutations | long | 已处理的变更(mutation)数量。 |
NumberOfUnrecoverableErrors | long | 处理提交日志期间发生的不可恢复错误数量。 |
连接器属性
属性
默认值
说明
INITIAL
指定 Cassandra 连接器代理启动时执行快照(例如初始同步)的条件。必须为 'INITIAL'、'ALWAYS' 或 'NEVER' 之一。默认快照模式为 'INITIAL'。
ALL
指定快照查询所使用的 {@link ConsistencyLevel}(一致性级别)。
8000
HTTP 服务器用于 ping、健康检查和构建信息的端口。
无默认值
Cassandra 节点所使用的 YAML 配置文件的绝对路径。
application.conf
Cassandra 驱动配置文件的路径。
commit.log.real.time.processing.enabled
false
仅适用于 Cassandra 4,且设置为 true 时,Cassandra 连接器代理将通过监视提交日志索引文件的更新来增量读取提交日志,并以 commit.log.marked.complete.poll.interval.ms 指定的频率实时流式传输数据。如果设置为 false,则 Cassandra 4 连接器会等待提交日志文件被标记为已完成后再进行处理。
commit.log.marked.complete.poll.interval.ms
10000
仅适用于 Cassandra 4,并且在通过 commit.log.real.time.processing.enabled 启用实时流处理时有效。此配置决定轮询提交日志索引文件以更新偏移量值的频率。
无默认值
处理完成后,提交日志从 cdc\_raw 目录迁移到的本地目录。
commit.log.post.processing.enabled
true
决定 CommitLogPostProcessor 是否运行,以便将已处理的提交日志从迁移目录中移出。如果禁用,提交日志将不会被移出迁移目录。
commit.log.relocation.dir.poll.interval.ms
10000
CommitLogPostProcessor 重新获取迁移目录中所有已处理提交日志的等待时间。
io.debezium.connector.cassandra.BlackHoleCommitLogTransfer
CommitLogPostProcessor 用于将已处理的提交日志从迁移目录中移出的类。内置的传输类是 BlackHoleCommitLogTransfer,它只是将所有已处理的提交日志从迁移目录中删除。如有需要,用户应自行实现定制的提交日志传输类。
commit.log.error.reprocessing.enabled
false
决定 CommitLogProcessor 是否重新处理出错的提交日志。
COMMITLOG\_FILE
指定如何保证变更事件的顺序。每个选项表示用于哈希计算的属性。具有相同哈希值的事件将保持其顺序。建议使用 'PARTITION\_VALUES',以使哈希策略与 Kafka 中的消息保持一致。必须为 'COMMITLOG\_FILE' 或 'PARTITION\_VALUES' 之一。
无默认值
列出连接器可以使用的自定义转换器实例的符号名,以逗号分隔。例如,
isbn
必须设置 converters 属性,才能让连接器使用自定义转换器。
对于你为连接器配置的每个转换器,都必须添加一个 .type 属性,该属性指定实现转换器接口的类的完全限定名。.type 属性使用以下格式:
<converterSymbolicName>.type
例如:
isbn.type: io.debezium.test.IsbnConverter如果你想进一步控制已配置的转换器的行为,可以添加一个或多个配置参数,以便向转换器传递值。要将任何附加配置参数与某个转换器关联起来,需要在参数名前加上该转换器的符号名作为前缀。例如:
isbn.schema.name: io.debezium.cassandra.type.Isbn无默认值
用于存储偏移量跟踪文件的目录。
0
提交偏移量之前需要等待的最短时间。默认值 0 表示每次都会刷新偏移量。
100
在必须将偏移量刷新到磁盘之前允许处理的最大记录数。此配置仅在 offset_flush_interval_ms != 0 时有效。
8192
正整数值,用于指定在写入 Kafka 之前,从提交日志读取并放入其中的变更事件所构成的阻塞队列的最大大小。当写入 Kafka 的速度较慢或 Kafka 不可用时,该队列可以对提交日志读取器提供背压。出现在队列中的事件不包含在此连接器定期记录的偏移量中。默认值为 8192,且应始终大于 max.batch.size 属性中指定的最大批处理大小。该队列在反序列化记录被转换为 Kafka Connect 结构并发送到 Kafka 之前所能容纳的容量。
2048
每次从队列中取出的变更事件的最大数量。
0
长整型值,用于指定阻塞队列的最大字节容量。默认情况下,未对阻塞队列设置容量限制。若要指定队列可占用的字节数,请将此属性设置为一个正的长整型值。
如果同时设置了 max.queue.size,则当队列的大小达到任一属性所指定的限制时,向队列的写入操作将被阻塞。例如,如果你设置 max.queue.size=1000、max.queue.size.in.bytes=5000,那么当队列中包含 1000 条记录后,或者队列中记录的总字节数达到 5000 字节后,向队列的写入操作就会被阻塞。
1000
正整数值,用于指定提交日志处理器在每次迭代中等待队列中出现新变更事件的毫秒数。默认值为 1000 毫秒,即 1 秒。
10000
一个正整数值,指定架构处理器在刷新已缓存的 Cassandra 表架构之前应等待的毫秒数。
10000
每次轮询在重试之前等待的最大时长。
10000
一个正整数值,指定快照处理器在重新扫描表以查找新启用 CDC 的表之前应等待的毫秒数。默认为 10000 毫秒,即 10 秒。
false
删除事件之后是否应生成后续的墓碑事件(true 表示生成,false 表示不生成)。需要注意的是,在 Cassandra 中,具有相同键的两个事件可能会更新同一张表的不同列。因此,如果记录尚未被消费者消费,这可能会导致记录在压缩(compaction)过程中丢失。换句话说,如果你启用了 Kafka 压缩,切勿将此项设置为 true。
无默认值
以逗号分隔的字段全限定名列表,这些字段将被排除在变更事件消息的值之外。字段的全限定名格式为 keyspace_name>.<field_name>.<nested_field_name>。
1
变更事件队列及队列处理器的数量。默认为 1。
t
以逗号分隔的操作类型列表,这些操作将在流式传输过程中被跳过。操作类型包括:c 表示插入/创建,u 表示更新,d 表示删除,t 表示清空(truncate),none 表示不跳过任何操作。默认情况下,清空操作会被跳过(不由此连接器发出)。
io.debezium.schema.SchemaTopicNamingStrategy
应使用的 TopicNamingStrategy 类的名称,用于确定数据变更、架构变更、事务、心跳事件等的 topic 名称,默认为 SchemaTopicNamingStrategy。
.
指定 topic 名称的分隔符,默认为 .。
无默认值
用于所有主题的名称前缀。
| 请勿更改此属性的值。如果更改了该名称值,重新启动后,连接器将不再继续向原有主题发送事件,而是将后续事件发送到基于新值命名的主题。连接器还将无法恢复其数据库模式历史主题。 |
|---|
10000
在有界并发哈希映射中存放主题名称所使用的容量大小。此缓存用于确定与给定数据集合对应的主题名称。
__debezium-heartbeat
控制连接器发送心跳消息所使用的主题名称。主题名称遵循以下模式:
topic.heartbeat.prefix.topic.prefix
例如,如果数据库服务器名称或主题前缀是 fulfillment,则默认主题名称为 __debezium-heartbeat.fulfillment。
如果设置了 topic.heartbeat.name,则忽略此属性。
空
为连接器发送心跳消息所使用的主题指定一个明确的完整名称,从而覆盖由 topic.heartbeat.prefix 和 topic.prefix 推导出的基于前缀的命名方式。
设置后,无论连接器的 topic.prefix 是什么,所有心跳消息都会路由到这个确切的主题名称。当运行多个连接器并希望将心跳事件汇总到一个共享主题中时,此选项非常有用,可避免创建大量单分区的心跳主题。
例如,将其设置为 debezium-heartbeat 会将所有心跳消息路由到名为 debezium-heartbeat 的主题。
如果此属性为空或未设置,连接器将回退到默认行为:topic.heartbeat.prefix.topic.prefix。
long
指定变更事件中 varint 列的表示方式。可选设置如下:
long(默认值)使用 Java 的 long 来表示值,这种方式可能精度不足,但便于使用者使用。
precise 使用 java.math.BigDecimal 表示值,在变更事件中通过二进制表示形式和 Kafka Connect 的 org.apache.kafka.connect.data.Decimal 类型进行编码。
string 将值编码为格式化的字符串,便于消费。
double
指定 decimal 列在变更事件中的表示方式。可选设置有:
double(默认值)使用 Java 的 double 表示值,虽然可能无法提供所需的精度,但便于消费者使用。
precise 使用 java.math.BigDecimal 表示值,在变更事件中通过二进制表示形式和 Kafka Connect 的 org.apache.kafka.connect.data.VariableScaleDecimal 类型进行编码。
string 将值编码为格式化的字符串,便于消费。
none
指定如何调整架构名称以兼容连接器使用的消息转换器。可选设置:
none不进行任何调整。avro将不能用于 Avro 类型名称的字符替换为下划线。avro_unicode将下划线或不能用于 Avro 类型名称的字符替换为对应的 Unicode 转义,如 _uxxxx。注意:_ 是 Java 中类似反斜杠的转义序列。
none
指定如何调整字段名称以兼容连接器使用的消息转换器。可选设置:
none不进行任何调整。avro将不能用于 Avro 类型名称的字符替换为下划线。avro_unicode将下划线或不能用于 Avro 类型名称的字符替换为对应的 Unicode 转义,如 _uxxxx。注意:_ 是 Java 中类似反斜杠的转义序列。
更多详情请参阅 Avro 命名。
无默认值
自定义指标标签接受键值对,用于自定义 MBean 对象名称,这些键值对应追加到常规名称末尾。每个键表示 MBean 对象名称的一个标签,对应的值即为该标签的值。例如:k1=v1,k2=v2。
.*secret$|.*password$|.*sasl\.jaas\.config$|.*basic\.auth\.user\.info$|.*registry\.auth\.client-secret
一个可选的正则表达式,用于指定屏蔽敏感配置键的自定义模式。
默认情况下,Debezium 会按照预定义的模式对配置项键进行屏蔽,以隐藏密码、认证令牌等已知的敏感属性的值。若要自定义连接器屏蔽配置项键值的方式,请将此属性设置为一个正则表达式,用于匹配你想要屏蔽的特定配置项键值。
例如,要屏蔽包含 api.key 或 token 字符串的配置项键,请将此属性设置为以下值:
"custom.sanitize.pattern": ".*api\\.key.*\|.*token.*"如果你为连接器设置了此属性,Debezium 会在返回显示内容、日志记录或响应 API 调用时,对其中匹配的配置键值进行遮蔽。原始值会被一串星号替换,例如 ********。
你为 custom.sanitize.pattern 属性指定的模式会覆盖 Debezium 的默认遮蔽模式;指定的值不会扩展或补充默认模式。 |
|---|
-1
指定当操作因可重试错误(如连接错误)而失败时,连接器的响应方式。
请设置以下选项之一:
-1
无限制。无论之前失败了多少次,连接器始终会自动重启并重试该操作。
0
已禁用。连接器会立即失败,且永不重试该操作。需要人工干预才能重新启动连接器。
> 0
连接器会自动重启,直到达到指定的最大重试次数。在下一次失败后,连接器将停止运行,需要人工干预才能重新启动。
true
此属性指定 Debezium 是否为其发出的消息添加前缀为 __debezium.context. 的上下文标头。OpenLineage 集成需要这些标头,它们提供了元数据,使下游处理系统能够跟踪并识别变更事件的来源。
该属性会添加以下标头:
__debezium.context.connectorLogicalName
Debezium 连接器的逻辑名称。
__debezium.context.taskId
连接器任务的唯一标识符。
__debezium.context.connectorName
Debezium 连接器的名称。
如果 Cassandra 代理使用 SSL 连接 Cassandra 节点,则需要一个 SSL 配置文件。以下示例展示了如何编写 SSL 配置文件:
keyStore.location=/var/private/ssl/cassandra.keystore.jks
keyStore.password=cassandra
keyStore.type=JKS
trustStore.location=/var/private/ssl/cassandra.truststore.jks
trustStore.password=cassandra
trustStore.type=JKS
keyManager.algorithm=SunX509
trustManager.algorithm=SunX509
cipherSuites=TLS_ECDHE_RSA_WITH_AES_256_GCM_SHA384| cipherSuites 字段不是必填的,它仅用于添加一个(或多个)当前不存在的加密套件。trustStore.type 和 keyStore.type 的默认值是 JKS。keyManager.algorithm 和 trustManager.algorithm 的默认值是 SunX509。 |
|---|
该连接器还支持在创建 Kafka 生产者时使用的传递配置属性。具体来说,所有以 kafka.producer. 前缀开头的连接器配置属性,在创建用于将事件写入 Kafka 的生产者时都会被使用(使用时会去掉该前缀)。
例如,可以使用以下连接器配置属性来保护与 Kafka 代理的连接:
kafka.producer.security.protocol=SSL
kafka.producer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
kafka.producer.ssl.keystore.password=test1234
kafka.producer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
kafka.producer.ssl.truststore.password=test1234
kafka.producer.ssl.key.password=test1234
kafka.consumer.security.protocol=SSL
kafka.consumer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
kafka.consumer.ssl.keystore.password=test1234
kafka.consumer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
kafka.consumer.ssl.truststore.password=test1234
kafka.consumer.ssl.key.password=test1234关于 Kafka 生产者的所有配置属性,请务必查阅 Kafka 文档。
该连接器支持以下用于键/值序列化的 Kafka Connect 转换器:
io.confluent.connect.avro.AvroConverter
org.apache.kafka.connect.storage.StringConverter
org.apache.kafka.connect.json.JsonConverter
com.blueapron.connect.protobuf.ProtobufConverter评论
登录后参与评论
KnowForge