Flink DataStream API 总结
流程如下 Environment(执行环境) -> Source(数据源) -> Transformation(转换操作) -> Sink(输出)
1.Environment(执行环境)
package cn.coreqi.env;
import org.apache.flink.api.common.JobExecutionResult;
import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.RestOptions;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class EnvDemo {
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration(); //conf对象可以修改一些配置参数
conf.set(RestOptions.BIND_PORT,"8082"); // 默认端口8081
StreamExecutionEnvironment env = StreamExecutionEnvironment
//.getExecutionEnvironment(); // 自动识别 是远程集群 还是本地环境
//.createLocalEnvironment(); // 本地环境
//.createRemoteEnvironment("192.168.58.130",8081,"/home/FlinkTutorial-1.0-SNAPSHOT.jar"); // 远程集群环境
.getExecutionEnvironment(conf);
// 废除了DataSet API - ExecutionEnvironment 同一套代码通过配置来决定是以批还是流的方式处理,默认是以流的方式处理
// 可以在提交时通过参数 -Dexecution.runtime-mode=BATCH | STREAMING 来定义是以批还是流的方式处理
env.setRuntimeMode(RuntimeExecutionMode.BATCH); //批的方式处理
JobExecutionResult execute = env.execute(); //触发程序执行
// 1.默认 env.execute() 触发一个flink job,一个main方法中可以调用多个execute(),但是没意义,在第一个execute()会阻塞程序的运行
// 2.env.executeAsync(),异步触发,不阻塞,每个 executeAsync()都会生成一个flink job
//env.executeAsync();
}
}
2.Source(数据源)
Flink 从各种来源获取数据,然后构建DataStream进行转换处理。
1.分类
1.Version < Flink1.12
在 Flink1.12 以前,旧的添加 source 的方式,是调用执行环境的 addSource()方法:
DataStream<String> stream = env.addSource(...);
方法传入的参数是一个“源函数”(source function),需要实现 SourceFunction 接口。
2.Version >= Flink1.12
从 Flink1.12 开始,主要使用流批统一的新 Source 架构:
DataStreamSource<String> stream = env.fromSource(…);
Flink 直接提供了很多预实现的接口,此外还有很多外部连接工具也帮我们实现了对应的Source,通常情况下足以应对我们的实际需求。
2.实操
1.从集合中读取数据
最简单的读取数据的方式,就是在代码中直接创建一个 Java 集合,然后调用执行环境的fromCollection 方法进行读取。这相当于将数据临时存储到内存中,形成特殊的数据结构后,作为数据源使用,一般用于测试。
List<Integer> data = Arrays.asList(1, 3, 5, 7, 9);
DataStreamSource<Integer> ds = env.fromCollection(data);
ds.print();
DataStreamSource<Integer> ds = env.fromElements(1,3,5,7,9);
ds.print();
2.从文件中读取
从存储介质中获取数据,一个比较常见的方式就是读取日志文件。这也是批处理中最常见的读取方式。
1.旧的添加 source 的方式 - addSource()
// 读取数据
DataSource<String> lineDS = env.readTextFile("input/word.txt");
2.新的
1.POM添加依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-files</artifactId>
<version>1.18.0</version>
</dependency>
2.代码
FileSource<String> fileSource = FileSource.forRecordStreamFormat(
new TextLineInputFormat(),
new Path("input/word.txt"))
.build();
DataStreamSource<String> source = env.fromSource(fileSource, WatermarkStrategy.noWatermarks(), "fileSource");
source.print();
3.从Socket中读取
集合还是文件,读取的其实都是有界数据。在流处理的场景中,数据往往是
无界的。
读取 socket 文本流,就是流处理场景。但是这种方式由于吞吐量小、稳定性较差,一般也是用于测试。
// 读取socket数据
DataStreamSource<String> socketDS = env.socketTextStream("localhost", 7878);
4.从 Kafka 读取数据
Flink 官方提供了连接工具 flink-connector-kafka,直接帮我们实现了一个消费者 FlinkKafkaConsumer,它就是用来读取 Kafka 数据的 SourceFunction。
想要以 Kafka 作为数据源获取数据,只需要引入 Kafka 连接器的依赖。Flink 官方提供的是一个通用的 Kafka 连接器,它会自动跟踪最新版本的 Kafka 客户端。目前最新版本只支持 0.10.0 版本以上的 Kafka。
1.POM添加依赖'
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>3.0.2-1.18</version>
</dependency>
2.代码
KafkaSource<String> kafkaSource =
KafkaSource.<String>builder()
.setBootstrapServers("192.168.58.130:9092,192.168.58.131:9092,192.168.58.132:9092") //指定kafka节点的地址和端口
.setGroupId("coreqi") //指定消费者组的ID
.setTopics("topic_1") // 指定消费者的Topic
.setValueOnlyDeserializer(new SimpleStringSchema()) //指定value的反序列化器
//kafka消费者的参数
// auto.reset.offsets
// earliest: 如果有offset,从offset继续消费;如果没有offset,从最早消费
// latest: 如果有offset,从offset继续消费;如果没有offset, 从最新消费
// flink的kafkasource,offset消费策略;offsetInitalize,默认是earliest
//ealiest: 一定从 最早 消费
// latest: 一定从 最新 消费
.setStartingOffsets(OffsetsInitializer.latest()) //offset初始化器,flink消费kafka的策略
.build();
DataStreamSource<String> stream = env.fromSource(kafkaSource,
WatermarkStrategy.noWatermarks(), "kafkaSource");
stream.print("Kafka");