【Table/SQL Api】Flink Table/SQL Api表转流读取MySQL

引入依赖

jdbc依赖

flink-connector-jdbc + mysql-jdbc-driver 操作mysql数据库

        <!-- Flink-Connector-Jdbc --><dependency><groupId>org.apache.flink</groupId><artifactId>flink-connector-jdbc_${scala.binary.version}</artifactId></dependency><!-- mysql jdbc driver --><dependency><groupId>mysql</groupId><artifactId>mysql-connector-java</artifactId></dependency>

Table/SQL Api依赖

  1. Table/SQL Api 扩展依赖
  2. Table/SQL Api 基础依赖
  3. Table/SQL Api 和 DataStream Api 交互的依赖 bridge
  4. Flink Planner 依赖
        <!-- Table/SQL Api 依赖 --><dependency><groupId>org.apache.flink</groupId><artifactId>flink-table-api-java</artifactId></dependency><!-- Table/SQL Api 扩展依赖 --><dependency><groupId>org.apache.flink</groupId><artifactId>flink-table-common</artifactId></dependency><!-- bridge桥接器,主要负责Table API和 DataStream API的连接支持 --><dependency><groupId>org.apache.flink</groupId><artifactId>flink-table-api-java-bridge_${scala.binary.version}</artifactId></dependency><!-- Flink Planner 依赖 --><dependency><groupId>org.apache.flink</groupId><artifactId>flink-table-planner_${scala.binary.version}</artifactId></dependency>

对应版本在这 (项目Flink版本为1.14.5

image-20231210161727111

Flink读写MySQL工具类

Table Api 环境加载

Table API和SQL Api都是基于Table接口

Table Api上下文环境有3种类型

  1. TableEnvironment:只支持Batch作业
  2. BatchTableEnvironment:只支持Batch作业
  3. StreamTableEnvironment: 支持流计算【用这个】

Planner(查询处理器)

Planner(查询处理器):解析sql、优化sql和执行sql

Flink Planner的类型:

  1. Flink Planner (Old Planner)
  2. Blink Planner (Flink 1.14之前需要手动导入依赖)

Blink Planner从Flink 1.11版本开始为Flink-table的默认查询处理器

Blink Planner使得Table Api & SQL 层实现了流批统一

Catalog对象

Catalog对象是提供了元数据信息,数据源与数据表的信息则存储在Catalog中

// 创建Catalog对象
new JdbcCatalog(catalog_name, database, username, passwd, url);

Catalog对象是接口

Catalog接口的实现:(Flink 1.14版本之前)

  1. PG (PostgresSQL) Catalog
  2. HiveCatalog
  3. Mysql Catalog (Flink 1.15 才有)

DDL与数据库表结构必须一模一样,建立映射,这种方式数据库表结构如果变化,代码也必须随之变化重新打包,因此这种方式用的不多,一般catalog会用的比较多。

但由于项目Flink依赖用的是1.14.5,因此还是使用DDL语句实现。

代码实现

public class MysqlUtil {/*** 数据库连接对象*/private static Connection connection = null;/*** SQL语句对象*/private static PreparedStatement preparedStatement = null;/*** 结果集对象*/private static ResultSet rs = null;/*** 使用 Flink Table/SQL Api 读取Mysql** @param env:           流计算上下文环境* @param parameterTool: 参数工具* @param clazz:         流水线输出对象的类* @param tableName:     表名* @param ddlString:     DDL字符串* @param sql:           SQL查询语句* @return DataStream<T>:DataStream对象*/public static <T> DataStream<T> readWithTableOrSQLApi(StreamExecutionEnvironment env,ParameterTool parameterTool,Class<T> clazz,String tableName,String ddlString,String sql) throws Exception {// 创建TableApi运行环境EnvironmentSettings bsSettings =EnvironmentSettings.newInstance()// Flink 1.14不需要再设置 Planner//.useBlinkPlanner()// 设置流计算模式.inStreamingMode().build();// 创建StreamTableEnvironment实例StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, bsSettings);// 指定方言 (选择使用SQL语法还是HQL语法)tableEnv.getConfig().setSqlDialect(SqlDialect.DEFAULT);// 编写DDL ( 数据定义语言 )String ddl = buildMysqlDDL(parameterTool, tableName, ddlString);// StreamTableEnvironment注册虚拟表tableEnv.executeSql(ddl);// 查询结果是Table对象Table table = tableEnv.sqlQuery(sql);// 将Table对象转换为DataStream对象return tableEnv.toDataStream(table, clazz);}/*** 根据参数生成MySQL的DDL语句** @param parameterTool  参数工具,用于获取MySQL连接信息* @param tableName      要创建的表名* @param ddlFieldString 表字段的DDL语句* @return 生成的完整的MySQL DDL语句*/public static String buildMysqlDDL(ParameterTool parameterTool,String tableName,String ddlFieldString) {// 从参数工具中获取mysql连接的urlString url = parameterTool.get(ParameterConstants.Mysql_URL);// 从参数工具中获取mysql连接的用户名String username = parameterTool.get(ParameterConstants.Mysql_USERNAME);// 从参数工具中获取mysql连接的密码String passwd = parameterTool.get(ParameterConstants.Mysql_PASSWD);// 从参数工具中获取MySQL的驱动程序String driver = parameterTool.get(ParameterConstants.Mysql_DRIVER);// 返回完整的DDL语句return "CREATE TABLE IF NOT EXISTS " +tableName +" (\n" +ddlFieldString +")" +" WITH (\n" +"'connector' = 'jdbc',\n" +"'driver' = '" + driver + "',\n" +"'url' = '" + url + "',\n" +"'username' = '" + username + "',\n" +"'password' = '" + passwd + "',\n" +"'table-name' = '" + tableName + "'\n" +")";}/*** 初始化 jdbc Connection*/public static Connection init(ParameterTool parameterTool) {String _url = parameterTool.get(ParameterConstants.Mysql_URL);String _username = parameterTool.get(ParameterConstants.Mysql_USERNAME);String _passwd = parameterTool.get(ParameterConstants.Mysql_PASSWD);try {connection = DriverManager.getConnection(_url, _username, _passwd);} catch (Exception e) {throw new RuntimeException(e);}return connection;}/*** 生成 PreparedStatement*/public static PreparedStatement initPreparedStatement(String sql) {try {preparedStatement = connection.prepareStatement(sql);} catch (Exception e) {throw new RuntimeException(e);}return preparedStatement;}/*** 关闭 jdbc Connection*/public static void close() {try {if (preparedStatement != null) {preparedStatement.close();}if (connection != null) {connection.close();}} catch (Exception e) {throw new RuntimeException(e);}}/*** 关闭 PreparedStatement*/public static void closePreparedStatement() {try {if (preparedStatement != null) {preparedStatement.close();}} catch (Exception e) {throw new RuntimeException(e);}}/*** 关闭 ResultSet*/public static void closeResultSet() {try {if (rs != null) {rs.close();}} catch (Exception e) {throw new RuntimeException(e);}}/*** 执行 sql 语句*/public static ResultSet executeQuery(PreparedStatement ps) {preparedStatement = ps;try {rs = preparedStatement.executeQuery();} catch (Exception e) {throw new RuntimeException(e);}return rs;}}

测试一下

测试库中有个tb_user表

image-20231210174346826

创建与表映射的实体类

@Data
public class UserPO {private Long id;private String name;
}
class MysqlUtilTest {@DisplayName("测试使用 Flink Table/SQL Api 读取Mysql")@Testpublic void testReadWithTableOrSQLApi() throws Exception {// 初始化环境StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();// 设置并行度1env.setParallelism(1);// 获取参数工具实例ParameterTool parameterTool = ParameterUtil.getParameters();/* ************************ CREATE 语句用于向当前或指定的 Catalog 中注册表。* 注册后的表、视图和函数可以在 SQL 查询中使用** *********************/// 表名String tableName = "tb_user";// 表字段ddlString ddlFieldString ="id BIGINT,\n" +"name STRING \n";// 查询表的全部字段String sql = "SELECT * FROM " + tableName;DataStream<UserPO> rowDataStream =MysqlUtil.readWithTableOrSQLApi(env,parameterTool,UserPO.class,tableName,ddlFieldString,sql);rowDataStream.print("mysql");env.execute();}
}

image-20231210174720832

查询成功!

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.hqwc.cn/news/263795.html

如若内容造成侵权/违法违规/事实不符,请联系编程知识网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

金融行业文件摆渡,如何兼顾安全和效率?

金融行业是数据密集型产业&#xff0c;每时每刻都会产生海量的数据&#xff0c;业务开展时&#xff0c;数据在金融机构内部和内外部快速流转&#xff0c;进入生产的各个环节。 为了保障基础的数据安全和网络安全&#xff0c;金融机构采用网络隔离的方式来隔绝外部网络的有害攻击…

苹果IOS在Safari浏览器中将网页添加到主屏幕做伪Web App,自定义图标,启动动画,自定义名称,全屏应用pwa

在ios中我们可以使用Safari浏览自带的将网页添加到主屏幕上&#xff0c;让我们的web页面看起来像一个本地应用程序一样&#xff0c;通过桌面APP图标一打开&#xff0c;直接全屏展示&#xff0c;就像在APP中效果一样&#xff0c;完全体会不到你是在浏览器中。 1.网站添加样式 在…

《算法与数据结构》答疑

答疑 问题一问题二问题三问题四 问题一 在匹配成功时&#xff0c;在返回子串位置那里&#xff0c;为什么不是i-t的长度啊&#xff0c;为什么还要加一 问题二 问题三 问题四 问&#xff1a;如果题目让我们构造一个哈夫曼树&#xff0c;像我发的这个例题的话&#xff0c;我画成我…

详解异常 ! !(对异常有一个全面的认识)

【本章目标】 1. 异常概念与体系结构 2. 异常的处理方式 3. 异常的处理流程 4. 自定义异常类 1. 异常的概念与体系结构 1.1 异常的概念 在生活中&#xff0c;一个人表情痛苦&#xff0c;出于关心&#xff0c;可能会问&#xff1a;你是不是生病了&#xff0c;需要我陪你去看医…

零基础一看就会?Python实现性能自动化测试竟然如此简单

一、思考❓❔ 1.什么是性能自动化测试? 性能 系统负载能力超负荷运行下的稳定性系统瓶颈自动化测试 使用程序代替手工提升测试效率性能自动化 使用代码模拟大批量用户让用户并发请求多页面多用户并发请求采集参数&#xff0c;统计系统负载能力生成报告 2.Python中的性能自动化…

[GPT]Andrej Karpathy微软Build大会GPT演讲(下)--该如何使用GPT助手

该如何使用GPT助手--将GPT助手模型应用于问题 现在我要换个方向,让我们看看如何最好地将 GPT 助手模型应用于您的问题。 现在我想在一个具体示例的场景里展示。让我们在这里使用一个具体示例。 假设你正在写一篇文章或一篇博客文章,你打算在最后写这句话。 加州的人口是阿拉…

上班必备——项目部署环境

大家都知道&#xff0c;互联网行业有很多的岗位&#xff0c;前端&#xff0c;后端&#xff0c;产品&#xff0c;测试&#xff0c;ui等。 ui&#xff0c;产品和测试的同事在前端开发的过程中&#xff0c;都会时刻关注着进度&#xff0c;是要看页面效果的&#xff0c;这个时候怎…

java+springboot+ssm学生社团管理系统76c2e

本系统包括前台和后台两个部分。前台主要是展示社团列表、社团风采、社团活动、新闻列表等&#xff0c;前台登录后进入个人中心&#xff0c;在个人中心能申请加入社团、查看参加的社团活动等&#xff1b;后台为管理员与社团负责人使用&#xff0c;应用于对社团的管理及内容发布…

软件无线电SDR-频谱采集python实现

sdr做的频谱采集&#xff0c;保存的500张频谱图&#xff0c;能看出来是什么东西吗&#xff1f;

kafka入门(四):消费者

消费者 (Consumer ) 消费者 订阅 Kafka 中的主题 (Topic) &#xff0c;并 拉取消息。 消费者群组&#xff08; Consumer Group&#xff09; 每一个消费者都有一个对应的 消费者群组。 一个群组里的消费者订阅的是同一个主题&#xff0c;每个消费者接收主题的一部分分区的消息…

Minio保姆级教程

转载自&#xff1a;www.javaman.cn Minio服务器搭建和整合 1、centos安装minio 1.1、创建安装目录 mkdir -p /home/minio1.2、在线下载minio #进入目录 cd /home/minio #下载 wget https://dl.minio.io/server/minio/release/linux-amd64/minio1.3、minio配置 1.3.1、添加…

Shell数组函数:函数

一、概述 概念&#xff1a; 函数是一段完成特定功能的代码片段&#xff08;块&#xff09;在shell中定义了函数&#xff0c;就可以使代码模块化&#xff0c;使于复用代码注意函数必须先定义才可以使用。 重点&#xff1a; 传参 $1,$2局部变量 local返回值 return 即$? 二、定…