YuF*_*hen 8 apache-flink flink-streaming
我只是得到下面关于并行性的示例,并且有一些相关的问题:
setParallelism(5)将Parallelism 5设置为求和或flatMap和求和?
是否可以分别为flatMap和sum等不同的运算符设置不同的Parallelism?例如将Parallelism 5设置为sum和10设置为flatMap。
根据我的理解,keyBy正在根据不同的密钥将DataStream划分为逻辑Stream \分区,并假设有10,000个不同的键值,因此有10,000个不同的分区,那么有多少个线程可以处理10,000个分区?只有5个线程?如果不设置setParallelism(5)怎么办?
https://ci.apache.org/projects/flink/flink-docs-release-1.3/dev/parallel.html
final StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = [...]
DataStream<Tuple2<String, Integer>> wordCounts = text
.flatMap(new LineSplitter())
.keyBy(0)
.timeWindow(Time.seconds(5))
.sum(1).setParallelism(5);
wordCounts.print();
env.execute("Word Count Example");
Run Code Online (Sandbox Code Playgroud)
调用setParallelism运算符时,它将更改此特定运算符的并行性。因此,在您的示例中,仅窗口运算符将以的并行度执行,5而前一个flatMap运算符将以默认的并行度执行。
因此,可以为每个运算符设置不同的并行度。但是,请注意,不能将具有不同并行度的运算符链接在一起,并且需要进行重新平衡(类似于随机播放)操作。
如果要为所有运算符设置并行性,则必须通过ExecutionEnvironment#setParallelismAPI调用来实现。
keyBy输入流中的操作分区将与您具有并行运算符实例的分区一样多。这样可以确保所有具有相同键的元素都位于同一分区中。因此,在将并行度设置为的示例中5,最终将有5个分区。每个分区可以包含具有不同键的元素。
| 归档时间: |
|
| 查看次数: |
3113 次 |
| 最近记录: |