Flink 开发一些代码

Coreqi / 2024-01-20 / 原文

本文基于:Flink Java Demo

1.开发中开启WEB UI

1.添加依赖

        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-runtime-web</artifactId>
            <version>1.18.0</version>
            <scope>provided</scope>
        </dependency>

2.创建带WEBUI的执行环境

        // 创建带WEBUI的执行环境,一般用于本地测试,需要引入flink-runtime-web依赖
        StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());

3.访问

http://localhost:8081/

2.设置任务的并行度

在算子方法后使用setParallelism()方法设置并行度
在idea运行,不指定并行度,默认就是本机的线程数

        // 处理数据
        SingleOutputStreamOperator<Tuple2<String, Integer>> wordToOne = socketDS
                .flatMap((String s, Collector<Tuple2<String, Integer>> collector) -> {
                    // 按照空格切分单词
                    String[] words = s.split(" ");
                    for (String word : words) {
                        Tuple2<String, Integer> wordTuple2 = Tuple2.of(word, 1);
                        // 使用 Collector 向下游发送数据
                        collector.collect(wordTuple2);
                    }
                })
                .setParallelism(2)  //设置并行度
                .returns(Types.TUPLE(Types.STRING, Types.INT));

3.全局配置并行度

算子方法中配置的并行度优先级要高于全局配置的并行度

        // 创建带WEBUI的执行环境,一般用于本地测试,需要引入flink-runtime-web依赖
        StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());

        env.setParallelism(7);  // 全局配置并行度

4.提交作业时指定并行度

并行度优先级:
算子指定 > env > 提交时指定 > flink-conf.yaml[parallelim.default]

1.在WEB UI中提交任务时设置并行度[略]

2.在命令行中提交任务时设置并行度

使用-p参数指定即可。

5.