Flink 开发一些代码
本文基于: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参数指定即可。