OceanBase
  • 产品
  • 解决方案
  • 客户
  • 合作伙伴
  • 资源与服务
  • 文档
  • 社区
云控制台登录 / 注册
  • 免费试用
OceanBase AI 数据平台

基于湖库一体架构,统一管理结构化、半结构化与非结构化等多模态数据,一个系统承载事务处理、实时分析与 AI 工作负载。

一体化能力

TP 事务处理

关键业务稳定运行,保障数据零丢失

AP 实时分析

事务分析一体,驱动智能决策与运营

AI 现代负载

统一多模态数据,支撑生产级应用

关键产品
OceanBase 分布式数据库

首批通过安全可靠测评,面向关键业务

OceanBase 集中式数据库

集中式架构,兼具性能与成本优势

OceanBase AI 数据库

面向 AI 应用与 Agent 的多模数据库

OceanBase AI 湖库

湖库一体的 AI 多模态数据系统

OceanBase DataPilot

企业级 AI 数据分析 Agent

OceanBase Agentbase

企业级 Agent 后端基础设施

OceanBase OAgent

数据库 AI 运维 Agent

OB Cloud

一体化云数据库,提供多云一致体验

OceanBase 数据库一体机

软硬一体,极致性能与高可靠性保障

OceanBase AI 一体机

开箱即用的 AI 数据库一体机

OMA 迁移评估工具

全链路数据库迁移评估

OMS 数据迁移工具

一站式数据传输与同步

OCP 运维管理工具

数据库全生命周期管理

ODC 开发者工具

数据库开发与管控协同

OAS 自治服务工具

数据库智能诊断与自治

通用场景
全场景业务系统 OLTP
实时分析混合负载
异地多活
多基础设施部署
一站式传统数据库升级
混合云部署
一库多芯软硬件混合部署
分布式数据库单机部署
大存储类数据库降本
冷数据归档降本
多实例资源整合
分库分表一体化升级
高并发场景
数据中台
行业解决方案
国有大行和股份制银行核心系统解决方案
区域性银行核心系统解决方案
寿险核心系统解决方案
产险核心系统解决方案
资管交易类系统解决方案
资管 TA 清算类系统解决方案
运营商核心系统解决方案
人社核心系统解决方案
电力核心系统解决方案
行业专区
银行专区

助力银行完成各类核心业务系统升级

保险专区

寿险、产险核心系统升级的更佳选择

零售专区

助力200+零售行业客户规模化落地

DB 大咖说
oceanbase爱奇艺

百亿级卡券业务的“单库双擎”架构升级

oceanbase四川银行

800个测试用例选定分布式数据库

oceanbase太平洋保险

先难后易,核心系统数据库升级复盘

行业案例
oceanbase交通银行

核心数据库的“分布式革命”

oceanbase中国移动

B域核心CRM&BOSS近乎零改造分布式升级

oceanbase理想

打造领先的智能制造系统和自动驾驶体验

演讲实录
oceanbase中国联通

集团应用分布式数据库覆盖B/O/M域

oceanbase国泰海通

智能推送系统稳定支撑单日亿级消息处理量

oceanbase中国联合航空

中国首款机票盲盒背后的数据库力量

用户实践
oceanbase北京银行

最快速度完成40余套系统国产数据库升级

oceanbaseVIVO

替换 MySQL 分库分表,探索成本效益最优

oceanbase滴滴

数据库大规模运维体系建设及落地实践

合作伙伴
合作伙伴类型
联合解决方案
产业生态伙伴
经销商伙伴
技术服务伙伴
培训认证伙伴
生态联合解决方案
神州信息 x OceanBase 银行核心系统
长亮科技 x OceanBase 新核心系统
中电金信 x OceanBase 金融分布式核心系统
天阳科技 x OceanBase 贷记卡方案
易诚互动 x OceanBase 手机银行方案
恒生 x OceanBase UF3.0/O45/TA/估值方案
商业发行版
云树®数据库软件 ActionDB
服务
支持与服务
提交工单
软件下载
OceanBase 企业版
OceanBase 社区版
OB Cloud
学习
培训与认证
在线课堂
在线体验
开发者
开发者中心
资料
行业报告与白皮书
官方博客
年度发布会资料
开发者大会资料
oceanbase白皮书

