Flink的DataSource三部曲之一:直接API

欢迎访问我的GitHub

https://github.com/zq2599/blog_demos

内容:所有原创文章分类汇总及配套源码,涉及Java、Docker、Kubernetes、DevOPS等;

本文是《Flink的DataSource三部曲》系列的第一篇,该系列旨在通过实战学习和了解Flink的DataSource,为以后的深入学习打好基础,由以下三部分组成:

  1. 直接API:即本篇,除了准备环境和工程,还学习了StreamExecutionEnvironment提供的用来创建数据来的API;
  2. 内置connector:StreamExecutionEnvironment的addSource方法,入参可以是flink内置的connector,例如kafka、RabbitMQ等;
  3. 自定义:StreamExecutionEnvironment的addSource方法,入参可以是自定义的SourceFunction实现类;

Flink的DataSource三部曲文章链接

  1. 《Flink的DataSource三部曲之一:直接API》
  2. 《Flink的DataSource三部曲之二:内置connector》
  3. 《Flink的DataSource三部曲之三:自定义》

关于Flink的DataSource

官方对DataSource的解释:Sources are where your program reads its input from,即DataSource是应用的数据来源,如下图的两个红框所示: 在这里插入图片描述

DataSource类型

对于常见的文本读入、kafka、RabbitMQ等数据来源,可以直接使用Flink提供的API或者connector,如果这些满足不了需求,还可以自己开发,下图是我按照自己的理解梳理的: 在这里插入图片描述

环境和版本

熟练掌握内置DataSource的最好办法就是实战,本次实战的环境和版本如下:

  1. JDK:1.8.0_211
  2. Flink:1.9.2
  3. Maven:3.6.0
  4. 操作系统:macOS Catalina 10.15.3 (MacBook Pro 13-inch, 2018)
  5. IDEA:2018.3.5 (Ultimate Edition)

源码下载

