温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

Flink批处理之读写Mysql

发布时间:2020-07-06 11:01:39 来源:网络 阅读:3121 作者:兴趣e族 栏目:大数据

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,不然数据库中会出现乱码。
向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI