---
title: "同步 OceanBase 数据库的数据至 Kafka | OceanBase 文档中心"
description: 同步 OceanBase 数据库的数据至 Kafka Kafka 是目前广泛应用的高性能分布式流计算平台，OceanBase 迁移服务（OceanBase Migration Service，OMS）支持 OceanBase 两种租户与自建 Kafka 数据源之间的数据实时同步，扩展消息处理能力，广泛应用于实时数据仓…
---
切换语言

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

文档反馈![](https://mdn.alipayobjects.com/huamei_22khvb/afts/img/A*L03BS6f-o40AAAAAAAAAAAAADiGDAQ/original) 迁移服务 OMSV 4.2.3 企业版

# 同步 OceanBase 数据库的数据至 Kafka

更新时间：2026-04-14 15:35:54

Kafka 是目前广泛应用的高性能分布式流计算平台，OceanBase 迁移服务（OceanBase Migration Service，OMS）支持 OceanBase 两种租户与自建 Kafka 数据源之间的数据实时同步，扩展消息处理能力，广泛应用于实时数据仓库搭建、数据查询和报表分流等业务场景。本文为您介绍如何同步 OceanBase 数据库的数据至 Kafka。

OMS 支持数据同步至消息队列产品，扩展业务在监控数据聚合、流式数据处理、在线和离线分析等大数据领域的全方位应用。同步 OceanBase 数据库的数据至 Kafka 时，两种租户对应的数据格式说明请参见 [数据格式说明](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988075)。

## 前提条件

已为源端 OceanBase 数据库创建专用于数据同步项目的数据库用户，并为其赋予了相关权限。详情请参见 [创建数据库用户](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000987931)。

## 使用限制

- 数据同步的对象仅支持物理表，不支持其它对象。
 - OMS 支持的 Kafka 版本为 V0.9、V1.0 和 V2.x。

  #### 注意

  当 Kafka 版本为 V0.9 时，不支持结构同步。
 - 数据同步过程中，如果您在源端修改了同步范围内的表名称，且重命名后的名称不在同步对象中，则该部分数据将不被同步至目标 Kafka 实例中。
 - 待同步的表名和其中的列名不能包含中文字符。
 - 数据源标识和用户账号等，在 OMS 系统内是全局唯一的。
 - OMS 仅支持同步库名、表名和列名为 ASCII 码且不包含特殊字符（包括换行、空格，以及 .|"'`()=;/&\）的对象。
 - OMS 不支持 OceanBase 备库作为源端。

## 注意事项

- 在源端为 OceanBase 数据库并开启同步 DDL 的数据同步项目中，如果源端库表发生重命名（`RENAME`）操作，建议您重新启动项目，避免增量同步丢失数据。
 - 当 OceanBase 数据库为 V4.0.0 ~ V4.3.x 之间的版本（V4.2.5 BP1 除外），并且选择了增量同步时，请为生成列配置 [STORED 属性](https://www.oceanbase.com/docs/common-oceanbase-database-cn-1000000000752593)。否则增量日志中将不保存生成列的信息，可能导致增量同步数据异常的问题。
 - 当更新的行包括 LOB 列时：

     - 如果 LOB 列为更新列，请勿依赖 LOB 列在 `UPDATE` 或 `DELETE` 操作前的值。

      目前使用 LOB 列进行存储的数据类型包括 JSON、GIS、XML、UDT（用户定义类型），以及 LONGTEXT、MEDIUMTEXT 等各类 TEXT。
     - 如果 LOB 列为非更新列，LOB 列在 `UPDATE` 或 `DELETE` 操作前或操作后的值均为 NULL。
 - 节点之间的时钟不同步，或者电脑终端和服务器之间的时钟不同步，均可能导致增量同步的延迟时间不准确。

  例如，如果时钟早于标准时间，可能导致延迟时间为负数。如果时钟晚于标准时间，可能导致延迟。
 - 当项目意外中断进行断点续传时，Kafka 实例中可能会存在部分重复数据（最近一分钟内），因此下游系统需要具备排重能力。
 - 如果创建数据同步项目时，您仅配置了 **增量同步**，OMS 要求源端数据库的本地增量日志保存 48 小时以上。

  ​如果创建数据同步项目时，您配置了 **全量同步** + **增量同步**，OMS 要求源端数据库的本地增量日志至少保留 7 天以上。否则数据同步可能因为无法获取增量日志导致数据同步项目失败，甚至导致源端和目标端数据不一致。
 - 同步 OceanBase 数据库的数据至 Kafka 时，源端执行创建唯一索引语句失败，Kafka 会消费到创建 DDL 语句和删除 DDL 语句。如果传到下游的创建索引 DDL 执行失败，请忽略该异常。

  #### 注意

     - 同步 OceanBase 数据库 V2.2.73 及之后版本、V3.x 之前版本的数据时，请使用 OBCDC V2.2.73 之后、V3.x 之前版本或 OBCDC V3.2.4.5 之后版本，以确保事务内 DML 行变更的顺序。
     - 同步 OceanBase 数据库 V2.2.73 之前版本的数据时，无法保证事务内 DML 行变更的顺序。

## 同步 DDL 支持的范围

- 创建表 `CREATE TABLE`

  #### 注意

  创建的表需要在同步对象范围之内。目前仅支持对已经同步的表进行 `DROP TABLE` 操作后，再执行 `CREATE TABLE`。
 - 修改表 `ALTER TABLE`
 - 删除表 `DROP TABLE`
 - 清空表 `TRUNCATE TABLE`

  延迟删除场景下，同一个事务中会有两条一样的 `TRUNCATE TABLE` DDL。此时，下游消费需要按幂等方式处理。
 - 从指定的分区中删除数据 `ALTER TABLE…TRUNCATE PARTITION`
 - 创建索引 `CREATE INDEX`
 - 删除索引 `DROP INDEX`
 - 添加表的注释 `COMMENT ON TABLE`

  #### 注意

  同步 OceanBase 数据库 Oracle 兼容模式的数据至 Kafka 时，不支持该 DDL。
 - 表重命名 `RENAME TABLE`

  #### 注意

  重命名后的表名需要在同步对象范围之内。

## 操作步骤

1. 新建数据同步项目。

      1. 登录 OMS 控制台。
      2. 在左侧导航栏，单击 **数据同步**。
      3. 在 **数据同步** 页面，单击右上角的 **新建同步项目**。
 2. 在 **选择源和目标** 页面，配置各项参数。

   | 参数 | 描述 |
   | --- | --- |
   | 同步项目名称 | 建议使用中文、数字和字母的组合。名称中不能包含空格，长度不能超过 64 个字符。 |
   | 标签（可选） | 单击文本框，在下拉列表中选择目标标签。您也可以单击 **管理标签**，进行新建、修改和删除。详情请参见 [通过标签管理数据同步项目](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988018)。 |
   | 源端 | 如果您已创建 OceanBase 数据源（包括物理数据源和公有云数据源），请从下拉列表中进行选择。如果未创建，请单击下拉列表中的 **新建数据源**，在右侧对话框进行新建。参数详情请参见 [新建 OceanBase 物理数据源](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988151) 或 [新建 OceanBase 公有云数据源](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988148)。 |
   | 目标端 | 如果您已新建 Kafka 数据源，请从下拉列表中进行选择。如果未新建，请单击下拉列表中的 **新建数据源**，在右侧对话框进行添加。参数详情请参见 [新建 Kafka 数据源](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988032)。 |
 3. 单击 **下一步**，在 **选择同步类型** 页面，选择当前数据同步项目的同步类型。

   同步类型包括 **结构同步**、**全量同步** 和 **增量同步**。**全量同步** 支持无主键表的同步，**增量同步** 支持 **同步 DML**（包括 `Insert`、`Delete` 和 `Update`）和 **同步 DDL**。详情请参见 [DML 过滤](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988070) 和 [同步 DDL](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988082)。
 4. （可选）单击 **下一步**。

   如果您选择了 **增量同步**，但源端 OceanBase 数据源未配置相应参数，则会弹出 **补充数据源信息** 对话框，提醒您进行配置。参数详情请参见 [新建 OceanBase 物理数据源](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988151) 或 [新建 OceanBase 公有云数据源](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988148)。

   补充完成后，单击 **测试连接**。测试连接成功后，单击 **确定**。
 5. 单击 **下一步**，在 **选择同步对象** 页面，选择同步范围。

   同步 OceanBase 数据库的数据至 Kafka 时，支持多表到多 Topic 的同步。

      1. 在选择区域左侧选中需要同步的对象。
      2. 单击 **>**。
      3. 在 **将对象映射至 Topic** 对话框中，选择需要的映射方式进行配置。

        如果选择同步类型时未选择 **结构同步**，则仅支持选择 **已有 Topic**。如果选择同步类型时已选择 **结构同步**，则仅支持选择一种映射方式进行 Topic 的创建或选择。

        例如，已选择结构同步的情况下，您使用了新建 Topic 和选择已有 Topic 两种映射方式，或通过重命名的方式更改了 Topic 的名称，会因为选项冲突导致预检查报错。

        | 参数 | 描述 |
        | --- | --- |
        | 新建 Topic | 在文本框中输入新建 Topic 的名称。支持 3~64 位字符, 且只能包含英文、数字、短横线（-）和下划线（_）。 |
        | 选择 Topic | OMS 提供查询 Kafka Topic 的能力，您可以单击 **选择 Topic**，在 **已有 Topic** 下拉列表中，搜索并选中需要同步的 Topic。 您也可以输入已有 Topic 后选中显示的 Topic 名称。 |
        | 批量生成 Topic | 批量生成 Topic 的规则为 `Topic_${Database Name}_${Table Name}`。 |

        如果您选择 **新建 Topic** 或 **批量生成 Topic**，结构同步成功后，在 Kafka 侧能够查询到新建的 Topic（分区数量默认为 3 个，分区副本数量默认为 1 个，且不支持修改）。如果不符合业务要求，请自行在目标端创建。
      4. 单击 **确定**。

        #### 说明

        OMS 会自动过滤不支持的表，查询表对象的 SQL 语句请参见 [查询表对象 SQL](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000987972)。

   OMS 支持通过文本导入对象，并支持对目标端对象进行更改 Topic、设置行过滤、移除单个对象或全部对象等操作。目标端对象的结构为 Topic>DataBase>Table。

   | 操作 | 步骤 |
   | --- | --- |
   | 导入对象 | 1. 在选择区域的右侧列表中，单击右上角的 **导入对象**。   2. 在对话框中，单击 **确定**。         **注意：**         导入会覆盖之前的操作选择，请谨慎操作。   3. 在 **导入同步对象** 对话框中，导入需要同步的对象。         您可以通过导入 CSV 文件的方式进行库表重命名、设置行过滤条件等操作。详情请参见 [下载和导入同步对象配置](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988015)。   4. 单击 **检验合法性**。   5. 通过合法性的检验后，单击 **确定**。 |
   | 更改 Topic | OMS 支持对目标对象进行更改 Topic 操作。详情请参见 [更改 Topic](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988072)。 |
   | 设置 | OMS 支持配置行过滤、选择分片列和需要同步的列。    1. 在选择区域的右侧列表中，鼠标悬停至目标对象。   2. 单击显示的 **设置**。   3. 在 **设置** 对话框中，您可以进行以下操作。           - 在 **行过滤条件** 区域的文本框中，输入标准的 SQL 语句中的 `WHERE` 子句，来配置行过滤。详情请参见 [SQL 条件过滤数据](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988078)。          - 在 **分片列** 下拉列表中，选择目标分片列。您可以选择多个字段作为分片列，该参数为选填。               选择分片列时，如果没有特殊情况，默认选择主键即可。如果存在主键负载不均衡的情况，请选择唯一性标识且负载相对均衡的字段作为分片列，避免潜在的性能问题。分片列的主要作用如下：                   - 负载均衡：在目标端可以进行并发写入的情况下，通过分片列区分发送消息需要使用的特定线程。                  - 有序性：由于存在并发写入可能导致的无序问题，OMS 确保在分片列的值相同的情况下，用户接收到的消息是有序的。此处的有序是指变更顺序（DML 对于一列的执行顺序）。          - 在 **选择列** 区域，选择需要同步的列。详情请参见 [列过滤](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988074)。   4. 单击 **确定**。 |
   | 移除/全部移除 | OMS 支持在数据映射时，对暂时选中到目标端的单个或多个对象进行移除操作。    - 移除单个同步对象        在选择区域的右侧列表中，鼠标悬停至目标对象，单击显示的 **移除**，即可移除该同步对象。   - 移除全部同步对象        在选择区域的右侧列表中，单击右上角的 **全部移除**。在对话框中，单击 **确定**，即可移除全部同步对象。 |
 6. 单击 **下一步**，在 **同步选项** 页面，配置各项参数。

      - 全量同步

       在 **选择同步类型** 页面，选中 **全量同步**，才会显示下述参数。

       | 参数 | 描述 |
       | --- | --- |
       | 全量同步资源配置 | 您可以选择 **小**、**中**、**大** 的默认读取并发、写入并发和内存，也可以自定义全量同步的资源配置。通过全量导入组件 Full-Import 的资源配置，可以限制项目全量同步阶段的资源消耗。   #### 注意    自定义配置时，最小值为 1，且仅支持配置为整数。 |
      - 增量同步

       在 **选择同步类型** 页面，选中 **增量同步**，才会显示下述参数。

       | 参数 | 描述 |
       | --- | --- |
       | 增量日志拉取资源配置 | 您可以选择 **小**、**中**、**大** 的默认内存，也可以自定义增量日志拉取的资源配置。通过增量拉取组件 Store 的资源配置，可以限制项目增量同步阶段日志拉取的资源消耗。   #### 注意    自定义配置时，最小值为 1，且仅支持配置为整数。 |
       | 增量数据写入资源配置 | 您可以选择 **小**、**中**、**大** 的默认写入并发和内存，也可以自定义增量数据写入的资源配置。通过增量同步组件 Incr-Sync 的资源配置，可以限制项目增量同步阶段数据写入的资源消耗。   #### 注意    自定义配置时，最小值为 1，且仅支持配置为整数。 |
       | 增量记录保存时间 | OMS 中增量解析文件缓存的时长。配置的保存时间越长，Store 组件需要消耗的磁盘空间越大。 |
       | 增量同步起始位点 | - 如果选择同步类型时已选择 **全量同步**，此处默认为项目启动时间，不支持修改。     - 如果选择同步类型时未选择 **全量同步**，请在此处指定同步某个时间节点之后的数据，默认为当前系统时间。详情请参见 [设置增量同步位点](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988079)。 |
      - 高级选项

       | 参数 | 描述 |
       | --- | --- |
       | 序列化方式 | 控制数据同步至 Kafka 的消息格式，目前支持 **Default**、**Canal**、**Dataworks**（支持 2.0 版本）、**SharePlex**、**DefaultExtendColumnType**、**Debezium**、**DebeziumFlatten**、**DebeziumSmt** 和 **Avro**。详情请参见 [数据格式说明](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988075)。    **注意：**     - 目前仅 OceanBase 数据库 MySQL 租户支持 **Debezium**、**DebeziumFlatten**、**DebeziumSmt** 和 **Avro**。     - 当选择 **DataWorks** 时，同步 DDL 不支持 `COMMENT ON TABLE` 和 `ALTER TABLE…TRUNCATE PARTITION`。 |
       | 分区规则 | 同步 OceanBase 数据库的数据至 Kafka Topic 的规则，目前支持 **Hash**、**Table** 和 **One**。不同场景下的 DDL 语句投递和示例，请参见表格下方的说明。      - **Hash** 表示 OMS 使用一定的 Hash 算法，根据主键值或分片列值 Hash 选择 Kafka Topic 的分区。         **注意：**         Hash 仅支持投递到所有的 partition 中。     - **Table** 表示 OMS 将一张表中的全部数据投递至同一个分区中，以表名作为 Hash 键。     - **One** 表示 JSON 消息会投递至 Topic 下的某个分区，目的是为了保持排序。 |
       | 业务系统标识（可选） | 用于标识数据的业务系统来源，以便您后续进行自定义处理。该业务系统标识的长度限制为 1~20 个字符。 |

       下表为不同场景的 DDL 语句投递说明。

       | 分区规则 | DDL 语句涉及多张表（例如 `RENAME TABLE`） | DDL 语句无法确认相关表（例如 `DROP INDEX`） | DDL 语句涉及单张表 |
       | --- | --- | --- | --- |
       | Hash | DDL 语句投递至相关表所在 Topic 的所有分区。   例如，DDL 语句涉及 A、B 和 C 三张表，如果 A 在 Topic 1、B 在 Topic 2、C 不在本项目中，则该 DDL 语句投递至 Topic 1 和 Topic 2 下的所有分区。 | DDL 语句投递至本项目所有 Topic 的所有分区。   例如，DDL 语句无法被 OMS 识别，如果本项目存在三个 Topic，则该 DDL 语句被投递至这三个 Topic 的所有分区。 | DDL 语句投递至该表所属 Topic 下的所有分区。 |
       | Table | DDL 语句投递至相关表所在 Topic 的对应表名 Hash 值所在的分区。   例如，DDL 语句涉及 A、B 和 C 三张表，如果 A 在 Topic 1、B 在 Topic 2、C 不在本项目中，则该 DDL 语句投递至 Topic 1 和 Topic 2 下相关表 Hash 值所在的分区。 | DDL 语句投递至本项目所有 Topic 的所有分区。   例如，DDL 语句无法被 OMS 识别，如果本项目存在三个 Topic，则该 DDL 语句被投递至这三个 Topic 的所有分区。 | 根据 Table Name 进行 Hash，投递至该表所属 Topic 内的某个分区。 |
       | One | DDL 语句投递至相关表所在 Topic 的固定分区。   例如，DDL 语句涉及 A、B 和 C 三张表，如果 A 在 Topic 1、B 在 Topic 2、C 不在本项目中，则该 DDL 语句投递至 Topic 1 和 Topic 2 下某个固定的分区。 | DDL 语句投递至本项目所有 Topic 的某个固定分区。   例如，DDL 语句无法被 OMS 识别，如果本项目存在三个 Topic，则该 DDL 语句被投递至这三个 Topic 的某个固定分区。 | DDL 语句投递至该表所属 Topic 下的某个固定分区。 |

   如果页面的配置参数无法满足需求，您可以单击页面下方的 **参数配置**，进行更加具体的配置。如果您有已配置的项目模板或组件模板，还可以在此处进行引用。

   ![sync-1-zh](https://obbusiness-private.oss-cn-shanghai.aliyuncs.com/doc/img/oms/oms-enterprise/sync-1-zh.png)
 7. 单击 **预检查**，系统对数据同步项目进行预检查。

   在 **预检查** 环节，OMS 会检测和目标端 Kafka 实例的连接情况。如果预检查报错：

      - 您可以排查并处理问题后，重新执行预检查，直至预检查成功。
      - 您也可以单击失败预检查项操作列中的 **跳过**，会弹出对话框提示您跳过本操作的具体影响，确认可以跳过后，请单击对话框中的 **确定**。
 8. 单击 **启动项目**。如果您暂时无需启动项目，请单击 **保存**，跳转至数据同步项目的详情页面，您可以根据需要手动启动数据同步项目。

   OMS 支持在数据同步项目运行过程中修改同步对象，详情请参见 [查看和修改同步对象](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988020)。数据同步项目启动后，会根据选择的同步类型依次执行，详情请参见 [查看数据同步项目的详情](https://www.oceanbase.com/docs/enterprise-oms-doc-cn-1000000000988017) 中《查看同步详情》模块的内容。

   如果数据同步项目运行报错（通常由于网络不通或进程启动过慢导致），您可以在数据同步项目的列表或详情页面，单击 **恢复**。

 上一篇 下一篇 ![有帮助](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) 咨询热线
