---
title: "OceanBase 数据库与 Kafka 集成最佳实践 - OB Cloud 云数据库 master | OceanBase 文档中心"
description: OceanBase 数据库与 Kafka 集成最佳实践 Kafka 是一个高性能的分布式流处理平台。Kafka 通常用于处理实时数据流，已经被广泛应用于构建实时数据管道、流处理应用和实时分析系统等场景。 您可以使用 Kafka 及 Kafka Connect 模块，然后引入第三方的 Connector 实现相同的功能…
---
切换语言

- 中文站 - 简体中文
- International - English
- 日本站 - 日本語

文档反馈![](https://mdn.alipayobjects.com/huamei_22khvb/afts/img/A*qbZXRo_94ZEAAAAAAAAAAAAADiGDAQ/original) OB Cloud 云数据库

# 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 工作流程

1. 配置 Source Connector，本文中为Debezium MySQL Connector，从 OceanBase 数据库读取数据。
 2. 将读取到的数据写入到 Kafka 中的特定 Topic。
 3. 配置 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 数据库：

1. 开启 OB Cloud 的 Binlog 服务。 开启 Binlog 日志服务的路径：实例列表 -> 租户管理 -> Binlog 服务，点击开通即可。更多信息，参考 [开通 Binlog 日志服务](https://www.oceanbase.com/docs/common-oceanbase-cloud-1000000004154488)。
 2. 连接 OceanBase 数据库，在 OceanBase 数据库中创建 Debezium Connector 的专用账号。

   ```sql
   CREATE USER 'debezium_user'@'localhost' IDENTIFIED BY 'debezium_password';

   ```

   如果您使用的是 OceanBase Cloud，登录 OceanBase 云服务控制台创建用户并授权。更多信息，参考 [创建账号（数据库用户）](https://www.oceanbase.com/docs/common-oceanbase-cloud-1000000001018091)。
 3. 为 `debezium_user` 账号赋予权限。

   ```sql
   GRANT SELECT, CREATE, RELOAD, SHOW DATABASES ON *.* TO 'debezium_user' IDENTIFIED BY 'debezium_password';

   ```
 4. 为 OceanBase 数据库开启 Binlog 服务。

   如果您使用的是 OceanBase Cloud，通过以下方式开启 Binlog：

   开启 Binlog 日志服务的路径：实例列表 -> 租户管理 -> Binlog 服务，点击开通即可。更多信息，参考 [开通 Binlog 日志服务](https://www.oceanbase.com/docs/common-oceanbase-cloud-1000000004154488)。

   运行以下命令，确认 Binlog 已开启：

   ```sql
   SHOW MASTER STATUS;

   ```

   预期返回结果：

   ```sql
   +------------------+----------+--------------+------------------+------------------------------------------+
   | 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)

   ```
 5. 查看 `interactive_timeout` 和 `wait_timeout` 的参数值。

   ```sql
   SHOW VARIABLES LIKE 'interactive_timeout';
   SHOW VARIABLES LIKE 'wait_timeout';

   ```

   这两个参数的默认值均为 28800 秒，即 8 小时。当进行初始化快照时，可能时间很长，连接可能超时，您可根据实际情况决定是否调整这两个参数。

## 配置 Debezium MySQL Connector

按照以下步骤配置 Debezium MySQL Connector：

1. 下载 Debezium MySQL Connector plug-in 1.5.4.Final 版本至安装了 Kafka Connect 的机器上。

   更多信息，参考 [Debezium MySQL connector 下载地址](https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/1.5.4.Final/)。
 2. 解压 Debezium MySQL Connector。

## 从 OceanBase 数据库向 Kafka 发送数据

本节将配置 Kafka 和 Debezium，以同步 OceanBase 数据库的数据。

### 启动 Debezium MySQL Connector

按照以下步骤，以分布式方式启动 Debezium MySQL Connector：

1. 配置 Kafka Connect

   进入 Kafka 安装目录下的 `config` 文件夹。本例中，目录为 `/opt/kafka/config`。修改 `connect-distributed.properties` 文件，在文件的最后一行添加 Debezium MySQL connector 插件的路径。

   ```bash
   # 打开配置文件
   vi /opt/kafka/config/connect-distributed.properties
   # 添加以下内容：
   plugin.path=/opt/debezium/plugin

   ```

   #### 注意

   路径中不要包含 `plugin` 下的 `debezium-connector-mysql` 目录。
 2. 检查 Kafka Topics 状态。

   在您启动 Kafka Connect 之前，检查当前的 topic 状态。

   ```bash
   bin/kafka-topics.sh --list --bootstrap-server localhost:9092

   ```

   本文中目前是空的，没有 topic。
 3. 启动 Kafka Connect。

   使用下列命令启动 Kafka Connect：

   ```bash
   bin/connect-distributed.sh -daemon config/connect-distributed.properties

   ```
 4. 配置 OceanBase 数据库的同步信息。

   在任意目录下创建 `register-oceanbase-debezium.json` 文件，并配置 OceanBase 数据库同步的相关信息。

   ```bash
   # 创建 register-oceanbase-debezium.json 文件
   vi register-oceanbase-debezium.json

   ```

   添加以下内容：

   ```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 文档](https://debezium.io/documentation/reference/stable/connectors/mysql.html)。

   #### 说明

   在上述配置中，`database.server.name` 为 `zhang-oceanbase`，`database.whitelist` 为 `test`。Debezium 会为白名单内数据库的每张表自动创建一个 Topic，命名规则为 `{database.server.name}.{数据库名}.{表名}`。因此，若您的 OceanBase `test` 库中存在表 `tb1`，将自动生成 Topic `zhang-oceanbase.test.tb1`。请先确保在 `test` 库中已创建需要同步的表，下文示例以表 `tb1` 为例。

   #### 注意

      - 将配置项修改为您数据库的相关信息。
      - 在进行初始快照阶段，只有在确信没有其他客户端会修改表结构时，将 `snapshot.locking.mode` 设置为 `none` 才是安全的。
 5. 通过 REST API 注册 Debezium MySQL Connector。

   使用以下命令将配置注册到 Debezium MySQL Connector：

   ```bash
   curl -X POST -H 'Content-Type: application/json' --data @register-oceanbase-debezium.json http://localhost:8083/connectors/

   ```

   #### 注意

   您必须在包含 `register-oceanbase-debezium.json` 文件的路径下执行此命令。
 6. 验证 Connector 是否添加成功。

   ```bash
   curl http://localhost:8083/connectors

   ```

   如果成功，将显示 Connector 列表。
 7. 检查 Kafka Topic。

   ```

   此时应显示相关的 Topic 列表。
 8. 检查 Kafka Topic 中的数据消息。

   执行以下命令，确认数据已经发送到 Kafka Topic：

   ```bash
   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` 中。

```sql
INSERT INTO tb1 (id,name,addtime) VALUES (15,'xiaoli','2024-06-18 17:20:00');

```

### 更新和删除 Connector

1. 更新 Connector 配置。

   更新和新增的 JSON 文件格式是不一样的。更新时，没有 `name` 和 `config` 这两个 key，更新时的文件格式如下：

   ```json
   {
       "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`：

   ```bash
   curl -X PUT -H 'Content-Type: application/json' --data @register-oceanbase-update.json http://localhost:8083/connectors/oceanbase-inventory-connector/config

   ```
 2. 删除 Connector。

   ```bash
   curl -X DELETE http://localhost:8083/connectors/oceanbase-inventory-connector

   ```

   删除 Connector 后 Kafka 的 Topic 没有被删除。
 3. 重新创建相同的 Connector。

   删除后重新创建一个相同的 Connector。新增数据以增量的形式，仍然进入到上次已建立的 Topic 中。

## Kafka 中的消息写入 OceanBase 数据库

按照以下步骤，使用 JDBC Sink 从 Kafka 的 Topic 中消费数据：

1. 准备 Confluent Kafka Connect JDBC。

      1. 下载 Confluent Kafka Connect JDBC。更多信息，参考 [JDBC Connector (Source and Sink)](https://www.confluent.io/hub/confluentinc/kafka-connect-jdbc)。部署方式选择 `Self-Hosted`。
      2. 将下载的压缩包解压，解压后的文件夹放到 `connect-distributed.properties` 文件中 `plugin.path` 配置的路径下。
      3. 由于 `lib` 没有 MySQL 的驱动，将 MySQL 的驱动放到 `lib` 目录下。执行以下命令，复制 `/opt/debezium/plugin/debezium-connector-mysql` 下的驱动：

        ```bash
        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/

        ```
 2. 重新启动 Kafka Connect。

      1. 使用 `jps` 命令找到 `ConnectDistributed` 进程，然后执行 `kill pid`，再重启 Kafka Connect。

        ```bash
        bin/connect-distributed.sh -daemon ./config/connect-distributed.properties

        ```
      2. 重启 Kafka Connect 后，查看 JDBC Sink Connector 是否添加成功：

        ```bash
        curl http://localhost:8083/connector-plugins

        ```

        预期返回结果：

        ```bash
        {"class":"io.confluent.connect.jdbc.JdbcSinkConnector","type":"sink","version":"10.7.6"}

        ```
 3. 编写 Kafka Connect JDBC Sink 配置文件。

   ```bash
   vi register-oceanbase-sink.json

   ```

   在 `register-oceanbase-sink.json` 文件中添加以下内容：

   ```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` 为事先建立的目标库。
 4. 向 Kafka Connect 中添加 JDBC Sink Connector。

   ```bash
   curl -X POST -H "Content-Type: application/json" --data @register-oceanbase-sink.json http://localhost:8083/connectors

   ```

   查看 JDBC Sink Connector 是否添加成功：

   ```

   预期返回结果：

   ```bash
   ["connect-oceanbase-sink","oceanbase-inventory-connector"]

   ```
 5. 同步全量数据

      1. 进入 `test_sink` 查看数据：

        ```sql
        SELECT * FROM tb2;

        ```

        预期返回结果：

        ```bash
        +----+-----------+----------------------+
        | 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)

        ```
 6. 同步增量数据

      1. 在源数据库 `test` 中增加一条 ID 等于 25 的一条数据：

        ```sql
        INSERT INTO tb1 (id, name, addtime) VALUES (25,'xiaozhang','2024-06-25 09:20:00');

        ```
      2. 在目的数据库 `test_sink` 中看到 ID 为 25 的增量数据：

        ```sql
        SELECT * FROM tb2 WHERE id = 25;

        ```
 7. 更改表结构

   在源表中增加 `age` 列：

   ```sql
   ALTER TABLE tb1 ADD COLUMN age INT;

   ```

   此时，原表已增加一列，目标表未变化，因为并没有新的数据到目标表中，只有新增数据时，目标表的表结构才会发生变化。
 8. 在源表中再增加一条 ID 为 26 的数据：

   ```sql
   INSERT INTO tb1 (id, name, addtime, age) VALUES (26,'xiaozhang','2024-06-25 10:20:00',20);

   ```

   可以看到目标表的表结构已经发生变化，并且 ID 为 26 的数据已同步到目标表。

### 更新表数据

目标表的 ID 为 26 的一行 `age` 已更改为 21。

```sql
UPDATE tb1 SET age = 21 WHERE id = 26;

```

### 删除表数据

1. 删除源表中 ID 等于 20 的数据：

   ```sql
   DELETE FROM tb1 WHERE id = 20;

   ```
 2. 目标表中对应的数据已被删除：

   ```sql
   SELECT * FROM tb2 WHERE id = 20;

   ```

 上一篇 下一篇 ![有帮助](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*y6ocSqN8cqsAAAAAAAAAAAAAARQnAQ)![无帮助](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*BG9IQJyLHF8AAAAAAAAAAAAAARQnAQ)![反馈](https://gw.alipayobjects.com/mdn/ob_asset/afts/img/A*eTWdQKCRKHwAAAAAAAAAAAAAARQnAQ)[AI](https://www.oceanbase.com/obi) 咨询热线
