第一章.Flink核心架构

1.运行架构

day02[Flink核心架构_流处理核心编程] - 图1

Flink运行时包含2种进程:1个JobManager和至少一个TaskManager

  1. # 客户端
  2. 严格上说, 客户端不是运行和程序执行的一部分, 而是用于准备和发送dataflow到JobManager. 然后客户端可以断开与JobManager的连接(detached mode), 也可以继续保持与JobManager的连接(attached mode)
  3. 客户端作为触发执行的java或者scala代码的一部分运行, 也可以在命令行运行:bin/flink run ...
  1. # JobManager
  2. 控制一个应用程序执行的主进程,也就是说,每个应用程序都会被一个的JobManager所控制执行。
  3. JobManager会先接收到要执行的应用程序,这个应用程序会包括:作业图(JobGraph)、逻辑数据流图(logical dataflow graph)和打包了所有的类、库和其它资源的JAR包。
  4. JobManager会把JobGraph转换成一个物理层面的数据流图,这个图被叫做“执行图”(ExecutionGraph),包含了所有可以并发执行的任务。JobManager会向资源管理器(ResourceManager)请求执行任务必要的资源,也就是任务管理器(TaskManager)上的插槽(slot)。一旦它获取到了足够的资源,就会将执行图分发到真正运行它们的TaskManager上。
  5. 而在运行过程中,JobManager会负责所有需要中央协调的操作,比如说检查点(checkpoints)的协调。
  6. JobManager包含ResourceManager,Dispatcher和JobMaster
  1. # ResourceManager
  2. 负责资源的管理,在整个 Flink 集群中只有一个 ResourceManager. 注意这个ResourceManager不是Yarn中的ResourceManager, 是Flink中内置的, 只是赶巧重名了而已.
  3. 主要负责管理任务管理器(TaskManager)的插槽(slot),TaskManger插槽是Flink中定义的处理资源单元。
  4. 当JobManager申请插槽资源时,ResourceManager会将有空闲插槽的TaskManager分配给JobManager。如果ResourceManager没有足够的插槽来满足JobManager的请求,它还可以向资源提供平台发起会话,以提供启动TaskManager进程的容器。另外,ResourceManager还负责终止空闲的TaskManager,释放计算资源。
  1. # Dispatcher
  2. 负责接收用户提供的作业,并且负责为这个新提交的作业启动一个新的JobMaster 组件. Dispatcher也会启动一个Web UI,用来方便地展示和监控作业执行的信息。Dispatcher在架构中可能并不是必需的,这取决于应用提交运行的方式。
  1. # JobMaster
  2. JobMaster负责管理单个JobGraph的执行.多个Job可以同时运行在一个Flink集群中, 每个Job都有一个自己的JobMaster
  1. # TaskManager
  2. Flink中的工作进程。通常在Flink中会有多个TaskManager运行,每一个TaskManager都包含了一定数量的插槽(slots)。插槽的数量限制了TaskManager能够执行的任务数量。
  3. 启动之后,TaskManager会向资源管理器注册它的插槽;收到资源管理器的指令后,TaskManager就会将一个或者多个插槽提供给JobManager调用。JobManager就可以向插槽分配任务(tasks)来执行了。
  4. 在执行过程中,一个TaskManager可以跟其它运行同一应用程序的TaskManager交换数据。

2.核心概念

一.TaskManager与Slots

day02[Flink核心架构_流处理核心编程] - 图2

  1. Flink中每一个worker(TaskManager)都是一个JVM进程,它可能会在独立的线程上执行一个Task。为了控制一个worker能接收多少个task,worker通过Task Slot来进行控制(一个worker至少有一个Task Slot)。
  2. 这里的Slot如何来理解呢?很多的文章中经常会和Spark框架进行类比,将Slot类比为Core,其实简单这么类比是可以的,可实际上,可以考虑下,当Spark申请资源后,这个Core执行任务时有可能是空闲的,但是这个时候Spark并不能将这个空闲下来的Core共享给其他Job使用,所以这里的Core是Job内部共享使用的。接下来我们再回想一下,之前在Yarn Session-Cluster模式时,其实是可以并行执行多个Job的,那如果申请两个Slot,而执行Job时,只用到了一个,剩下的一个怎么办?那我们自认而然就会想到可以将这个Slot给并行的其他Job,对吗?所以Flink中的Slot和Spark中的Core还是有很大区别的。
  3. 每个task slot表示TaskManager拥有资源的一个固定大小的子集。假如一个TaskManager有三个slot,那么它会将其管理的内存分成三份给各个slot。资源slot化意味着一个task将不需要跟来自其他job的task竞争被管理的内存,取而代之的是它将拥有一定数量的内存储备。需要注意的是,这里不会涉及到CPU的隔离,slot目前仅仅用来隔离task的受管理的内存。

