2.2、初始化 MySQL 源数据表
到这里 MySQL 环境已经配置完毕。接下来开始准备测试表和数据。
createdatabase demo characterset utf8mb4;
use demo;
createtable student(`id`int primary key, `name`varchar(128), `age`int);
这里创建了演示数据库 demo 和一张 student 表。
三、实战
3.1、使用 Java 代码读取 CDC 数据流
到这一步我们开始使用 Flink 程序来获取 CDC 数据流。
3.1.1、使用传统 MySQL 数据源方式
首先需要引入依赖。注意:本文参考 ververica 官网 的 CDC 教程
https://ververica.github.io/flink-cdc-
connectors/master/content/about.html
, 但是 ververica 被阿里收购后, groupId 变成了 com.alibaba.ververica
, 这个地方需要注意下。
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>1.3.0</version>
</dependency>
然后使用 Table API 编写程序。这里我们仅仅将 CDC 数据流配置为数据源,然
后将 CDC 数据流的内容打印出来。
val env = StreamExecutionEnvironment.getExecutionEnvironment
//
使用
MySQLSource
创建数据源
//
同时指定
StringDebeziumDeserializationSchema
,将
CDC
转换为
String
类型输出
评论