暂无图片
暂无图片
1
暂无图片
暂无图片
暂无图片
Flink 使用之 MySQL CDC.pdf
270
6页
1次
2023-08-15
5墨值下载
Flink 使用之 MySQL CDC
一、CDC 简介
CDC Change Data Capture
变更数据捕获,为 Flink 1.11 中一个新增功能。我们可以通过 CDC 得知数据源
表的更新内容(包含 Insert Update Delete),并将这些更新内容作为数据
流发送到下游系统。捕获到的数据操作具有一个标识符,分别对应数据的增
加,修改和删除。
+I:新增数据。
-U:一条数据的修改会产生两个 U
标识符数据。其-U
含义为修改前数据。
+U:修改之后的数据。
-D:删除的数据。
在我的文章 Flink CDC 中可以更详细的理解 CDC
二、CDC 需要的配置
2.1MySQL 启用 binlog
接下来以 MySQL CDC 为例,熟悉配置 Flink MySQL CDC
在使用 CDC 之前务必要开启 MySQL binlog。下面以 MySQL 5.7 版本为例
说明。
修改 my.cnf
文件,增加:
server_id=1
log_bin=mysql-bin
binlog_format=ROW
expire_logs_days=30
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
类型输出
of 6
5墨值下载
【版权声明】本文为墨天轮用户原创内容,转载时必须标注文档的来源(墨天轮),文档链接,文档作者等基本信息,否则作者和墨天轮有权追究责任。如果您发现墨天轮中有涉嫌抄袭或者侵权的内容,欢迎发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论

关注
最新上传
暂无内容,敬请期待...
下载排行榜
Top250 周榜 月榜