首页 / MYSQL / Flink批处理之读写Mysql
Flink批处理之读写Mysql
内容导读
互联网集市收集整理的这篇技术教程文章主要介绍了Flink批处理之读写Mysql,小编现在分享给大家,供广大互联网技能从业者学习和参考。文章包含3092字,纯文字阅读大概需要5分钟。
内容图文
![Flink批处理之读写Mysql](/upload/InfoBanner/zyjiaocheng/887/4668a7f19a524fca859289f31cb7376e.jpg)
1、添加Maven坐标
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.48</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-jdbc_2.12</artifactId>
<version>1.8.0</version>
</dependency>
2、建表
CREATE TABLE `temp` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`name` varchar(255) DEFAULT NULL,
`time` varchar(255) DEFAULT NULL,
`type` bigint(20) DEFAULT NULL,
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8
3、 Show Code
package com.fwmagic.flink.batch;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.typeinfo.BasicTypeInfo;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.io.jdbc.JDBCInputFormat;
import org.apache.flink.api.java.io.jdbc.JDBCOutputFormat;
import org.apache.flink.api.java.operators.DataSource;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.types.Row;
import java.util.concurrent.TimeUnit;
public class BatchDemoOperatorMysql {
public static void main(String[] args) throws Exception {
String driverClass = "com.mysql.jdbc.Driver";
String dbUrl = "jdbc:mysql://localhost:3306/test";
String userNmae = "root";
String passWord = "123456";
String sql = "insert into test.temp (name,time,type) values (?,?,?)";
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
/**
* 文件内容:
* 关羽,2019-10-14 00:00:01,1
* 张飞,2019-10-14 00:00:02,2
* 赵云,2019-10-14 00:00:03,3
*/
String filePath = "/Users/temp/data.csv";
//读csv文件内容,转成Row对象
DataSet<Row> outputData = env.readCsvFile(filePath).fieldDelimiter(",").types(String.class, String.class, Long.class).map(new MapFunction<Tuple3<String, String, Long>, Row>() {
@Override
public Row map(Tuple3<String, String, Long> t) throws Exception {
Row row = new Row(3);
row.setField(0, t.f0.getBytes("UTF-8"));
row.setField(1, t.f1.getBytes("UTF-8"));
row.setField(2, t.f2.longValue());
return row;
}
});
//将Row对象写到mysql
outputData.output(JDBCOutputFormat.buildJDBCOutputFormat()
.setDrivername(driverClass)
.setDBUrl(dbUrl)
.setUsername(userNmae)
.setPassword(passWord)
.setQuery(sql)
.finish());
//触发执行
env.execute("insert data to mysql");
System.out.println("mysql写入成功!");
TimeUnit.SECONDS.sleep(6);
//读mysql
DataSource<Row> dataSource = env.createInput(JDBCInputFormat.buildJDBCInputFormat()
.setDrivername(driverClass)
.setDBUrl(dbUrl)
.setUsername(userNmae)
.setPassword(passWord)
.setQuery("select * from temp")
.setRowTypeInfo(new RowTypeInfo(BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.STRING_TYPE_INFO, BasicTypeInfo.LONG_TYPE_INFO))
.finish());
//获取数据并打印
dataSource.map(new MapFunction<Row, String>() {
@Override
public String map(Row value) throws Exception {
System.out.println(value);
return value.toString();
}
}).print();
}
}
4、注意事项
- 数据写入mysql的DataSet泛型要求是row,需要转换;
- 数据读取的结果也是row类型,不能直接print,需要转换;
- 数据写入后一定要加上env.execute(),触发任务执行;
- 涉及到中文的,需要转换成UTF-8,不然数据库中会出现乱码。
内容总结
以上是互联网集市为您收集整理的Flink批处理之读写Mysql全部内容,希望文章能够帮你解决Flink批处理之读写Mysql所遇到的程序开发问题。 如果觉得互联网集市技术教程内容还不错,欢迎将互联网集市网站推荐给程序员好友。
内容备注
版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 gblab@vip.qq.com 举报,一经查实,本站将立刻删除。
内容手机端
扫描二维码推送至手机访问。