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>
@H_404_5@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
@H_404_5@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(); } }
@H_404_5@4、注意事项
- 数据写入MysqL的DataSet泛型要求是row,需要转换;
- 数据读取的结果也是row类型,不能直接print,需要转换;
- 数据写入后一定要加上env.execute(),触发任务执行;
- 涉及到中文的,需要转换成UTF-8,不然数据库中会出现乱码。
版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 [email protected] 举报,一经查实,本站将立刻删除。