二.Parallelism(并行度)

day02[Flink核心架构_流处理核心编程] - 图3

  1. 一个特定算子的子任务(subtask)的个数被称之为这个算子的并行度(parallelism),一般情况下,一个流程序的并行度,可以认为就是其所有算子中最大的并行度。一个程序中,不同的算子可能具有不同的并行度。
  2. Stream在算子之间传输数据的形式可以是one-to-one(forwarding)的模式也可以是redistributing的模式,具体是哪一种形式,取决于算子的种类。
  1. 1. One to One
  2. stream(比如在source和map operator之间)维护着分区以及元素的顺序。那意味着flatmap 算子的子任务看到的元素的个数以及顺序跟source 算子的子任务生产的元素的个数、顺序相同,map、fliter、flatMap等算子都是one-to-one的对应关系。类似于spark中的窄依赖
  3. 2. Redistributing
  4. stream(map()跟keyBy/window之间或者keyBy/window跟sink之间)的分区会发生改变。每一个算子的子任务依据所选择的transformation发送数据到不同的目标任务。例如,keyBy()基于hashCode重分区、broadcast和rebalance会随机重新分区,这些算子都会引起redistribute过程,而redistribute过程就类似于Spark中的shuffle过程。类似于spark中的宽依赖

三.Task与SubTask

  1. 一个算子就是一个Task,一个算子的并行度是几,这个Task就有几个SubTask

四.Operator Chains(任务链)

  1. 相同并行度的one to one操作,Flink将这样相连的算子链接在一起形成一个task,原来的算子成为里面的一部分。 每个subtask被一个线程执行.
  2. 将算子链接成task是非常有效的优化:它能减少线程之间的切换和基于缓存区的数据交换,在减少时延的同时提升吞吐量。链接的行为可以在编程API中进行指定。

五.ExecutionGraph(执行图)

  1. 由Flink程序直接映射成的数据流图是StreamGraph,也被称为逻辑流图,因为它们表示的是计算逻辑的高级视图。为了执行一个流处理程序,Flink需要将逻辑流图转换为物理数据流图(也叫执行图),详细说明程序的执行方式。
  2. Flink 中的执行图可以分成四层:StreamGraph -> JobGraph -> ExecutionGraph -> Physical Graph。
  1. 1. StreamGraph
  2. 是根据用户通过 Stream API 编写的代码生成的最初的图。用来表示程序的拓扑结构。
  3. 2. JobGraph
  4. StreamGraph经过优化后生成了 JobGraph,是提交给 JobManager 的数据结构。主要的优化为: 将多个符合条件的节点 chain 在一起作为一个节点,这样可以减少数据在节点之间流动所需要的序列化/反序列化/传输消耗。
  5. 3. ExecutionGraph
  6. JobManager 根据 JobGraph 生成ExecutionGraph。ExecutionGraph是JobGraph的并行化版本,是调度层最核心的数据结构。
  7. 4. Physical Graph
  8. JobManager 根据 ExecutionGraph 对 Job 进行调度后,在各个TaskManager 上部署 Task 后形成的“图”,并不是一个具体的数据结构。
  9. 2个并发度(Source为1个并发度)的 SocketTextStreamWordCount 四层执行图的演变过程
  10. env.socketTextStream().flatMap(…).keyBy(0).sum(1).print()

3.提交流程

一.通用提交流程

day02[Flink核心架构_流处理核心编程] - 图4

二.yarn-cluster提交流程per-job

day02[Flink核心架构_流处理核心编程] - 图5

  1. Flink任务提交后,Client向HDFS上传Flink的Jar包和配置
  2. 向Yarn ResourceManager提交任务,ResourceManager分配Container资源
  3. 通知对应的NodeManager启动ApplicationMaster,ApplicationMaster启动后加载Flink的 Jar包和配置构建环境,然后启动JobManager
  4. ApplicationMaster向ResourceManager申请资源启动TaskManager
  5. ResourceManager分配Container资源后,由ApplicationMaster通知资源所在节点的

    NodeManager启动TaskManager

  6. NodeManager加载Flink的Jar包和配置构建环境并启动TaskManager
  7. TaskManager启动后向JobManager发送心跳包,并等待JobManager向其分配任务。

