欢迎您访问365答案网,请分享给你的朋友!
生活常识 学习资料

新特性,可以替代Canal的数据同步方案—Flink-CDC

时间:2023-04-17

一、CDC简介

CDC是Change Data Capture(变更数据获取)的简称。核心思想是,监测并捕获数据库的变动(包括数据或数据表的插入、更新以及删除等),将这些变更按发生的顺序完整记录下来,写入到消息中间件中以供其他服务进行订阅及消费。

CDC主要分为基于查询和基于Binlog两种方式,我们主要了解一下这两种之间的区别:

Flink社区开发了 flink-cdc-connectors 组件,这是一个可以直接从 MySQL、PostgreSQL等数据库直接读取全量数据和增量变更数据的 source 组件。大数据培训目前也已开源,开源地址:https://github.com/ververica/flink-cdc-connectors

二、Flink DataStream方式应用的案例实操

1、在pom.xml中增加如下依赖

org.apache.flink flink-java 1.12.0 org.apache.flink flink-streaming-java_2.12 1.12.0 org.apache.flink flink-clients_2.12 1.12.0 org.apache.hadoop hadoop-client 3.1.3 mysql mysql-connector-java 5.1.49 com.alibaba.ververica flink-connector-mysql-cdc 1.2.0 com.alibaba fastjson 1.2.75 org.apache.maven.plugins maven-assembly-plugin 3.0.0 jar-with-dependencies make-assembly package single

2、编写代码

import com.alibaba.ververica.cdc.connectors.mysql.MySQLSource;import com.alibaba.ververica.cdc.debezium.DebeziumSourceFunction;import com.alibaba.ververica.cdc.debezium.StringDebeziumDeserializationSchema;import org.apache.flink.api.common.restartstrategy.RestartStrategies;import org.apache.flink.runtime.state.filesystem.FsStateBackend;import org.apache.flink.streaming.api.CheckpointingMode;import org.apache.flink.streaming.api.datastream.DataStreamSource;import org.apache.flink.streaming.api.environment.CheckpointConfig;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import java.util.Properties;public class FlinkCDC { public static void main(String[] args) throws Exception { //1.创建执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); //2.Flink-CDC将读取binlog的位置信息以状态的方式保存在CK, 如果想要做到断点续传,需要从Checkpoint或者Savepoint启动程序 //2.1 开启Checkpoint,每隔5秒钟做一次CK env.enableCheckpointing(5000L); //2.2 指定CK的一致性语义env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); //2.3 设置任务关闭的时候保留最后一次CK数据env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); //2.4 指定从CK自动重启策略 env.setRestartStrategy(RestartStrategies.fixedDelayRestart (3, 2000L)); //2.5 设置状态后端 env.setStateBackend(new FsStateBackend ("hdfs://hadoop102:8020/flinkCDC")); //2.6 设置访问HDFS的用户名 System.setProperty("HADOOP_USER_NAME", "atguigu"); //3.创建Flink-MySQL-CDC的Source //initial (default): Performs an initial snapshot on the monitored database tables upon first startup , and continue to read the latest binlog. //latest-offset: Never to perform snapshot on the monitored database tables upon first startup, just read from the end of the binlog which means only have the changes since the connector was started. //timestamp: Never to perform snapshot on the monitored database tables upon first startup, and directly read binlog from the specified timestamp、 The consumer will traverse the binlog from the beginning and ignore change events whose timestamp is smaller than the specified timestamp. //specific-offset: Never to perform snapshot on the monitored database tables upon first startup, and directly read binlog from the specified offset. DebeziumSourceFunction mysqlSource = MySQLSource.builder() .hostname("hadoop102") .port(3306) .username("root") .password("000000") .databaseList("gmall-flink") .tableList("gmall-flink.z_user_info") //可选配置项,如果不指定该参数,则会读取上一 个配置下的所有表的数据,注意:指定的时候需 要使用"db.table"的方式 .debeziumProperties(properties) .startupOptions(StartupOptions.initial()) .build(); //4.使用CDC Source从MySQL读取数据 DataStreamSource mysqlDS = env.addSource (mysqlSource); //5.打印数据 mysqlDS.print(); //6.执行任务 env.execute(); }}

3、案例测试

1)打包并上传至Linux

2)开启MySQL Binlog并重启MySQL

3)启动Flink集群

[atguigu@hadoop102 flink-standalone]$ bin/start-cluster.sh

4)启动HDFS集群

[atguigu@hadoop102 flink-standalone]$ start-dfs.sh

5)启动程序

[atguigu@hadoop102 flink-standalone]$ bin/flink run -c com.atguigu.FlinkCDCflink-1.0-SNAPSHOT-jar-with-dependencies.jar

6)在MySQL的gmall-flink.z_user_info表中添加、修改或者删除数据

7)给当前的Flink程序创建Savepoint

[atguigu@hadoop102 flink-standalone]$ bin/flink savepoint JobId hdfs://hadoop102:8020/flink/save

8)关闭程序以后从Savepoint重启程序