金融核心系统数据库升级路径与场景实践

oceanbase白皮书

人社关键业务数据库一体化升级实践

产品文档
oceanbaseOceanBase 数据库
数据库一体机
oceanbaseOceanBase AI 数据库
工具与组件
oceanbaseOB Cloud 云数据库
驱动和中间件
快速上手
OceanBase 数据库
OB Cloud 云数据库
知识库
汇聚常见产品使用问题案例
在线体验

Demo与实验,感受 OceanBase 的核心能力与应用场景

OceanBase 最佳实践
了解 OceanBase 分布式数据库的架构与系统原理
技术博客

技术解析 | 用户实践 | 社区月报

在线课堂

电子书 |视频课程|在线培训

Developer Hub
应用开发Demo | 数据开发与集成工具
问答论坛

快速答疑 | 常见问题 | 技术交流

社区活动

Meetup | 技术公开课

GitHub

查看源码 | 贡献代码 | 建议反馈

加入社区

社区组织 | 社区用户贡献 |开发者贡献

进入社区首页
oceanbase数据库大赛

第六届OceanBase数据库大赛

oceanbase免费课程

《Easy Data x AI》:面向所有 AI 爱好者的 Data 与 AI 基础知识入门教程

切换语言
  • 中文站 - 简体中文
  • International - English
  • 日本站 - 日本語

OceanBase

OceanBase 海扬数据库始创于 2010 年,是完全自主研发的数据库公司。2020年开始独立商业化运作,历经15年大规模核心场景验证,目前是中国数据库的领军企业之一。从分布式数据库到 AI 数据库,为企业提供安全、稳定、可扩展的数据底座,推动数据基础设施全面拥抱 AI 时代。

关于我们

关于 OceanBase最新动态资质荣誉客户专家委员会招贤纳士合作伙伴年度发布会开发者大会

资源与服务

支持与服务文档知识库软件与工具下载培训与认证在线体验数据库专题视频

社区

快速上手开发者中心博客活动学习问答GitHub

数据库百科

分布式数据库国产数据库OLTP 数据库OLAP 数据库HTAP 数据库数据库向量数据库向量检索

联系我们

服务热线:
400-109-0633
商务咨询
培训认证技术支持媒体合作
京公网安备11010802047223号京公网安备11010802047223号
京ICP备20024574号-1
合字B1.B2-20250395
网站服务协议隐私协议安全响应协议
OceanBase 版权所有 © 2026 基础资源和备案服务由阿里云提供
文档反馈
  1. 文档中心
  2. OceanBase 数据库
  3. 分布式版
  4. V5.0.1
  5. 数据迁移
  6. 从 MySQL 数据库迁移数据到 OceanBase 数据库
  7. 使用 Flink CDC 从 MySQL 数据库同步数据到 OceanBase 数据库
分布式版-V5.0.1
  • 简介
  • 快速上手
  • 应用开发
  • 部署
  • 升级
  • 数据迁移
    • 数据迁移概述
    • 从 MySQL 数据库迁移数据到 OceanBase 数据库
      • 使用 OMS 从 MySQL 数据库迁移数据到 OceanBase 数据库 MySQL 租户
      • 使用 mydumper 和 myloader 从 MySQL 数据库迁移数据到 OceanBase 数据库
      • 使用 DBCAT 迁移 MySQL 表结构到 OceanBase 数据库
      • 使用 DataX 迁移 MySQL 表数据到 OceanBase 数据库
      • 使用 CloudCanal 从 MySQL 数据库迁移数据到 OceanBase 数据库
      • 使用 Canal 从 MySQL 数据库同步数据到 OceanBase 数据库
      • 使用 Flink CDC 从 MySQL 数据库同步数据到 OceanBase 数据库
      • 使用 ChunJun 从 MySQL 数据库迁移数据到 OceanBase 数据库
    • 从 OceanBase 数据库迁移数据到 MySQL 数据库
    • 从 Oracle 数据库迁移数据到 OceanBase 数据库
    • 从 OceanBase 数据库迁移数据到 Oracle 数据库
    • 从 DB2 数据库迁移数据到 OceanBase 数据库
    • 从 OceanBase 数据库迁移数据到 DB2 数据库
    • 从 TiDB 数据库迁移数据到 OceanBase 数据库
    • 从 PostgreSQL 数据库迁移数据到 OceanBase 数据库
    • 从 CSV 文件迁移数据到 OceanBase 数据库
    • 从 SQL 文件导入数据到 OceanBase 数据库
    • 数据库之间迁移数据
    • 使用 SQL 语句迁移数据
    • 旁路导入概述
  • 管理数据库
  • AP
  • AI
  • 生态集成
  • 实践教程
  • 参考指南
  • 常见问题
  • 版本发布记录
  • 术语
  1. 文档中心
  2. OceanBase 数据库
  3. 分布式版
  4. V5.0.1
  5. 数据迁移
  6. 从 MySQL 数据库迁移数据到 OceanBase 数据库
  7. 使用 Flink CDC 从 MySQL 数据库同步数据到 OceanBase 数据库

