[基础架构] [Flink] Flink/Flink-CDC代码实现业务接入
简介DataStream 和 FlinkSQL 方式的对比DataStream 在 Flink1.12 和 1.13 都可以用而 FlinkSQL 只能在 Flink1.13 使用。DataStream 可以同时监控多库多表而 FlinkSQL 只能监控单表。方法 / 步骤一进行编码1.1 导入相关依赖dependenciesdependencygroupIdorg.apache.flink/groupIdartifactIdflink-java/artifactIdversion1.12.0/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-streaming-java_2.12/artifactIdversion1.12.0/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-clients_2.12/artifactIdversion1.12.0/version/dependencydependencygroupIdorg.apache.hadoop/groupIdartifactIdhadoop-client/artifactIdversion3.1.3/version/dependencydependencygroupIdmysql/groupIdartifactIdmysql-connector-java/artifactIdversion5.1.49/version/dependencydependencygroupIdorg.apache.flink/groupIdartifactIdflink-table-planner-blink_2.12/artifactIdversion1.12.0/version/dependencydependencygroupIdcom.ververica/groupIdartifactIdflink-connector-mysql-cdc/artifactIdversion2.0.0/version/dependencydependencygroupIdcom.alibaba/groupIdartifactIdfastjson/artifactIdversion1.2.75/version/dependency/dependenciesbuildpluginsplugingroupIdorg.apache.maven.plugins/groupId!-- 可以将依赖打到jar包中 --artifactIdmaven-assembly-plugin/artifactIdversion3.0.0/versionconfigurationdescriptorRefsdescriptorRefjar-with-dependencies/descriptorRef/descriptorRefs/configurationexecutionsexecutionidmake-assembly/idphasepackage/phasegoalsgoalsingle/goal/goals/execution/executions/plugin/plugins/build1.2 业务编码1.2.1 入口类importcom.ververica.cdc.connectors.mysql.MySqlSource;importcom.ververica.cdc.connectors.mysql.table.StartupOptions;importcom.ververica.cdc.debezium.DebeziumSourceFunction;importorg.apache.flink.streaming.api.CheckpointingMode;importorg.apache.flink.streaming.api.datastream.DataStreamSource;importorg.apache.flink.streaming.api.environment.StreamExecutionEnvironment;/** * Description: * * author: YangGC */publicclassFlinkCDC2{publicstaticvoidmain(String[]args)throwsException{//1.获取Flink 执行环境StreamExecutionEnvironmentenvStreamExecutionEnvironment.getExecutionEnvironment();env.setParallelism(1);//1.1 开启 Checkpointenv.enableCheckpointing(5000);env.getCheckpointConfig().setCheckpointTimeout(10000);env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);//// env.setStateBackend(new FsStateBackend(hdfs://hadoop102:8020/cdc-test/ck));//2.通过FlinkCDC构建SourceFunctionDebeziumSourceFunctionStringsourceFunctionMySqlSource.Stringbuilder().hostname(192.168.1.220).port(3308).username(root).password(useradmin)//flinkcdc 下面的所有表.databaseList(flinkcdc.*)// .tableList(flinkcdc.user_info)//使用自定义的反序列化器.deserializer(newCustomerDeserializationSchema()).startupOptions(StartupOptions.initial()).build();DataStreamSourceStringdataStreamSourceenv.addSource(sourceFunction);//3.数据打印dataStreamSource.print();//4.启动任务env.execute(FlinkCDC2);}}1.2.2 自定义反序列化器importcom.alibaba.fastjson.JSONObject;importcom.ververica.cdc.debezium.DebeziumDeserializationSchema;importio.debezium.data.Envelope;importorg.apache.flink.api.common.typeinfo.BasicTypeInfo;importorg.apache.flink.api.common.typeinfo.TypeInformation;importorg.apache.flink.util.Collector;importorg.apache.kafka.connect.data.Field;importorg.apache.kafka.connect.data.Schema;importorg.apache.kafka.connect.data.Struct;importorg.apache.kafka.connect.source.SourceRecord;importjava.util.List;/** * 自定义反序列化器 * Description: * * author: YangGC */publicclassCustomerDeserializationSchemaimplementsDebeziumDeserializationSchemaString{/** * { * db:, * tableName:, * before:{id:1001,name:...}, * after:{id:1001,name:...}, * op: * } */Overridepublicvoiddeserialize(SourceRecordsourceRecord,CollectorStringcollector)throwsException{//创建JSON对象用于封装结果数据JSONObjectresultnewJSONObject();//获取库名表名StringtopicsourceRecord.topic();String[]fieldstopic.split(\\.);result.put(db,fields[1]);result.put(tableName,fields[2]);//获取before数据Structvalue(Struct)sourceRecord.value();Structbeforevalue.getStruct(before);JSONObjectbeforeJsonnewJSONObject();if(before!null){//获取列信息Schemaschemabefore.schema();ListFieldfieldListschema.fields();for(Fieldfield:fieldList){beforeJson.put(field.name(),before.get(field));}}result.put(before,beforeJson);//获取after数据Structaftervalue.getStruct(after);JSONObjectafterJsonnewJSONObject();if(after!null){//获取列信息Schemaschemaafter.schema();ListFieldfieldListschema.fields();for(Fieldfield:fieldList){afterJson.put(field.name(),after.get(field));}}result.put(after,afterJson);//获取操作类型Envelope.OperationoperationEnvelope.operationFor(sourceRecord);result.put(op,operation);//输出数据collector.collect(result.toJSONString());}OverridepublicTypeInformationStringgetProducedType(){returnBasicTypeInfo.STRING_TYPE_INFO;}1.3 业务打包打包完成二Flink 作业到任务面板2.1 通过命令行添加任务把cdc-connector-1.0-SNAPSHOT-jar-with-dependencies.jar 包上传到 到flink主目录并运行下面命令行# 主要是配置入口类 指定flink的运行地址bin/flink run-m127.0.0.1:8081-ccom.yanggc.cdc.FlinkCDC2 ./cdc-connector-1.0-SNAPSHOT-jar-with-dependencies.jar作业面板查看job正在运行查看业务进行正常监控输出2.2 上传Jar包进行任务添加添加相关参数和命令行启动相关效果, 正常成功启动进行监控参考资料 致谢[1] flink-cdc-connectors