flink.pdf


第二章.Flink流处理核心编程

1.准备工作

一.pom.xml

  1. <!-- https://mvnrepository.com/artifact/org.projectlombok/lombok -->
  2. <dependency>
  3. <groupId>org.projectlombok</groupId>
  4. <artifactId>lombok</artifactId>
  5. <version>1.18.16</version>
  6. <scope>provided</scope>
  7. </dependency>
  8. <dependency>
  9. <groupId>org.apache.hadoop</groupId>
  10. <artifactId>hadoop-client</artifactId>
  11. <version>3.1.3</version>
  12. <scope>provided</scope>
  13. </dependency>
  14. <dependency>
  15. <groupId>org.apache.flink</groupId>
  16. <artifactId>flink-connector-kafka_2.12</artifactId>
  17. <version>1.13.1</version>
  18. </dependency>

二.提供一个JavaBean(WaterSensor)类方便展示

  1. package com.atguigu.flink.day02.pojo;
  2. import lombok.AllArgsConstructor;
  3. import lombok.Data;
  4. import lombok.NoArgsConstructor;
  5. /**
  6. * 水位传感器:用于接收水位数据
  7. *
  8. * id:传感器编号
  9. * ts:时间戳
  10. * vc:水位
  11. */
  12. @Data
  13. @NoArgsConstructor
  14. @AllArgsConstructor
  15. public class WaterSensor {
  16. private String id;
  17. private Long ts;
  18. private Integer vc;
  19. }

2.Source

  1. Flink框架可以从不同的来源获取数据,将数据提交给框架进行处理,我们将获取数据的来源称之为数据源(Source)

一.从Java的集合中读取数据

  1. package com.atguigu.flink.day02.source;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. import java.util.Arrays;
  5. import java.util.List;
  6. /**
  7. * 从java的集合中读取数据
  8. */
  9. public class Flink_Source_Collection {
  10. public static void main(String[] args) throws Exception {
  11. List<WaterSensor> waterSensors = Arrays.asList(
  12. new WaterSensor("ws001", 1111111L, 45),
  13. new WaterSensor("ws002", 2222222L, 45),
  14. new WaterSensor("ws003", 3333333L, 45)
  15. );
  16. //创建执行环境
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(1);
  19. env.fromCollection(waterSensors)
  20. .print();
  21. env.fromElements(
  22. new WaterSensor("ws001", 1111111L, 45),
  23. new WaterSensor("ws002", 2222222L, 45),
  24. new WaterSensor("ws003", 3333333L, 45)
  25. ).print();
  26. env.execute();
  27. }
  28. }

二.从文件读取数据

  1. package com.atguigu.flink.day02.source;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. /**
  4. * 从文件中读取数据
  5. */
  6. public class Flink_Source_File {
  7. public static void main(String[] args) throws Exception {
  8. //创建执行环境
  9. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  10. env.setParallelism(1);
  11. env.readTextFile("hdfs://hadoop102:9820/flink/data/word.txt").print();
  12. env.execute();
  13. }
  14. }

三.从kafka中读取数据

  1. package com.atguigu.flink.day02.source;
  2. import org.apache.flink.api.common.serialization.SimpleStringSchema;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
  5. import org.apache.flink.streaming.util.serialization.JSONKeyValueDeserializationSchema;
  6. import java.util.Properties;
  7. //{"id":1001,"name":"zhangsan"}
  8. public class Flink_Source_Kafka {
  9. public static void main(String[] args) throws Exception {
  10. //获取执行环境
  11. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  12. env.setParallelism(1);
  13. Properties properties = new Properties();
  14. properties.setProperty("bootstrap.servers","hadoop102:9092,hadoop103:9092,hadoop104:9092");
  15. properties.setProperty("group.id","Flink_Source_Kafka");
  16. properties.setProperty("auto.offset.reset","latest");
  17. //env.addSource(new FlinkKafkaConsumer<>("first",new SimpleStringSchema(),properties)).print();
  18. env.addSource(new FlinkKafkaConsumer<>("first",new JSONKeyValueDeserializationSchema(false),properties)).print();
  19. env.execute();
  20. }
  21. }