使用 Flink CDC 从 MySQL 数据库同步数据到 OceanBase 数据库

更新时间:2026-07-29 10:39:50

github-fill编辑
编组分享

Flink CDC (CDC Connectors for Apache Flink) 是 Apache Flink 的一组 Source 连接器,它支持从大多数据库中实时地读取存量历史数据和增量变更数据。Flink CDC 能够将数据库的全量和增量数据同步到消息队列和数据仓库中。Flink CDC 也可以用于实时数据集成,您可以使用它将数据库数据实时导入数据湖或者数据仓库。同时,Flink CDC 还支持数据加工,您可以通过它的 SQL Client 对数据库数据做实时关联、打宽、聚合,并将结果写入到各种存储中。CDC (Change Data Capture,即变更数据捕获)能够帮助您监测并捕获数据库的变动。CDC 提供的数据可以做很多事情,比如:做历史库、做近实时缓存、提供给消息队列(MQ),用户消费 MQ 做分析和审计等。

以下将介绍使用 Flink CDC 从 MySQL 数据库同步数据到 OceanBase 数据库。

Flink CDC 环境准备

下载 Flink 和所需要的依赖包:

  1. 通过 下载地址 下载 Flink。本文档使用的是 Flink 1.15.3,并将其解压至目录 /FLINK_HOME/flink-1.15.3。

  2. 下载下面列出的依赖包,并将它们放到目录 /FLINK_HOME/flink-1.15.3/lib/ 下。

    • flink-sql-connector-mysql-cdc-2.1.1.jar

    • flink-connector-jdbc-1.15.3.jar

    • mysql-connector-java-5.1.47.jar

准备数据

准备 MySQL 数据库数据

在 MySQL 数据库中准备测试数据,作为导入 OceanBase 数据库的源数据。

  1. 进入 MySQL 数据库。

    [xxx@xxx /...]
    $mysql -hxxx.xxx.xxx.xxx -P3306 -uroot -p******
    <Omit echo information>
    
    MySQL [(none)]>
    
  2. 创建数据库 test_mysql_to_ob,表 tbl1 和 tbl2,并插入数据。

    MySQL [(none)]> CREATE DATABASE test_mysql_to_ob;
    Query OK, 1 row affected
    
    MySQL [(none)]> USE test_mysql_to_ob;
    Database changed
    MySQL [test_mysql_to_ob]> CREATE TABLE tbl1(col1 INT PRIMARY KEY, col2 VARCHAR(20),col3 INT);
    Query OK, 0 rows affected
    
    MySQL [test_mysql_to_ob]> INSERT INTO tbl1 VALUES(1,'China',86),(2,'Taiwan',886),(3,'Hong Kong',852),(4,'Macao',853),(5,'North Korea',850);
    Query OK, 5 rows affected
    Records: 5  Duplicates: 0  Warnings: 0
    
    MySQL [test_mysql_to_ob]> CREATE TABLE tbl2(col1 INT PRIMARY KEY,col2 VARCHAR(20));
    Query OK, 0 rows affected
    
    MySQL [test_mysql_to_ob]> INSERT INTO tbl2 VALUES(86,'+86'),(886,'+886'),(852,'+852'),(853,'+853'),(850,'+850');
    Query OK, 5 rows affected
    Records: 5  Duplicates: 0  Warnings: 0
    