如果您不想写代码,整个系列的源码可在GitHub下载到,地址和链接信息如下表所示(https://github.com/zq2599/blog_demos):

名称

链接

备注

项目主页

https://github.com/zq2599/blog_demos

该项目在GitHub上的主页

git仓库地址(https)

https://github.com/zq2599/blog_demos.git

该项目源码的仓库地址,https协议

git仓库地址(ssh)

git@github.com:zq2599/blog_demos.git

该项目源码的仓库地址,ssh协议

这个git项目中有多个文件夹,本章的应用在<font color="blue">flinkdatasourcedemo</font>文件夹下,如下图红框所示: 在这里插入图片描述

环境和版本

本次实战的环境和版本如下:

  1. JDK:1.8.0_211
  2. Flink:1.9.2
  3. Maven:3.6.0
  4. 操作系统:macOS Catalina 10.15.3 (MacBook Pro 13-inch, 2018)
  5. IDEA:2018.3.5 (Ultimate Edition)

创建工程

  1. 在控制台执行以下命令就会进入创建flink应用的交互模式,按提示输入gourpId和artifactId,就会创建一个flink应用(我输入的groupId是<font color="blue">com.bolingcavalry</font>,artifactId是<font color="blue">flinkdatasourcedemo</font>):

    mvn
    archetype:generate
    -DarchetypeGroupId=org.apache.flink
    -DarchetypeArtifactId=flink-quickstart-java
    -DarchetypeVersion=1.9.2

  2. 现在maven工程已生成,用IDEA导入这个工程,如下图: 在这里插入图片描述

  3. 以maven的类型导入: 在这里插入图片描述

  4. 导入成功的样子: 在这里插入图片描述

  5. 项目创建成功,可以开始写代码实战了;

辅助类Splitter

实战中有个功能常用到:将字符串用空格分割,转成Tuple2类型的集合,这里将此算子做成一个公共类Splitter.java,代码如下:

1package com.bolingcavalry; 2 3import org.apache.flink.api.common.functions.FlatMapFunction; 4import org.apache.flink.api.java.tuple.Tuple2; 5import org.apache.flink.util.Collector; 6import org.apache.flink.util.StringUtils; 7 8public class Splitter implements FlatMapFunction<String, Tuple2<String, Integer>> { 9 @Override 10 public void flatMap(String s, Collector<Tuple2<String, Integer>> collector) throws Exception { 11 12 if(StringUtils.isNullOrWhitespaceOnly(s)) { 13 System.out.println("invalid line"); 14 return; 15 } 16 17 for(String word : s.split(" ")) { 18 collector.collect(new Tuple2<String, Integer>(word, 1)); 19 } 20 } 21}

准备完毕,可以开始实战了,先从最简单的Socket开始。

Socket DataSource

Socket DataSource的功能是监听指定IP的指定端口,读取网络数据;

  1. 在刚才新建的工程中创建一个类Socket.java:

    package com.bolingcavalry.api;

    import com.bolingcavalry.Splitter; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.windowing.time.Time;

    public class Socket { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    1 //监听本地9999端口,读取字符串 2 DataStream<String> socketDataStream = env.socketTextStream("localhost", 9999); 3 4 //每五秒钟一次,将当前五秒内所有字符串以空格分割,然后统计单词数量,打印出来 5 socketDataStream 6 .flatMap(new Splitter()) 7 .keyBy(0) 8 .timeWindow(Time.seconds(5)) 9 .sum(1) 10 .print(); 11 12 env.execute("API DataSource demo : socket"); 13}

    }

从上述代码可见,StreamExecutionEnvironment.socketTextStream就可以创建Socket类型的DataSource,在控制台执行命令<font color="blue">nc -lk 9999</font>,即可进入交互模式,此时输出任何字符串再回车,都会将字符串传输到本机9999端口;

  1. 在IDEA上运行Socket类,启动成功后再回到刚才执行<font color="blue">nc -lk 9999</font>的控制台,输入一些字符串再回车,可见Socket的功能已经生效:

在这里插入图片描述

集合DataSource(generateSequence)

  1. 基于集合的DataSource,API如下图所示:

在这里插入图片描述 2. 先试试最简单的generateSequence,创建指定范围内的数字型的DataSource:

1package com.bolingcavalry.api; 2 3import org.apache.flink.api.common.functions.FilterFunction; 4import org.apache.flink.streaming.api.datastream.DataStream; 5import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; 6 7public class GenerateSequence { 8 public static void main(String[] args) throws Exception { 9 final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); 10 11 //并行度为1 12 env.setParallelism(1); 13 14 //通过generateSequence得到Long类型的DataSource 15 DataStream<Long> dataStream = env.generateSequence(1, 10); 16 17 //做一次过滤,只保留偶数,然后打印 18 dataStream.filter(new FilterFunction<Long>() { 19 @Override 20 public boolean filter(Long aLong) throws Exception { 21 return 0L==aLong.longValue()%2L; 22 } 23 }).print(); 24 25 env.execute("API DataSource demo : collection"); 26 } 27}

3. 运行时会打印偶数:

4.

集合DataSource(fromElements+fromCollection)

  1. fromElements和fromCollection就在一个类中试了吧,创建<font color="blue">FromCollection</font>类,里面是这两个API的用法:

    package com.bolingcavalry.api;

    import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    import java.util.ArrayList; import java.util.List;

    public class FromCollection { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

    1 //并行度为1 2 env.setParallelism(1); 3 4 //创建一个List,里面有两个Tuple2元素 5 List<Tuple2<String, Integer>> list = new ArrayList<>(); 6 list.add(new Tuple2("aaa", 1)); 7 list.add(new Tuple2("bbb", 1)); 8 9 //通过List创建DataStream 10 DataStream<Tuple2<String, Integer>> fromCollectionDataStream = env.fromCollection(list); 11 12 //通过多个Tuple2元素创建DataStream 13 DataStream<Tuple2<String, Integer>> fromElementDataStream = env.fromElements( 14 new Tuple2("ccc", 1), 15 new Tuple2("ddd", 1), 16 new Tuple2("aaa", 1) 17 ); 18 19 //通过union将两个DataStream合成一个 20 DataStream<Tuple2<String, Integer>> unionDataStream = fromCollectionDataStream.union(fromElementDataStream); 21 22 //统计每个单词的数量 23 unionDataStream 24 .keyBy(0) 25 .sum(1) 26 .print(); 27 28 env.execute("API DataSource demo : collection"); 29}

    }

  2. 运行结果如下: 在这里插入图片描述

文件DataSource

  1. 下面的ReadTextFile类会读取绝对路径的文本文件,并对内容做单词统计:

    package com.bolingcavalry.api;

    import com.bolingcavalry.Splitter; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

    public class ReadTextFile { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); //设置并行度为1 env.setParallelism(1);

    1 //用txt文件作为数据源 2 DataStream<String> textDataStream = env.readTextFile("file:///Users/zhaoqin/temp/202003/14/README.txt", "UTF-8"); 3 4 //统计单词数量并打印出来 5 textDataStream 6 .flatMap(new Splitter()) 7 .keyBy(0) 8 .sum(1) 9 .print(); 10 11 env.execute("API DataSource demo : readTextFile"); 12}

    }

  2. 请确保代码中的绝对路径下存在名为README.txt文件,运行结果如下:

在这里插入图片描述 3. 打开StreamExecutionEnvironment.java源码,看一下刚才使用的readTextFile方法实现如下,原来是调用了另一个同名方法,该方法的第三个参数确定了文本文件是一次性读取完毕,还是周期性扫描内容变更,而第四个参数就是周期性扫描的间隔时间:

1public DataStreamSource<String> readTextFile(String filePath, String charsetName) { 2 Preconditions.checkArgument(!StringUtils.isNullOrWhitespaceOnly(filePath), "The file path must not be null or blank."); 3 4 TextInputFormat format = new TextInputFormat(new Path(filePath)); 5 format.setFilesFilter(FilePathFilter.createDefaultFilter()); 6 TypeInformation<String> typeInfo = BasicTypeInfo.STRING_TYPE_INFO; 7 format.setCharsetName(charsetName); 8 9 return readFile(format, filePath, FileProcessingMode.PROCESS_ONCE, -1, typeInfo); 10 }

4. 上面的FileProcessingMode是个枚举,源码如下:

1@PublicEvolving 2public enum FileProcessingMode { 3 4 /** Processes the current contents of the path and exits. */ 5 PROCESS_ONCE, 6 7 /** Periodically scans the path for new data. */ 8 PROCESS_CONTINUOUSLY 9}

5. 另外请关注<font color="blue">readTextFile</font>方法的<font color="red">filePath</font>参数,这是个URI类型的字符串,除了本地文件路径,还可以是HDFS的地址:<font color="blue">hdfs://host:port/file/path</font>

至此,通过直接API创建DataSource的实战就完成了,后面的章节我们继续学习内置connector方式的DataSource;

欢迎关注公众号:程序员欣宸

微信搜索「程序员欣宸」,我是欣宸,期待与您一同畅游Java世界... https://github.com/zq2599/blog_demos

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

jackson学习之三:常用API操作

欢迎访问我的GitHubhttps://github.com/zq2599/blog\_demos(https://www.oschina.net/action/GoToLink?urlhttps%3A%2F%2Fgithub.com%2Fzq2599%2Fblog_demos)内容:所有原创文章分类汇总及配套源码,涉及Java、Doc

K8S的StorageClass实战(NFS)

欢迎访问我的GitHubhttps://github.com/zq2599/blog\_demos(https://www.oschina.net/action/GoToLink?urlhttps%3A%2F%2Fgithub.com%2Fzq2599%2Fblog_demos)内容:所有原创文章分类汇总及配套源码,涉及Java、Doc

Flink处理函数实战之五:CoProcessFunction(双流处理)

欢迎访问我的GitHubhttps://github.com/zq2599/blog\_demos(https://www.oschina.net/action/GoToLink?urlhttps%3A%2F%2Fgithub.com%2Fzq2599%2Fblog_demos)内容:所有原创文章分类汇总及配套源码,涉及Java、Doc