flink 将mysql作为Source和Sink的代码示例
来源:SegmentFault
2023-02-19 10:11:29
0浏览
收藏
所属专题:Apache Flink 2.3 实时流处理工程实践专题
- 从事件时间、窗口与状态,到 CDC、Exactly-once 和生产运维
在数据库实战开发的过程中,我们经常会遇到一些这样那样的问题,然后要卡好半天,等问题解决了才发现原来一些细节知识点还是没有掌握好。今天golang学习网就整理分享《flink 将mysql作为Source和Sink的代码示例》,聊聊MySQL、flink,希望可以帮助到正在努力赚钱的你。
1.maven导入
mysql mysql-connector-java 5.1.34
2.SourceFromMySQL工具类java代码
import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.source.RichSourceFunction; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.ResultSet; /** * @Description TODO * @Author zytshijack * @Date 2019-05-14 08:32 * @Version 1.0 */ public class SourceFromMySQL extends RichSourceFunction{ //Student就是一个自定义类 PreparedStatement ps; private Connection connection; /** * open() 方法中建立连接,这样不用每次 invoke 的时候都要建立连接和释放连接。 * * @param parameters * @throws Exception */ @Override public void open(Configuration parameters) throws Exception { super.open(parameters); connection = getConnection(); String sql = "select * from Student;"; ps = this.connection.prepareStatement(sql); } /** * 程序执行完毕就可以进行,关闭连接和释放资源的动作了 * * @throws Exception */ @Override public void close() throws Exception { super.close(); if (connection != null) { //关闭连接和释放资源 connection.close(); } if (ps != null) { ps.close(); } } /** * DataStream 调用一次 run() 方法用来获取数据 * * @param ctx * @throws Exception */ @Override public void run(SourceContext ctx) throws Exception { ResultSet resultSet = ps.executeQuery(); while (resultSet.next()) { Student student = new Student( resultSet.getInt("id"), resultSet.getString("name").trim(), resultSet.getString("password").trim(), resultSet.getInt("age")); ctx.collect(student); } } @Override public void cancel() { } private static Connection getConnection() { Connection con = null; try { Class.forName("com.mysql.jdbc.Driver"); con = DriverManager.getConnection("jdbc:mysql://localhost:3306/flink?useUnicode=true&characterEncoding=UTF-8", "root", "123"); } catch (Exception e) { System.out.println("-----------mysql get connection has exception , msg = "+ e.getMessage()); } return con; } }
3.SinkToMySQL工具类java代码
import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.sink.RichSinkFunction; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; /** * @Description TODO * @Author zytshijack * @Date 2019-05-14 09:10 * @Version 1.0 */ public class SinkToMySQL extends RichSinkFunction{ PreparedStatement ps; private Connection connection; /** * open() 方法中建立连接,这样不用每次 invoke 的时候都要建立连接和释放连接 * * @param parameters * @throws Exception */ @Override public void open(Configuration parameters) throws Exception { super.open(parameters); connection = getConnection(); String sql = "insert into Student(id, name, password, age) values(?, ?, ?, ?);"; ps = this.connection.prepareStatement(sql); } @Override public void close() throws Exception { super.close(); //关闭连接和释放资源 if (connection != null) { connection.close(); } if (ps != null) { ps.close(); } } /** * 每条数据的插入都要调用一次 invoke() 方法 * * @param value * @param context * @throws Exception */ @Override public void invoke(Student value, Context context) throws Exception { //组装数据,执行插入操作 ps.setInt(1, value.getId()); ps.setString(2, value.getName()); ps.setString(3, value.getPassword()); ps.setInt(4, value.getAge()); ps.executeUpdate(); } private static Connection getConnection() { Connection con = null; try { Class.forName("com.mysql.jdbc.Driver"); con = DriverManager.getConnection("jdbc:mysql://localhost:3306/flink?useUnicode=true&characterEncoding=UTF-8&useSSL=false", "root", "123"); } catch (Exception e) { System.out.println("-----------mysql get connection has exception , msg = "+ e.getMessage()); } return con; } }
4.主函数使用工具类
env.addSource(new SourceFromMySQL()) ... stream.addSink(new SinkToMySQL())
本篇关于《flink 将mysql作为Source和Sink的代码示例》的介绍就到此结束啦,但是学无止境,想要了解学习更多关于数据库的相关知识,请关注golang学习网公众号!
版本声明
本文转载于:SegmentFault 如有侵犯,请联系study_golang@163.com删除
介绍一款免费好用的可视化数据库管理工具
- 上一篇
- 介绍一款免费好用的可视化数据库管理工具
- 下一篇
- Mysql从会用到用好
评论列表
-
- 昏睡的爆米花
- 这篇技术贴出现的刚刚好,太详细了,真优秀,码住,关注师傅了!希望师傅能多写数据库相关的文章。
- 2023-03-27 21:47:52
查看更多
最新文章
-
- 数据库 · MySQL | 21小时前 |
- MySQL 8.4 隐藏索引怎么试:不删索引也能验证查询计划
- 133浏览 收藏
-
- 数据库 · MySQL | 1天前 |
- MySQL 分区表怎么处理跨分区唯一键:分区列约束与建表取舍
- 501浏览 收藏
-
- 数据库 · MySQL | 2天前 | MySQL · 数据库设计 · 索引优化 · mysql UUID_TO_BIN BIN_TO_UUID swap_flag 二进制UUID
- MySQL UUID_TO_BIN 如何保持时间有序:swap_flag 与 BIN_TO_UUID 还原边界
- 261浏览 收藏
-
- 数据库 · MySQL | 2天前 |
- MySQL 隐形列怎么兼容 SELECT *:INVISIBLE COLUMN 的灰度加字段方法
- 122浏览 收藏
-
- 数据库 · MySQL | 2天前 |
- MySQL 分页越翻越慢怎么改:从深分页到覆盖索引的一个列表接口
- 475浏览 收藏
-
- 数据库 · MySQL | 2天前 | MySQL · 索引 · 数据库 · 运维 · mysql lock ALTER TABLE ALGORITHM metadata lock
- MySQL 大表新增索引如何降低阻塞:在线变更算法、锁窗口与回滚检查
- 224浏览 收藏
-
- 数据库 · MySQL | 2天前 | MySQL · 主从复制 · 运维排障 · MySQL 8.4 replicate-wild-do-table 复制过滤
- MySQL 8.4 复制过滤为何漏掉目标表:replicate-wild-do-table、库名匹配与配置验收
- 405浏览 收藏
-
- 数据库 · MySQL | 2天前 | MySQL · 递归 · SQL查询 · 故障排查 · CTE · WITH RECURSIVE cte_max_recursion_depth MySQL CTE 递归查询 递归成员
- MySQL CTE 递归查询为什么会提前停止:锚点、递归成员与 cte_max_recursion_depth
- 438浏览 收藏
查看更多
课程推荐
-
- 前端进阶之JavaScript设计模式
- 设计模式是开发人员在软件开发过程中面临一般问题时的解决方案,代表了最佳的实践。本课程的主打内容包括JS常见设计模式以及具体应用场景,打造一站式知识长龙服务,适合有JS基础的同学学习。
- 543次学习
-
- GO语言核心编程课程
- 本课程采用真实案例,全面具体可落地,从理论到实践,一步一步将GO核心编程技术、编程思想、底层实现融会贯通,使学习者贴近时代脉搏,做IT互联网时代的弄潮儿。
- 516次学习
-
- 简单聊聊mysql8与网络通信
- 如有问题加微信:Le-studyg;在课程中,我们将首先介绍MySQL8的新特性,包括性能优化、安全增强、新数据类型等,帮助学生快速熟悉MySQL8的最新功能。接着,我们将深入解析MySQL的网络通信机制,包括协议、连接管理、数据传输等,让
- 500次学习
-
- JavaScript正则表达式基础与实战
- 在任何一门编程语言中,正则表达式,都是一项重要的知识,它提供了高效的字符串匹配与捕获机制,可以极大的简化程序设计。
- 487次学习
-
- 从零制作响应式网站—Grid布局
- 本系列教程将展示从零制作一个假想的网络科技公司官网,分为导航,轮播,关于我们,成功案例,服务流程,团队介绍,数据部分,公司动态,底部信息等内容区块。网站整体采用CSSGrid布局,支持响应式,有流畅过渡和展现动画。
- 485次学习
查看更多
AI推荐
-
- SuperCLUE
- SuperCLUE是权威的中文大语言模型综合评测基准,涵盖语言理解、知识应用、AI Agent智能体及安全性等12项核心能力。通过多轮对话与客观测试,定期发布榜单与技术报告,为模型研发、优化及行业选型提供科学依据。
- 27次使用
-
- Gradio
- Gradio是一个用于构建机器学习和数据科学Web应用的开源Python库。支持快速创建交互界面,获Google、Meta等大厂青睐,适合模型演示、部署反馈及调试。
- 25次使用
-
- AutoGPT
- AutoGPT是基于GPT-4的开源AI代理平台,拥有超10万GitHub星标。本文介绍其低代码界面、自动化工作流功能、系统配置要求及安装步骤,助您高效部署和管理AI Agent。
- 27次使用
-
- 腾讯扣叮
- 腾讯扣叮是腾讯推出的6-18岁青少年编程学习平台,依托游戏与AI技术,提供图形化编程、3D创作、虚拟实验室及丰富赛事课程,助力培养计算思维与创新能力。
- 25次使用
-
- 堆友AI学习
- 堆友AI学习是堆友推出的专业AI设计教育平台,提供从基础到进阶的线上课程及线下实训营。结合阿里国际AITIC认证,通过视频教程、笔记分享和实战案例,帮助设计师掌握AIGC技能,提升职业竞争力。
- 25次使用
查看更多
相关文章
-
- MySQL 明明加了索引,为什么查询还是很慢?先查这 6 个点
- 2026-06-27 374浏览
-
- golang MySQL实现对数据库表存储获取操作示例
- 2022-12-22 499浏览
-
- golang 基于 mysql 简单实现分布式读写锁
- 2023-01-07 384浏览
-
- 详解如何利用GORM实现MySQL事务
- 2023-01-07 184浏览
-
- Go语言实现操作MySQL的基础知识总结
- 2023-01-23 265浏览

