【问题标题】:Apache Flink: What is the difference of setParallelism() and setMaxParallelism()Apache Flink:setParallelism() 和 setMaxParallelism() 有什么区别
【发布时间】:2019-02-06 20:06:25
【问题描述】:

我尝试使用ExecutionConfig.setMaxParallelism() 方法为 Flink 作业设置最大并行度,但它似乎不起作用。

我还修改了标准 WordCount 示例以运行一些测试,似乎setMaxParallelism() 方法对本地环境或独立集群都没有任何影响。

setMaxParallelism() 是如何工作的?

【问题讨论】:

  • 您期望setMaxParallelism() 的哪种行为?您可能会将其与setParallelsim() 混淆吗?
  • 没有 setMaxParallelism 行,我在执行过程中看到 parallelism = 8,这就是 env 的配置方式。使用 setMaxParallelism(4) 行,我希望在执行期间看到并行度 = 4。这是正确的吗?我还在具有 16 个内核的本地环境上进行了测试,并且日志文件显示了 1/16 任务等,有或没有 setMaxParallelism 行。我正在尝试为其他一些工作设置最大平行度的上限,这是正确的理解吗?另一方面,setParallelsim 已按预期工作,但这不是我们想要做的。谢谢

标签: apache-flink flink-streaming


【解决方案1】:

Flink 提供了两种设置:

  • setParallelism(x) 将作业或算子的并行度设置为x,即算子的并行任务数。
  • setMaxParallelism(y) 控制可以分配 keyed state 的最大任务数,即 operator 的最大有效并行度。操作员仍然可以有更多任务,但其中只有y 将分配有键状态并可用于处理。分配密钥状态的单位称为密钥组。

documentation 更详细地解释了这些概念。

【讨论】:

  • 感谢 Fabian 的解释。因此,在幕后,最大并行度限制了所涉及的密钥组的最大数量。是不是因为我运行的第一个测试没有使用 KeyedStream(使用 DataSet),所以 setMaxParallelism 没有效果?
  • setMaxParallelism() 不影响任务数。它只是限制密钥组的数量。由于任务多于关键组并没有什么意义,因此 Flink 会抛出您在答案中发布的异常。
  • @FabianHueske 顺便说一句,setParallelism() 应用于StreamExecutionEnvironment 和单独应用于StreamOperators 之间有什么区别?假设我申请了env.setParallelism(10),但将setParallelism(1) 用于来源和运营商
【解决方案2】:

我今天又进行了一些测试,使用流而不是数据集。这次看到了setMaxParallelism()的效果。

    public static void main(String[] args) throws Exception
    {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.getConfig().setMaxParallelism(4); // <-- effect

        DataStream<String> text = env.fromElements(WORDS);

        DataStream<Tuple2<String, Integer>> counts = text.flatMap(new Tokenizer()).keyBy(0).sum(1);

        counts.writeAsCsv("test.dat");

        env.execute("WordCount Example");
    }

客户看到的有趣错误,

Caused by: org.apache.flink.runtime.JobException: Vertex Flat Map's parallelism (8) is higher than the max parallelism (4). Please lower the parallelism or increase the max parallelism.
        at org.apache.flink.runtime.executiongraph.ExecutionJobVertex.<init>(ExecutionJobVertex.java:188)
        at org.apache.flink.runtime.executiongraph.ExecutionGraph.attachJobGraph(ExecutionGraph.java:830)
        at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilder.buildGraph(ExecutionGraphBuilder.java:232)
        at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilder.buildGraph(ExecutionGraphBuilder.java:100)
        at org.apache.flink.runtime.jobmaster.JobMaster.createExecutionGraph(JobMaster.java:1152)
        at org.apache.flink.runtime.jobmaster.JobMaster.createAndRestoreExecutionGraph(JobMaster.java:1132)
        at org.apache.flink.runtime.jobmaster.JobMaster.<init>(JobMaster.java:294)
        at org.apache.flink.runtime.jobmaster.JobManagerRunner.<init>(JobManagerRunner.java:157)
        ... 10 more

谢谢

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-05-04
    • 2015-08-28
    • 2016-06-23
    • 2016-06-05
    • 2018-11-04
    • 2015-03-21
    • 2016-08-30
    • 1970-01-01
    相关资源
    最近更新 更多