欢迎光临
我们一直在努力

FlinkCDC_MySQL同步案例

水善利万物而不争,处众人之所恶,故几于道💦

文章目录

      • 一、 Flink集群搭建
        • 1. 编辑conf目录下的config.yaml文件。
        • 2. 编辑conf目录下的masters文件:
        • 3. 编辑conf目录下的workers文件:
        • 4. 把flink安装目录分发到其他机器
        • 5. 启动集群
      • 二、 mysql-cdc准备
        • 1. 驱动jar包
        • 2. 建源表和目标表
        • 3. 编写同步脚本
        • 4. 提交执行
        • 5. Web页面查看执行情况
        • 6. 修改源表中的数据进行测试

一、 Flink集群搭建

下载Flink安装包,然后解压。我这里以1.19.3版本为例。

1. 编辑conf目录下的config.yaml文件。

指定Job Manager的地址信息: 在这里插入图片描述 指定Task Manager的地址信息: 在这里插入图片描述 指定Web页面访问地址: 在这里插入图片描述

2. 编辑conf目录下的masters文件:

配置JobManager的地址 在这里插入图片描述

3. 编辑conf目录下的workers文件:

配置Task Manager的地址 在这里插入图片描述

4. 把flink安装目录分发到其他机器

分发到其他机器,然后将taskmanager下的地址改成对应的机器名: 在这里插入图片描述 在这里插入图片描述

5. 启动集群

执行bin目录下的start-cluster.sh脚本,启动flink集群 在这里插入图片描述 启动后访问jobmanager的8081端口: 在这里插入图片描述 执行bin/sql-client.sh脚本,可以启动sql客户端,然后可以在sql-client客户端中写sql提交任务: 在这里插入图片描述 我这里就不在sql客户端写了,直接写在sql文件中,用sql客户端去执行文件就可以了。 bin/sql-client.sh -f mysql_to_mysql.sql

二、 mysql-cdc准备

1. 驱动jar包

已上传资源。将这两个jar包放到Flink的lib目录下,然后重启Flink集群 在这里插入图片描述

2. 建源表和目标表

源表和目标表必须存在,实测如果没有目标表会提示目标表不存在 源表,test库下的a表:

CREATE TABLE `a` (
`sku_id` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NOT NULL,
`price` decimal(10, 2) NULL DEFAULT NULL,
`category_id` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
`from_date` datetime NULL DEFAULT NULL,
`ddd` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
`str` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
PRIMARY KEY (`sku_id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_0900_ai_ci ROW_FORMAT = Dynamic;

目标表,test库下的a_flink_cdc表:

CREATE TABLE `a_flink_cdc` (
`sku_id` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NOT NULL,
`price` decimal(10, 2) NULL DEFAULT NULL,
`category_id` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
`from_date` datetime NULL DEFAULT NULL,
`ddd` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
`str` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci NULL DEFAULT NULL,
PRIMARY KEY (`sku_id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_0900_ai_ci ROW_FORMAT = Dynamic;

3. 编写同步脚本

在flink安装目录下新建job目录,在job目录下新建mysql_to_mysql.sql文件,内容如下:

CREATE TABLE b_source (
sku_id VARCHAR(255) primary key not enforced,
price DECIMAL(10,2),
category_id VARCHAR(255),
from_date TIMESTAMP(3),
ddd VARCHAR(255),
str VARCHAR(255)
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '192.168.1.11',
'port' = '3306',
'username' = 'root',
'password' = 'xxxxxx',
'database-name' = 'test',
'table-name' = 'a',
'scan.startup.mode' = 'initial'
);

CREATE TABLE a_sink (
sku_id VARCHAR(255) primary key not enforced,
price DECIMAL(10,2),
category_id VARCHAR(255),
from_date TIMESTAMP(3),
ddd VARCHAR(255),
str VARCHAR(255)
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://192.168.1.11:3306/test',
'driver' = 'com.mysql.cj.jdbc.Driver',
'username' = 'root',
'password' = 'xxxxxx',
'table-name' = 'a_flink_cdc'
);
insert into a_sink select * from b_source;

4. 提交执行

[qcln@hadoop102 flink-1.19.3]$ bin/sql-client.sh -f job/mysql_to_mysql.sql

在这里插入图片描述

5. Web页面查看执行情况

在这里插入图片描述

6. 修改源表中的数据进行测试

对源表中的数据进行增删改,刷新目标表可以看到实时同步成功 在这里插入图片描述

赞(0)
未经允许不得转载:171主机测评 » FlinkCDC_MySQL同步案例
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址