准备 OceanBase 数据库数据

在 OceanBase 数据库中创建存放源数据的表。

  1. 登录 OceanBase 数据库。

    使用 user001 用户登录集群的 mysql001 租户。

    [xxx@xxx /...]
    $obclient -h10.10.10.2 -P2881 -uuser001@mysql001 -p -A
    Enter password:
    Welcome to the OceanBase.  Commands end with ; or \g.
    Your OceanBase connection id is 3221536981
    Server version: OceanBase 4.0.0.0 (r100000302022111120-7cef93737c5cd03331b5f29130c6e80ac950d33b) (Built Nov 11 2022 20:38:33)
    
    Copyright (c) 2000, 2018, OceanBase and/or its affiliates. All rights reserved.
    
    Type 'help;' or '\h' for help. Type '\c' to clear the current input statement.
    
    obclient [(none)]>
    
  2. 创建数据库 test_mysql_to_ob 和表 mysql_tbl1_and_tbl2。

    obclient [(none)]> CREATE DATABASE test_mysql_to_ob;
    Query OK, 1 row affected
    
    obclient [(none)]> USE test_mysql_to_ob;
    Database changed
    obclient [test_mysql_to_ob]> CREATE TABLE mysql_tbl1_and_tbl2(col1 INT PRIMARY KEY,col2 INT,col3 VARCHAR(20),col4 VARCHAR(20));
    Query OK, 0 rows affected
    

启动 Flink 集群和 Flink SQL CLI

  1. 使用下面的命令跳转至 Flink 目录下。

    [xxx@xxx /FLINK_HOME]
    #cd flink-1.15.3
    
  2. 使用下面的命令启动 Flink 集群。

    [xxx@xxx /FLINK_HOME/flink-1.15.3]
    #./bin/start-cluster.sh
    

    启动成功的话,可以在 http://localhost:8081/ 访问到 Flink Web UI,如下所示:

    Flink_Web_UI

    说明

    执行 ./bin/start-cluster.sh 后,如果提示:bash: ./bin/start-cluster.sh: Permission denied。需要把 flink-1.15.3 目录下的所有 -rw-rw-r-- 权限的文件的权限都设置为 -rwxrwxrwx 权限。

    示例如下:

    
      [xxx@xxx /FLINK_HOME/flink-1.15.3]
      # chmod -R 777 /FLINK_HOME/flink-1.15.3/*
      
  3. 使用下面的命令启动 Flink SQL CLI。

    [xxx@xxx /FLINK_HOME/flink-1.15.3]
    #./bin/sql-client.sh
    

    启动成功后,可以看到如下的页面:

    Flink_SQL_CLI

设置 checkpoint

在 Flink SQL CLI 中开启 checkpoint,每隔 3 秒做一次 checkpoint。

Flink SQL> SET execution.checkpointing.interval = 3s;
[INFO] Session property has been set.

创建 MySQL CDC 表

在 Flink SQL CLI 中创建 MySQL 数据库对应的表。

对于 MySQL 数据库中 test_mysql_to_ob 的表 tbl1 和 tbl2 使用 Flink SQL CLI 创建对应的表,用于同步这些底层数据库表的数据。

Flink SQL> CREATE TABLE mysql_tbl1 (
    col1 INT PRIMARY KEY,
    col2 VARCHAR(20),
    col3 INT) 
    WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'xxx.xxx.xxx.xxx',
    'port' = '3306',
    'username' = 'root',
    'password' = '******',
    'database-name' = 'test_mysql_to_ob',
    'table-name' = 'tbl1');
[INFO] Execute statement succeed.

Flink SQL> CREATE TABLE mysql_tbl2 (col1 INT PRIMARY KEY,
    col2 VARCHAR(20))
    WITH ('connector' = 'mysql-cdc',
    'hostname' = 'xxx.xxx.xxx.xxx',
    'port' = '3306',
    'username' = 'root',
    'password' = '******',
    'database-name' = 'test_mysql_to_ob',
    'table-name' = 'tbl2');
[INFO] Execute statement succeed.

有关 MySQL CDC Connector WITH 选项的详细信息,请参见 Connector Options。