四.自定义Source

  1. package com.atguigu.flink.day02.source;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. import org.apache.flink.streaming.api.functions.source.SourceFunction;
  5. import java.io.BufferedReader;
  6. import java.io.InputStream;
  7. import java.io.InputStreamReader;
  8. import java.net.Socket;
  9. import java.nio.charset.StandardCharsets;
  10. /**
  11. * 自定义source
  12. * //sensor1 1607527992000 20
  13. */
  14. public class Flink_Source_Custom {
  15. public static void main(String[] args) throws Exception {
  16. //获取执行环境
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.addSource(new MySocketSource("hadoop102",9999)).print();
  19. env.execute();
  20. }
  21. static class MySocketSource implements SourceFunction<WaterSensor>{
  22. String host;
  23. Integer port;
  24. boolean flag = true;
  25. public MySocketSource(String host, Integer port) {
  26. this.host = host;
  27. this.port = port;
  28. }
  29. /**
  30. * 数据生产的代码,ctx就是用来放我们的数据的
  31. * @param ctx
  32. * @throws Exception
  33. */
  34. @Override
  35. public void run(SourceContext<WaterSensor> ctx) throws Exception {
  36. Socket socket = new Socket(host,port);
  37. InputStream inputStream = socket.getInputStream();
  38. InputStreamReader inputStreamReader = new InputStreamReader(inputStream, StandardCharsets.UTF_8);
  39. BufferedReader bufferedReader = new BufferedReader(inputStreamReader);
  40. String line = bufferedReader.readLine();
  41. while(line != null && flag){
  42. String[] data = line.split(" ");
  43. ctx.collect(new WaterSensor(data[0],Long.valueOf(data[1]),Integer.valueOf(data[2])));
  44. line = bufferedReader.readLine();
  45. }
  46. }
  47. /**
  48. * 程序结束的时候进行调用,或者是进行外部调用来停止整个流
  49. */
  50. @Override
  51. public void cancel() {
  52. flag = false;
  53. }
  54. }
  55. }

3.Transform

转换算子可以把一个或多个DataStream转成一个新的DataStream,程序可以把多个复杂的转换组合成复杂的数据流拓扑

一.map

作用:将数据流中的数据进行转换,形成新的数据流,消费一个元素并产生一个元素

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. /**
  5. * 使用map算子将数字转换为数字的平方
  6. */
  7. public class $01_MapFunction {
  8. public static void main(String[] args) throws Exception {
  9. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  10. env.setParallelism(2);
  11. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4);
  12. ds.map(ele -> ele * ele).print();
  13. env.execute();
  14. }
  15. }

二.RichMap

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.api.common.functions.RichMapFunction;
  3. import org.apache.flink.configuration.Configuration;
  4. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. /**
  7. * 所有Flink函数类都有其Rich版本,它与常规函数的不同在于,可以获取运行环境的上下文,并拥有一些生命周期方法,所以可以实现更复杂的功能,也
  8. * 就意味着提供了更多的,丰富的功能,例如:RichMapFunction
  9. */
  10. public class $02_RichMapFunction {
  11. public static void main(String[] args) throws Exception {
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(2);
  14. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4);
  15. ds.map(new RichMapFunction<Integer, Integer>() {
  16. @Override
  17. public void open(Configuration parameters) throws Exception {
  18. /**
  19. * 一般可以在open方法中,创建一个数据库的连接
  20. * 算子的并行度为几,那么open方法就会执行几次,close与open类似
  21. */
  22. System.out.println("$02_RichMapFunction.open");
  23. }
  24. @Override
  25. public void close() throws Exception {
  26. //一般可以在close方法中,可以归还一个数据库连接
  27. System.out.println("$02_RichMapFunction.close");
  28. }
  29. @Override
  30. public Integer map(Integer value) throws Exception {
  31. System.out.println("$02_RichMapFunction.map");
  32. return value * value;
  33. }
  34. }).print();
  35. env.execute();
  36. }
  37. }

三.FlatMap

消费一个元素并产生零个或多个元素

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.api.common.typeinfo.Types;
  3. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.util.Collector;
  7. /**
  8. * 需求:使用flatmap算子将数字转换为数字,数字的平方,数字的立方并重新放回流中
  9. */
  10. public class $03_FlatMapFunction {
  11. public static void main(String[] args) throws Exception {
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(2);
  14. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4);
  15. SingleOutputStreamOperator<Integer> results = ds.flatMap((Integer ele, Collector<Integer> out) -> {
  16. out.collect(ele);
  17. out.collect(ele * ele);
  18. out.collect(ele * ele * ele);
  19. }
  20. ).returns(Types.INT);
  21. results.print();
  22. env.execute();
  23. }
  24. }

