5分钟速通Flink CDC:MySQL到TiDB实时同步全流程
·
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_image | MySQL全字段日志 | FULL |
scan.startup.mode | 同步起点(initial全量同步) | latest-offset |
5. 部署与验证
- 提交作业:
./bin/flink run -d -c org.apache.flink.table.api.SqlRunner flink-sql-job.jar - 验证同步:
- 在MySQL插入数据:
INSERT INTO users VALUES (1, 'Alice', 'alice@example.com'); - 查询TiDB是否同步:
SELECT * FROM users WHERE id=1;
- 在MySQL插入数据:
6. 故障处理
- 数据延迟:调整Flink并行度(
parallelism.default=4) - 主键冲突:确保TiDB表结构与MySQL一致
- 断点续传:启用Checkpoint(Flink配置中设置)
提示:TiDB通过JDBC写入,兼容MySQL协议,无需额外驱动。全流程依赖Flink SQL引擎,无需编写自定义代码。
更多推荐
所有评论(0)