[atguigu@hadoop102 flink-standalone]$ bin/flink run -s hdfs://hadoop102:8020/flink/save/...-c com.atguigu.FlinkCDC flink-1.0-SNAPSHOT-jar-with-dependencies.jar

三、Flink SQL方式应用的案例实操

1、在pom.xml中增加如下依赖

org.apache.flink flink-table-planner-blink_2.12 1.12.0

2、代码实现

import org.apache.flink.api.common.restartstrategy.RestartStrategies;import org.apache.flink.runtime.state.filesystem.FsStateBackend;import org.apache.flink.streaming.api.CheckpointingMode;import org.apache.flink.streaming.api.environment.CheckpointConfig;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;public class FlinkSQL_CDC { public static void main(String[] args) throws Exception { //1.创建执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); StreamTableEnvironment tableEnv =StreamTableEnvironment.create(env); //2.创建Flink-MySQL-CDC的Source tableEnv.executeSql("CREATE TABLE user_info (" + " idINT," + " name STRING," + " phone_num STRING" + ") WITH (" + " 'connector' = 'mysql-cdc'," + " 'hostname' = 'hadoop102'," + " 'port' = '3306'," + " 'username' = 'root'," + " 'password' = '000000'," + " 'database-name' = 'gmall-flink'," + " 'table-name' = 'z_user_info'" + ")"); tableEnv.executeSql("select * from user_info").print(); env.execute(); }}

4、自定义反序列化器

代码实现

import com.alibaba.fastjson.JSONObject;import com.alibaba.ververica.cdc.connectors.mysql.MySQLSource;import com.alibaba.ververica.cdc.debezium.DebeziumDeserializationSchema;import com.alibaba.ververica.cdc.debezium.DebeziumSourceFunction;import io.debezium.data.Envelope;import org.apache.flink.api.common.restartstrategy.RestartStrategies;import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.runtime.state.filesystem.FsStateBackend;import org.apache.flink.streaming.api.CheckpointingMode;import org.apache.flink.streaming.api.datastream.DataStreamSource;import org.apache.flink.streaming.api.environment.CheckpointConfig;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.util.Collector;import org.apache.kafka.connect.data.Field;import org.apache.kafka.connect.data.Struct;import org.apache.kafka.connect.source.SourceRecord;import java.util.Properties;public class Flink_CDCWithCustomerSchema { public static void main(String[]args) throws Exception { //1.创建执行环境 StreamExecutionEnvironmentenv = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); //2.创建Flink-MySQL-CDC的Source Properties properties= new Properties(); //initial (default):Performs an initial snapshot on the monitored database tables upon firststartup, and continue to read the latest binlog、 //latest-offset: Never to performsnapshot on the monitored database tables upon first startup, just read fromthe end of the binlog which means only have the changes since the connector wasstarted、 //timestamp: Never to performsnapshot on the monitored database tables upon first startup, and directly readbinlog from the specified timestamp、The consumer will traverse the binlog fromthe beginning and ignore change events whose timestamp is smaller than thespecified timestamp、 //specific-offset: Never toperform snapshot on the monitored database tables upon first startup, anddirectly read binlog from the specified offset、 DebeziumSourceFunction mysqlSource =MySQLSource.builder() .hostname("hadoop102") .port(3306) .username("root") .password("000000") .databaseList("gmall-flink") .tableList("gmall-flink.z_user_info") //可选配置项,如果不指定该参数,则会读取上一个配置下的所有表的数据,注意:指定的时候需要使用"db.table"的方式 .debeziumProperties(properties).startupOptions(StartupOptions.initial()) .deserializer(new DebeziumDeserializationSchema(){ //自定义数据解析器 @Override public void deserialize(SourceRecord sourceRecord,Collector collector) throws Exception { //获取主题信息,包含着数据库和表名 mysql_binlog_source.gmall-flink.z_user_info String topic = sourceRecord.topic(); String[] arr =topic.split("\."); String db = arr[1]; String tableName= arr[2]; //获取操作类型 READ DELETE UPDATE CREATE Envelope.Operation operation = Envelope.operationFor(sourceRecord); //获取值信息并转换为Struct类型 Struct value = (Struct) sourceRecord.value(); //获取变化后的数据 Struct after = value.getStruct("after"); //创建JSON对象用于存储数据信息 JSonObject data = new JSonObject(); for (Field field : after.schema().fields()) { Object o =after.get(field); data.put(field.name(), o); } //创建JSON对象用于封装最终返回值数据信息 JSonObject result = new JSonObject(); result.put("operation", operation.toString().toLowerCase()); result.put("data", data); result.put("database", db); result.put("table", tableName); //发送数据至下游 collector.collect(result.toJSonString()); } @Override public TypeInformation getProducedType(){ return TypeInformation.of(String.class); } }) .build(); //3.使用CDC Source从MySQL读取数据 DataStreamSource mysqlDS =env.addSource(mysqlSource); //4.打印数据 mysqlDS.print(); //5.执行任务 env.execute(); }}

Copyright © 2016-2020 www.365daan.com All Rights Reserved. 365答案网 版权所有 备案号:

部分内容来自互联网,版权归原作者所有,如有冒犯请联系我们,我们将在三个工作时内妥善处理。