首页/新闻资讯/正文详情

[基础架构] [Flink] Flink/Flink-CDC代码实现业务接入

发布时间:2026/9/23 8:50:15 来源:云帆数科 栏目:资讯中心
[基础架构] [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

相关推荐

搞懂秒表的读法:这高频面试题坑了多少人
搞懂秒表的读法:这高频面试题坑了多少人

搞懂秒表的读法:这高频面试题坑了多少人 看了一堆教程还是不会写项目?别急着骂自己笨,很多时候是基础概念没吃透。最近整理后端高频面试题,发现“秒表的读法”这个看似简单的点,居然能把一堆自诩熟练的开发者问懵。不是让你去体育场上看表,而是在编程里… · 2026/9/23 8:49:58

基于Spring Cloud和Vue3的智慧云停车场系统设计与实践
基于Spring Cloud和Vue3的智慧云停车场系统设计与实践

1. 项目背景与核心价值停车难问题已经成为现代城市管理的痛点。传统停车场管理系统存在信息孤岛、资源利用率低、用户体验差等问题。我们团队基于Spring Cloud微服务架构和Vue3前端技术栈,开发了一套智慧云停车场服务管理系统,实现了停车场资源的智能化管… · 2026/9/23 8:49:39

3个维度拆解卡通可爱壁纸生成,面试必问的技术选型指南
3个维度拆解卡通可爱壁纸生成,面试必问的技术选型指南

3个维度拆解卡通可爱壁纸生成,面试必问的技术选型指南 官方文档翻了三遍还是头大?别慌,你不是一个人。做技术选型最怕的就是陷在文档海洋里找不到北,特别是像【卡通可爱壁纸】这种看似简单实则坑多的需求。面试官最爱拿这种小需求考你:如果让你从零实现… · 2026/9/23 8:49:26

3天搞定网众无盘教程图解原理,拒绝堆砌
3天搞定网众无盘教程图解原理,拒绝堆砌

3天搞定网众无盘教程图解原理,拒绝堆砌 报错一堆看不懂 StackTrace?别慌。 很多刚接触网众无盘的朋友,一看到满屏红色的 Error 信息就头大,根本不知道从哪下手。 今天咱们不整虚的,直接上 图解原理 。… · 2026/9/23 12:08:49

基于Python的学业预警系统:从数据清洗到Flask接口的完整实战
基于Python的学业预警系统:从数据清洗到Flask接口的完整实战

简介:这份Python项目源码面向高校教务管理者、辅导员以及计算机相关专业的课程设计学习者,核心目标是借助数据分析提前识别学业困难学生,降低辍学风险。系统围绕数据采集与整合、基于逻辑回归或随机森林等算法的风险评估模型、实时监控、预警… · 2026/9/23 12:08:49

杨瑞凯速查手册:3步搞定项目搭建避坑指南
杨瑞凯速查手册:3步搞定项目搭建避坑指南

杨瑞凯速查手册:3步搞定项目搭建避坑指南 官方文档翻了三遍还是懵?别慌,我直接上干货。这份【杨瑞凯】实战项目的【速查手册】,就是为了解决你“看文档像看天书”的痛点。… · 2026/9/23 12:08:49

NVIDIA Triton Inference Server 架构解析与核心特性全景指南
NVIDIA Triton Inference Server 架构解析与核心特性全景指南

NVIDIA Triton Inference Server 架构解析与核心特性全景指南 【免费下载链接】server The Triton Inference Server provides an optimized cloud and edge inferencing solution. 项目地址: https://gitcode.com/gh_mirrors/server117/server Triton Inference Serve… · 2026/9/23 12:08:43

5步搞定联想s720运维,最佳实践让项目落地不再难
5步搞定联想s720运维,最佳实践让项目落地不再难

5步搞定联想s720运维,最佳实践让项目落地不再难 看了一堆教程还是不会写项目?别急,这其实是90%初学者的通病。理论背得滚瓜烂熟,一到真实场景就卡壳。 真正的 最佳实践 ,不是让你背更多命令,而是建立一套可复用的运维思维。今天我们就以… · 2026/9/23 12:08:37

化纤面料全解析:从聚酯纤维到混纺,教你选对衣服不踩坑
化纤面料全解析:从聚酯纤维到混纺,教你选对衣服不踩坑

很多人在挑选衣服时,第一反应就是翻看吊牌上的成分表,看到“聚酯纤维”“锦纶”这些字眼,眉头就皱起来了。我做了十来年面料采购和开发,几乎每周都会被人问到“化纤面料到底是什么”“是不是就是塑料”“穿上是不是闷得慌”。今天… · 2026/9/23 12:08:37

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码