创建 OceanBase CDC 表

在 Flink SQL CLI 中创建 OceanBase 数据库对应的表。创建 mysql_tbl1_and_tbl2 表,用来将关联后的数据写入 OceanBase 数据库中。

Flink SQL> CREATE TABLE mysql_tbl1_and_tbl2(
    col1 INT PRIMARY KEY,
    col2 INT,col3 VARCHAR(20),
    col4 VARCHAR(20))
    WITH ('connector' = 'jdbc',
    'url' = 'jdbc:mysql://10.10.10.2:2881/test_mysql_to_ob',
    'username' = 'root@mysql001',
    'password' = '******',
    'table-name' = 'mysql_tbl1_and_tbl2');
[INFO] Execute statement succeed.

有关 JDBC SQL Connector WITH 选项的详细信息,请参见 Connector Options。

在 Flink SQL CLI 中将数据写入 OceanBase 数据库中

使用 Flink SQL 将表 tbl1 与表 tbl2 关联,并将关联后的信息写入 OceanBase 数据库中。

Flink SQL> INSERT INTO mysql_tbl1_and_tbl2 
    SELECT t1.col1,t1.col3,t1.col2,t2.col2 
    FROM mysql_tbl1 t1,mysql_tbl2 t2 
    WHERE t1.col3=t2.col1;
[INFO] Submitting SQL update statement to the cluster...
Loading class `com.mysql.jdbc.Driver'. This is deprecated. The new driver class is `com.mysql.cj.jdbc.Driver'. The driver is automatically registered via the SPI and manual loading of the driver class is generally unnecessary.
[INFO] SQL update statement has been successfully submitted to the cluster:
Job ID: c5ee92498addf813858e448ec25e85af

说明

本文档测试示例使用的 MySQL 驱动(com.mysql.jdbc.Driver)是 MySQL Connector/J 5.1.47 版本。新版本 MySQL 驱动(com.mysql.cj.jdbc.Driver)请使用 MySQL Connector/J 8.x 版本。

查看关联数据写入 OceanBase 数据库情况

登录 OceanBase 数据数,在 test_mysql_to_ob 库中查看表 mysql_tbl1_and_tbl2 的数据。

obclient [test_mysql_to_ob]> SELECT * FROM mysql_tbl1_and_tbl2;
+------+------+-------------+------+
| col1 | col2 | col3        | col4 |
+------+------+-------------+------+
|    1 |   86 | China       | +86  |
|    2 |  886 | Taiwan      | +886 |
|    3 |  852 | Hong Kong   | +852 |
|    4 |  853 | Macao       | +853 |
|    5 |  850 | North Korea | +850 |
+------+------+-------------+------+
5 rows in set

查看数据更新情况

  1. 在 MySQL 数据库的表 tbl1 和 tbl2 中插入分别插入一条数据。

    MySQL [test_mysql_to_ob]> INSERT INTO tbl1 VALUES(6,'code',673);
    Query OK, 1 row affected
    
    MySQL [test_mysql_to_ob]> INSERT INTO tbl2 VALUES(673,'+673');
    Query OK, 1 row affected
    
  2. 在 OceanBase 数据库中查看数据是否同步。

    obclient [test_mysql_to_ob]> SELECT * FROM mysql_tbl1_and_tbl2;
    +------+------+-------------+------+
    | col1 | col2 | col3        | col4 |
    +------+------+-------------+------+
    |    1 |   86 | China       | +86  |
    |    2 |  886 | Taiwan      | +886 |
    |    3 |  852 | Hong Kong   | +852 |
    |    4 |  853 | Macao       | +853 |
    |    5 |  850 | North Korea | +850 |
    |    6 |  673 | code        | +673 |
    +------+------+-------------+------+
    6 rows in set
    

本文目录

Flink CDC 环境准备准备数据准备 MySQL 数据库数据准备 OceanBase 数据库数据启动 Flink 集群和 Flink SQL CLI设置 checkpoint创建 MySQL CDC 表创建 OceanBase CDC 表在 Flink SQL CLI 中将数据写入 OceanBase 数据库中查看关联数据写入 OceanBase 数据库情况查看数据更新情况
有帮助
无帮助
反馈
AI