Flink CDC:MySQL到TiDB实时同步全流程

以下流程可在5分钟内完成部署,使用Flink CDC实现MySQL到TiDB的实时数据同步。


1. 环境准备
  • Flink 1.13+(需支持CDC连接器)
  • MySQL 5.7+(开启Binlog)
  • TiDB 5.0+
  • 依赖包
    <!-- pom.xml -->
    <dependency>
        <groupId>com.ververica</groupId>
        <artifactId>flink-connector-mysql-cdc</artifactId>
        <version>2.3.0</version>
    </dependency>
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-jdbc</artifactId>
        <version>1.15.4</version>
    </dependency>
    


2. MySQL配置

确保MySQL开启Binlog(在my.cnf中添加):

[mysqld]
server_id=1
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL


3. Flink SQL同步代码

通过SQL实现全量+增量同步(无需Java代码):

-- 创建MySQL CDC源表
CREATE TABLE mysql_users (
    id INT PRIMARY KEY,
    name STRING,
    email STRING
) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = 'localhost',
    'port' = '3306',
    'username' = 'root',
    'password' = '123456',
    'database-name' = 'test_db',
    'table-name' = 'users'
);

-- 创建TiDB目标表(兼容MySQL协议)
CREATE TABLE tidb_users (
    id INT PRIMARY KEY,
    name STRING,
    email STRING
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:mysql://tidb:4000/test_db',
    'table-name' = 'users',
    'username' = 'root',
    'password' = ''
);

-- 启动同步任务
INSERT INTO tidb_users SELECT * FROM mysql_users;


4. 关键参数说明
参数作用示例值
connector数据源类型mysql-cdc
binlog_row_imageMySQL全字段日志FULL
scan.startup.mode同步起点(initial全量同步)latest-offset

5. 部署与验证
  1. 提交作业
    ./bin/flink run -d -c org.apache.flink.table.api.SqlRunner flink-sql-job.jar
    

  2. 验证同步
    • 在MySQL插入数据:
      INSERT INTO users VALUES (1, 'Alice', 'alice@example.com');
      

    • 查询TiDB是否同步:
      SELECT * FROM users WHERE id=1;
      


6. 故障处理
  • 数据延迟:调整Flink并行度(parallelism.default=4
  • 主键冲突:确保TiDB表结构与MySQL一致
  • 断点续传:启用Checkpoint(Flink配置中设置)

提示:TiDB通过JDBC写入,兼容MySQL协议,无需额外驱动。全流程依赖Flink SQL引擎,无需编写自定义代码。

更多推荐