FlinkCDC从Mysql数据写入Kafka
环境安装:
1.jdk
2.Zookeeper
3.Kafka
4.maven
5.
一、binlog监控Mysql的库
二、编写FlinkCDC程序
1.添加pom文件
<?xml version="1.0" encoding="UTF-8"?>"http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> 4.0.0 com.lxz gmall-logger 0.0.1-SNAPSHOT gmall-20210909 Demo project for Spring Boot 1.8 UTF-8 UTF-8 2.4.1 ${java.version} ${java.version} 1.12.0 2.12 3.1.3 org.apache.flink flink-java ${flink.version} org.apache.flink flink-streaming-java_${scala.version} ${flink.version} org.apache.flink flink-connector-kafka_${scala.version} ${flink.version} org.apache.flink flink-clients_${scala.version} ${flink.version} org.apache.flink flink-cep_${scala.version} ${flink.version} org.apache.flink flink-json ${flink.version} com.alibaba fastjson 1.2.68 com.alibaba.ververica flink-connector-mysql-cdc 1.2.0 org.apache.hadoop hadoop-client ${hadoop.version} org.slf4j slf4j-api 1.7.25 org.slf4j slf4j-log4j12 1.7.25 org.apache.logging.log4j log4j-to-slf4j 2.14.0 org.springframework.boot spring-boot-starter-web org.springframework.kafka spring-kafka org.projectlombok lombok true org.springframework.boot spring-boot-dependencies ${spring-boot.version} pom import org.apache.maven.plugins maven-compiler-plugin 3.8.1 1.8 1.8 UTF-8 org.springframework.boot spring-boot-maven-plugin 2.3.0.RELEASE repackage boot com.lxz.gamll20210909.Gamll20210909Application
2.MykafkaUtil工具类
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import java.util.Properties;
public class MyKafkaUtil { private static String KAFKA_SERVER = "hadoop201:9092,hadoop202:9092,hadoop203:9092"; private static Properties properties = new Properties(); static { properties.setProperty("bootstrap.servers",KAFKA_SERVER); } public static FlinkKafkaProducergetKafkaSink(String topic){ return new FlinkKafkaProducer (topic,new SimpleStringSchema(),properties); } }
3.FlinkCDC主程序
import com.alibaba.fastjson.JSONObject; import com.alibaba.ververica.cdc.connectors.mysql.MySQLSource; import com.alibaba.ververica.cdc.connectors.mysql.table.StartupOptions; import com.alibaba.ververica.cdc.debezium.DebeziumDeserializationSchema; import com.alibaba.ververica.cdc.debezium.DebeziumSourceFunction; import com.lxz.gamll20210909.util.MyKafkaUtil; import io.debezium.data.Envelope; import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.streaming.api.datastream.DataStreamSource; 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.Schema; import org.apache.kafka.connect.data.Struct; import org.apache.kafka.connect.source.SourceRecord; public class Flink_CDCWithCustomerSchema { public static void main(String[] args) throws Exception { //1.创建执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); //2.创建Flink-MySQL-CDC的Source DebeziumSourceFunctionmysqlSource = MySQLSource. builder() .hostname("hadoop201") .port(3306) .username("root") .password("000000") .databaseList("gmall-20210712") .startupOptions(StartupOptions.latest()) // .startupOptions(KafkaOptions.StartupOptions.class) .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(); if (after != null) { Schema schema = after.schema(); for (Field field : schema.fields()) { data.put(field.name(), after.get(field.name())); } } //创建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.addSink(MyKafkaUtil.getKafkaSink("ods_base_db")); //5.执行任务 env.execute(); } }
三、结果
1.启动FlinkCDC主程序
2.在服务器上开一个kafka的消费者
bin/kafka-console-consumer.sh --bootstrap-server hadoop201:9092 --topic ods_base_db
3.在Mysql中插入数据看Kafka会不会消费
Mysql端
Kafka端
成功消费。