截止目前,Flink CDC已经更新到3.X版本,因为笔者部署的Flink是1.17.1,最新版本的Flink CDC并不兼容,所以经过测试后得出能在Flink1.17.1下兼容的Flink CDC版本。

根据Flink官方文档,要使用mysql-cdc连接器,需要配置的JDBC Driver的版本为8.0.27。
在这里插入图片描述
从maven仓库中下载mysql-connector-java-8.0.27.jar,地址https://mvnrepository.com/artifact/mysql/mysql-connector-java/8.0.27,上传到Flink的lib目录。
下载flink-sql-connector-mysql-cdc-2.4.2.jar,上传到lib目录。
Flink官网中问答部分有介绍flink-sql-connector-mysql-cdc和flink-connector-mysql-cdc的区别以及使用场景,flink-sql-connector-mysql-cdc包括了Flink SQL执行所需的依赖,是个fat jar,只需导入这个包就可以执行Flink SQL连接mysql-cdc。而flink-connector-XX是用在DataStream,用户需要自己管理所需的三方包依赖,有冲突的依赖需要自己做 exclude, shade 处理。
两个jar包上传成功后,重启flink集群:

./bin/stop-cluster.sh
./bin/start-cluster.sh

打开Flink SQL客户端:

./bin/sql-client.sh 

mysql-cdc连接MySQL数据库的用户需要有select、slave权限,如果权限不足会报错:

Access denied; you need (at least one of) the SUPER, REPLICATION
CLIENT privilege(s) for this operation

执行以下SQL:

 CREATE TABLE sink_managecom_map (
   provincecomcode varchar(20),
   provincecomname varchar(255),
   primary key (provincecomcode) not enforced
 ) WITH (
    'connector' = 'mysql-cdc',
    'hostname' = '*****',
    'port' = '3306',
    'username' = '******',
    'password' = '*****',
    'database-name' = 'msbi',
    'table-name' = 'managecom_map',
    'server-id'='5401-5404',
    'debezium.snapshot.mode' = 'initial'
 );

 select * from sink_managecom_map;

执行结果如下:
在这里插入图片描述
验证insert、delete、update操作,能实时更新。
Flink CDC的更多用法可以参考官网https://flink.apache.org/

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