Flink DataStream API 总结

Coreqi / 2024-01-21 / 原文

流程如下 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");