首批通过分布式安全可靠测评,为关键业务系统打造
OceanBase 数据库与 Kafka 集成最佳实践
更新时间:2026-06-19 14:23:45
Kafka 是一个高性能的分布式流处理平台。Kafka 通常用于处理实时数据流,已经被广泛应用于构建实时数据管道、流处理应用和实时分析系统等场景。
您可以使用 Kafka 及 Kafka Connect 模块,然后引入第三方的 Connector 实现相同的功能。目前 Debezium 官方没有 OceanBase 数据库专用的 Connector,但是OceanBase 数据库兼容 MySQL,因此本文将以 Debezium MySQL Connector 和 Confluent JDBC Sink Connector 作为 Kafka 和 OceanBase 数据库的中间桥梁,进行 OceanBase 数据库和 Kafka 的数据集成。
背景信息
Debezium 是为 Kafka Connect 开发的系列 Source Connector,基于日志捕获数据库变更。作为分布式服务,Debezium 记录每个数据表的行级变更,并通过 Kafka 的 Topic 将数据变更呈现给应用进行处理。初始化时,可以发送所有既有数据到 Kafka 主题,再捕获新的数据变更。
Confluent JDBC Sink Connector 是一个用于将 Kafka 数据写入关系型数据库的 Kafka 连接器。它允许用户将 Kafka 主题中的消息流式传输到 JDBC 兼容的数据库中,从而实现了数据的实时同步和持久化存储。
Kafka Connect 工作流程
- 配置 Source Connector,本文中为Debezium MySQL Connector,从 OceanBase 数据库读取数据。
- 将读取到的数据写入到 Kafka 中的特定 Topic。
- 配置 Sink Connector,本文中为 Confluent JDBC Sink Connector,从 Kafka 中的特定 Topic 中读取数据,并写入到 OceanBase 数据库。
演示环境介绍
在您开始集成 OceanBase 数据库与 Kafka 之前,请确认:
- 您的 Debezium 版本为 V1.5.4.Final。
- 您的 Kafka 版本为 V2.12-2.5.0。
- 您的 Zookeeper 版本为 V3.6。
- 您的 OceanBase 数据库版本为 V4.2.3,且您已开启 Binlog。
- 您的 Confluent JDBC Sink Connector 版本为 V10.7.6。
说明
本文的演示环境版本仅供参考。您可以使用其他版本,确保兼容性即可。
配置 OceanBase 数据库
按照以下步骤配置 OceanBase 数据库:
开启 OB Cloud 的 Binlog 服务。 开启 Binlog 日志服务的路径:实例列表 -> 租户管理 -> Binlog 服务,点击开通即可。更多信息,参考 开通 Binlog 日志服务。
连接 OceanBase 数据库,在 OceanBase 数据库中创建 Debezium Connector 的专用账号。
CREATE USER 'debezium_user'@'localhost' IDENTIFIED BY 'debezium_password';如果您使用的是 OceanBase Cloud,登录 OceanBase 云服务控制台创建用户并授权。更多信息,参考 创建账号(数据库用户) 。
为
debezium_user账号赋予权限。GRANT SELECT, CREATE, RELOAD, SHOW DATABASES ON *.* TO 'debezium_user' IDENTIFIED BY 'debezium_password';为 OceanBase 数据库开启 Binlog 服务。
如果您使用的是 OceanBase Cloud,通过以下方式开启 Binlog:
开启 Binlog 日志服务的路径:实例列表 -> 租户管理 -> Binlog 服务,点击开通即可。更多信息,参考 开通 Binlog 日志服务。
运行以下命令,确认 Binlog 已开启:
SHOW MASTER STATUS;预期返回结果:
+------------------+----------+--------------+------------------+------------------------------------------+ | File | Position | Binlog_Do_DB | Binlog_Ignore_DB | Executed_Gtid_Set | +------------------+----------+--------------+------------------+------------------------------------------+ | mysql-bin.000001 | 2567 | | | a2750d9c-11da-11ef-81aa-0242ac110008:1-9 | +------------------+----------+--------------+------------------+------------------------------------------+ 1 row in set (0.125 sec)查看
interactive_timeout和wait_timeout的参数值。SHOW VARIABLES LIKE 'interactive_timeout'; SHOW VARIABLES LIKE 'wait_timeout';这两个参数的默认值均为 28800 秒,即 8 小时。当进行初始化快照时,可能时间很长,连接可能超时,您可根据实际情况决定是否调整这两个参数。
配置 Debezium MySQL Connector
按照以下步骤配置 Debezium MySQL Connector:
下载 Debezium MySQL Connector plug-in 1.5.4.Final 版本至安装了 Kafka Connect 的机器上。
更多信息,参考 Debezium MySQL connector 下载地址。
解压 Debezium MySQL Connector。
从 OceanBase 数据库向 Kafka 发送数据
本节将配置 Kafka 和 Debezium,以同步 OceanBase 数据库的数据。
启动 Debezium MySQL Connector
按照以下步骤,以分布式方式启动 Debezium MySQL Connector:
配置 Kafka Connect
进入 Kafka 安装目录下的
config文件夹。本例中,目录为/opt/kafka/config。修改connect-distributed.properties文件,在文件的最后一行添加 Debezium MySQL connector 插件的路径。# 打开配置文件 vi /opt/kafka/config/connect-distributed.properties # 添加以下内容: plugin.path=/opt/debezium/plugin注意
路径中不要包含
plugin下的debezium-connector-mysql目录。检查 Kafka Topics 状态。
在您启动 Kafka Connect 之前,检查当前的 topic 状态。
bin/kafka-topics.sh --list --bootstrap-server localhost:9092本文中目前是空的,没有 topic。
启动 Kafka Connect。
使用下列命令启动 Kafka Connect:
bin/connect-distributed.sh -daemon config/connect-distributed.properties配置 OceanBase 数据库的同步信息。
在任意目录下创建
register-oceanbase-debezium.json文件,并配置 OceanBase 数据库同步的相关信息。# 创建 register-oceanbase-debezium.json 文件 vi register-oceanbase-debezium.json添加以下内容:
{ "name": "oceanbase-inventory-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "xxx.xx.x.x", "database.port": "xxxx", "database.user": "debezium_user", "database.password": "xxxxxxxx", "database.server.id": "1", "database.server.name": "zhang-oceanbase", "database.whitelist": "test", "database.history.kafka.bootstrap.servers": "localhost:9092", "database.history.kafka.topic": "oceanbase-schema-changes-inventory", "snapshot.locking.mode": "none" } }更多配置项信息,参考 Debezium MySQL Connector 文档。
说明
在上述配置中,
database.server.name为zhang-oceanbase,database.whitelist为test。Debezium 会为白名单内数据库的每张表自动创建一个 Topic,命名规则为{database.server.name}.{数据库名}.{表名}。因此,若您的 OceanBasetest库中存在表tb1,将自动生成 Topiczhang-oceanbase.test.tb1。请先确保在test库中已创建需要同步的表,下文示例以表tb1为例。注意
- 将配置项修改为您数据库的相关信息。
- 在进行初始快照阶段,只有在确信没有其他客户端会修改表结构时,将
snapshot.locking.mode设置为none才是安全的。
通过 REST API 注册 Debezium MySQL Connector。
使用以下命令将配置注册到 Debezium MySQL Connector:
curl -X POST -H 'Content-Type: application/json' --data @register-oceanbase-debezium.json http://localhost:8083/connectors/注意
您必须在包含
register-oceanbase-debezium.json文件的路径下执行此命令。验证 Connector 是否添加成功。
curl http://localhost:8083/connectors如果成功,将显示 Connector 列表。
检查 Kafka Topic。
bin/kafka-topics.sh --list --bootstrap-server localhost:9092此时应显示相关的 Topic 列表。
检查 Kafka Topic 中的数据消息。
执行以下命令,确认数据已经发送到 Kafka Topic:
bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic zhang-oceanbase.test.tb1 --from-beginning在本示例中,若
test库包含表tb1,其中应有 14 条数据,ID最大值为 14。以上步骤展示了如何将 OceanBase 的数据同步到 Kafka 中由 Debezium 自动创建的 Topic(如本示例中的
zhang-oceanbase.test.tb1)。
新增数据
增加一条数据,查看新增数据是否传入到 Kafka 的 Topic zhang-oceanbase.test.tb1 中。
INSERT INTO tb1 (id,name,addtime) VALUES (15,'xiaoli','2024-06-18 17:20:00');
更新和删除 Connector
更新 Connector 配置。
更新和新增的 JSON 文件格式是不一样的。更新时,没有
name和config这两个 key,更新时的文件格式如下:{ "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "xxx.xx.x.x", "database.port": "xxxxx", "database.user": "debezium_user@zhang", "database.server.id": "1", "database.server.name": "zhang-oceanbase", "database.whitelist": "test", "database.history.kafka.bootstrap.servers": "localhost:9092", "database.history.kafka.topic": "oceanbase-schema-changes-inventory", "snapshot.locking.mode": "none" }使用 REST API 更新,方法为
PUT:curl -X PUT -H 'Content-Type: application/json' --data @register-oceanbase-update.json http://localhost:8083/connectors/oceanbase-inventory-connector/config删除 Connector。
curl -X DELETE http://localhost:8083/connectors/oceanbase-inventory-connector删除 Connector 后 Kafka 的 Topic 没有被删除。
重新创建相同的 Connector。
删除后重新创建一个相同的 Connector。新增数据以增量的形式,仍然进入到上次已建立的 Topic 中。
Kafka 中的消息写入 OceanBase 数据库
按照以下步骤,使用 JDBC Sink 从 Kafka 的 Topic 中消费数据:
准备 Confluent Kafka Connect JDBC。
下载 Confluent Kafka Connect JDBC。更多信息,参考 JDBC Connector (Source and Sink)。部署方式选择
Self-Hosted。将下载的压缩包解压,解压后的文件夹放到
connect-distributed.properties文件中plugin.path配置的路径下。由于
lib没有 MySQL 的驱动,将 MySQL 的驱动放到lib目录下。执行以下命令,复制/opt/debezium/plugin/debezium-connector-mysql下的驱动:cp /opt/debezium/plugin/debezium-connector-mysql/mysql-connector-java-8.0.21.jar /opt/debezium/plugin/confluentinc-kafka-connect-jdbc-10.7.6/lib/
重新启动 Kafka Connect。
使用
jps命令找到ConnectDistributed进程,然后执行kill pid,再重启 Kafka Connect。bin/connect-distributed.sh -daemon ./config/connect-distributed.properties重启 Kafka Connect 后,查看 JDBC Sink Connector 是否添加成功:
curl http://localhost:8083/connector-plugins预期返回结果:
{"class":"io.confluent.connect.jdbc.JdbcSinkConnector","type":"sink","version":"10.7.6"}
编写 Kafka Connect JDBC Sink 配置文件。
vi register-oceanbase-sink.json在
register-oceanbase-sink.json文件中添加以下内容:{ "name": "connect-oceanbase-sink", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "connection.url": "jdbc:mysql://xxx.xx.x.x:xxxx/test_sink?useUnicode=true&characterEncoding=UTF-8&useSSL=false", "connection.user": "xxxxx", "connection.password": "xxxxxx", "insert.mode": "upsert", "delete.enabled": "true", "pk.mode": "record_key", "auto.create": "true", "auto.evolve": "true", "topics": "zhang-oceanbase.test.tb1", "table.name.format": "tb2", "transforms": "ExtractField", "transforms.ExtractField.type": "org.apache.kafka.connect.transforms.ExtractField$Value", "transforms.ExtractField.field": "after" } }注意
- 示例没有为 Sink 单独建立用户,使用了
root用户,实际使用可根据情况来选择。 - 源库和目标库位于不同租户,如果为同一租户,可能造成 Binlog 记录的问题,
test_sink为事先建立的目标库。
- 示例没有为 Sink 单独建立用户,使用了
向 Kafka Connect 中添加 JDBC Sink Connector。
curl -X POST -H "Content-Type: application/json" --data @register-oceanbase-sink.json http://localhost:8083/connectors查看 JDBC Sink Connector 是否添加成功:
curl http://localhost:8083/connectors预期返回结果:
["connect-oceanbase-sink","oceanbase-inventory-connector"]同步全量数据
进入
test_sink查看数据:SELECT * FROM tb2;预期返回结果:
+----+-----------+----------------------+ | ID | name | addtime | +----+-----------+----------------------+ | 1 | xiaozhao | 2016-12-09T16:04:33Z | | 2 | xiaozhang | 2016-12-09T16:04:33Z | | 3 | xiaozhang | 2016-12-09T16:04:33Z | | 4 | xiaozhang | 2017-12-09T16:04:33Z | | 5 | xiaozhang | 2017-12-09T16:04:33Z | | 6 | xiaozhang | 2017-12-09T16:04:33Z | | 7 | xiaozhang | 2017-12-09T16:04:33Z | | 8 | xiaozhang | 2017-12-09T16:04:33Z | | 9 | xiaozhang | 2017-12-09T16:04:33Z | | 10 | xiaozhang | 2017-12-09T16:04:33Z | | 11 | xiaozhang | 2017-12-09T16:04:33Z | | 12 | xiaozhang | 2017-12-09T16:04:33Z | | 13 | xiaozhang | 2017-12-09T16:04:33Z | | 14 | xiaozhang | 2017-12-09T16:04:33Z | | 15 | xiaoli | 2024-06-18T09:20:00Z | | 16 | xiaoli | 2024-06-18T09:20:00Z | | 17 | xiaoli | 2024-06-18T09:20:00Z | | 18 | xiaoli | 2024-06-18T09:20:00Z | | 19 | xiaoli | 2024-06-18T09:20:00Z | | 20 | xiaoli | 2024-06-18T09:20:00Z | | 21 | xiaoli | 2024-06-18T09:20:00Z | | 22 | xiaoli | 2024-06-18T09:20:00Z | | 23 | xiaoli | 2024-06-18T09:20:00Z | | 24 | xiaozhang | 2024-06-24T09:20:00Z | +----+-----------+----------------------+ 24 rows in set (0.013 sec)
同步增量数据
在源数据库
test中增加一条 ID 等于 25 的一条数据:INSERT INTO tb1 (id, name, addtime) VALUES (25,'xiaozhang','2024-06-25 09:20:00');在目的数据库
test_sink中看到 ID 为 25 的增量数据:SELECT * FROM tb2 WHERE id = 25;
更改表结构
在源表中增加
age列:ALTER TABLE tb1 ADD COLUMN age INT;此时,原表已增加一列,目标表未变化,因为并没有新的数据到目标表中,只有新增数据时,目标表的表结构才会发生变化。
在源表中再增加一条 ID 为 26 的数据:
INSERT INTO tb1 (id, name, addtime, age) VALUES (26,'xiaozhang','2024-06-25 10:20:00',20);可以看到目标表的表结构已经发生变化,并且 ID 为 26 的数据已同步到目标表。
更新表数据
目标表的 ID 为 26 的一行 age 已更改为 21。
UPDATE tb1 SET age = 21 WHERE id = 26;
删除表数据
删除源表中 ID 等于 20 的数据:
DELETE FROM tb1 WHERE id = 20;目标表中对应的数据已被删除:
SELECT * FROM tb2 WHERE id = 20;