IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 大数据 -> flink 学习(三)mysql 作为数据源 -> 正文阅读

[大数据]flink 学习(三)mysql 作为数据源

一、jdbc连接

1、引入依赖

<!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-jdbc -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-jdbc_2.12</artifactId>
    <version>1.14.4</version>
    <scope>provided</scope>
</dependency>

2、实体类

public class UserInfo {
   private Integer userId;
   private String userName;
   private String userRealName;
   private String userPwd;
   private String userTel;
   private String userEmail;
   private Integer userStatus;
   private Date userCreateTime;
   private Date userUpdateTime;
}

二、自定义 SourceFunction

public class MysqlRichParallelSource extends RichParallelSourceFunction<User> {


    private boolean close = false;

    @Override
    public void run(SourceContext<User> out) throws Exception {
        String url = "jdbc:mysql://192.168.100.88:3306/newframe?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai&useSSL=false&allowMultiQueries=true&rewriteBatchedStatements=true";
        String sql = "select * from user_info";
        Connection conn = null;
        PreparedStatement ps = null;
        ResultSet rs = null;
        try {
            conn = DriverManager.getConnection(url, "root", "root");
            ps = conn.prepareStatement(sql);
        } catch (SQLException e) {
            e.printStackTrace();
        }

        while (!close) {
            rs = ps.executeQuery();
            while (rs.next()) {
                Integer userId = rs.getInt("user_id");
                String userName = rs.getString("user_name");
                String userRealName = rs.getString("user_real_name");
                User user = new User()
                        .setUserId(userId)
                        .setUserName(userName)
                        .setUserRealName(userRealName);
                //收集数据
                out.collect(user);
            }
            Thread.sleep(5000);
            cancel();
        }
        close(conn, ps, rs);
    }

    @Override
    public void cancel() {
        close = true;
    }

    public void close(Connection conn, PreparedStatement ps, ResultSet rs) throws Exception {
        if (conn != null) {
            conn.close();
        }
        if (ps != null) {
            ps.close();
        }
        if (rs != null) {
            rs.close();
        }
    }
}

三、测试

	@Test
    public void fromMysqlTest() throws Exception {
        // flink 流执行环境
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        //设置模式 STREAMING
        env.setRuntimeMode(RuntimeExecutionMode.STREAMING);
                //添加mysql数据源
        DataStreamSource<User> source = env.addSource(new MysqlRichParallelSource());
        //打印结果
        source.print();
        //开始执行
        env.execute("flink streaming from mysql");
    }

结果:多线程同时从mysql获取数据,启动了8个线程,每个数据获取了8次

2> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
8> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
3> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
5> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
4> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
1> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
7> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
4> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
5> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
3> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=23, userName=user2, userRealName=用户2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)

  大数据 最新文章
实现Kafka至少消费一次
亚马逊云科技:还在苦于ETL?Zero ETL的时代
初探MapReduce
【SpringBoot框架篇】32.基于注解+redis实现
Elasticsearch:如何减少 Elasticsearch 集
Go redis操作
Redis面试题
专题五 Redis高并发场景
基于GBase8s和Calcite的多数据源查询
Redis——底层数据结构原理
上一篇文章      下一篇文章      查看所有文章
加:2022-04-09 18:28:13  更:2022-04-09 18:32:31 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 -2025/1/16 13:08:46-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码