四.Filter

根据指定的规则将满足条件(true)的数据保留,不满足条件(false)的数据丢弃

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. /**
  5. * 需求:使用Filter算子将流中的全部偶数取出
  6. */
  7. public class $04_FilterFunction {
  8. public static void main(String[] args) throws Exception {
  9. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  10. env.setParallelism(2);
  11. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4);
  12. ds.filter(ele -> ele % 2 ==0).print();
  13. env.execute();
  14. }
  15. }

五.KeyBy

将流中的数据分到不同的分区(并行度)中,具有相同key的元素会分到同一个分区中,一个分区中可以有多个不同的key,在内部是使用key的hash分区来实现的

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. public class $05_KeyByFunction {
  5. public static void main(String[] args) throws Exception {
  6. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  7. env.setParallelism(2);
  8. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4, 5, 6, 7, 8, 9);
  9. ds.keyBy(ele -> {
  10. if(ele % 2 ==0){
  11. return "偶数";
  12. }else{
  13. return "奇数";
  14. }
  15. }).print();
  16. env.execute();
  17. }
  18. }

六.shuffle

将流中的元素随机打乱,对同一组数据,每次执行得到的结果都不同

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. /**
  5. * shuffle算子随机打散流中元素
  6. */
  7. public class $06_ShuffleFunction {
  8. public static void main(String[] args) throws Exception {
  9. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  10. env.setParallelism(2);
  11. DataStreamSource<Integer> ds = env.fromElements(1, 2, 3, 4, 5, 6, 7, 8, 9);
  12. ds.shuffle().print();
  13. env.execute();
  14. }
  15. }

七.connect

在某些情况下,我们需要将两个不同来源的数据流进行连接,实现数据匹配,比如订单支付和第三方交易信息,这两个信息的数据就来自于不同数据源,连接后,将订单支付和第三方交易信息进行对账,此时,才能算真正的支付完成

Flink中的connect算子可以连接两个保持他们类型的数据流,两个数据流被connect之后,只是被放在了同一个流中,内部依然保持各自的数据和形式不发生任何变化,两个流相互独立

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.ConnectedStreams;
  3. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  4. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  5. import org.apache.flink.streaming.api.functions.co.CoMapFunction;
  6. /**
  7. * 需求:连接两个不同的流
  8. * 两个流存储的数据类型可以不同
  9. * 只是机械的合并在一起,内部仍然是分离的2个流
  10. * 只能2个流进行connect,不能有第三个参与
  11. */
  12. public class $07_ConnectFunction {
  13. public static void main(String[] args) throws Exception {
  14. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  15. env.setParallelism(2);
  16. DataStreamSource<Integer> ds1 = env.fromElements(1, 2, 3, 4);
  17. DataStreamSource<String> ds2 = env.fromElements("a", "b", "c");
  18. ConnectedStreams<Integer, String> connect = ds1.connect(ds2);
  19. connect.map(new CoMapFunction<Integer, String, String>() {
  20. @Override
  21. public String map1(Integer value) throws Exception {
  22. return value + ":是数字";
  23. }
  24. @Override
  25. public String map2(String value) throws Exception {
  26. return value + ":是字母";
  27. }
  28. }).print();
  29. env.execute();
  30. }
  31. }

八.union

  1. package com.atguigu.flink.day02.transform;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. /**
  5. * connect与union的区别
  6. * Union之前两个或多个流的类型必须是一样的,connect可以不一样
  7. * connect只能操作两个流,union可以操作多个
  8. */
  9. public class $08_UnionFunction {
  10. public static void main(String[] args) throws Exception {
  11. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  12. env.setParallelism(2);
  13. DataStreamSource<Integer> ds1 = env.fromElements(1, 2, 3, 4);
  14. DataStreamSource<Integer> ds2 = env.fromElements(5, 6, 7, 8);
  15. DataStreamSource<Integer> ds3 = env.fromElements(9, 10, 11, 12);
  16. DataStreamSource<String> ds4 = env.fromElements("a", "b", "c");
  17. ds1.union(ds2).union(ds3).print();
  18. env.execute();
  19. }
  20. }