十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink 并行执行机制详解:从 setParallelism 到最大并行度与 Key-Group 状态约束

Flink 并行执行机制详解:从 setParallelism 到最大并行度与 Key-Group 状态约束 Flink 并行执行机制详解从 setParallelism 到最大并行度与 Key-Group 状态约束【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flinkFlink 作业的吞吐量与容错能力在很大程度上取决于并行度的配置。本文基于官方文档docs/content.zh/docs/dev/datastream/execution/parallel.md系统讲解 Flink 中并行度与最大并行度的四个配置层次算子、执行环境、客户端、系统并结合flink-core与flink-runtime模块源码剖析默认最大并行度的计算公式以及 key-group 重新分配的实现机制帮助你正确配置作业并行度、避免因状态不兼容导致的 savepoint 恢复失败。1. 核心概念task、并行实例与并行度一个 Flink 程序由多个 task 组成task 可以是转换算子transformation、数据源source或数据接收器sink。每个 task 又包含多个并行执行的实例每个实例只处理该 task 输入数据的一个子集。某个 task 的并行实例数就被称为该 task 的并行度parallelism。理解这个概念的另一个关键点是并行度不是全局唯一的一个数而是每个算子每个 Transformation各自可以单独设定的。这就是 Flink 支持按算子调优的基础——例如让 source 以高并行度拉取数据、让下游 CPU 密集的算子降低并行度。为什么需要最大并行度使用 Savepoint 时应该考虑设置最大并行度max parallelism作业从 savepoint 恢复时你可以改变特定算子乃至整个程序的并行度而最大并行度构成了这个变更的上限Flink 内部将状态划分为 key-groupkey-group 是状态重新分配的最小单元出于性能考虑key-group 的数目不能无限增加因此需要在建作业时就确定一个足够大、又不至于过大的最大并行度它决定了状态切分的粒度。2. 设置并行度的四个层次一个 task 的并行度可以从多个层次指定各层次之间存在一定的覆盖关系。2.1 算子层次setParallelism()单个算子、数据源和数据接收器的并行度可以通过调用setParallelism()方法指定这是最精细的控制粒度。Java 示例final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamString text [...]; DataStreamTuple2String, Integer wordCounts text .flatMap(new LineSplitter()) .keyBy(value - value.f0) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .sum(1).setParallelism(5); wordCounts.print(); env.execute(Word Count Example);Scala 等价写法val env StreamExecutionEnvironment.getExecutionEnvironment val text [...] val wordCounts text .flatMap{ _.split( ) map { (_, 1) } } .keyBy(_._1) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .sum(1).setParallelism(5) wordCounts.print() env.execute(Word Count Example)Python 版本env StreamExecutionEnvironment.get_execution_environment() text [...] word_counts text .flat_map(lambda x: x.split( )) \ .map(lambda i: (i, 1), output_typeTypes.TUPLE([Types.STRING(), Types.INT()])) \ .key_by(lambda i: i[0]) \ .window(TumblingEventTimeWindows.of(Time.seconds(5))) \ .reduce(lambda i, j: (i[0], i[1] j[1])) \ .set_parallelism(5) word_counts.print() env.execute(Word Count Example)从源码结构看DataStream.setParallelism()最终会落到 DAG 层面的 Transformation 上public void setParallelism(int parallelism) { setParallelism(parallelism, true); } public void setParallelism(int parallelism, boolean parallelismConfigured) { OperatorValidationUtils.validateParallelism(parallelism); this.parallelism parallelism; this.parallelismConfigured parallelismConfigured parallelism ! ExecutionConfig.PARALLELISM_DEFAULT; }这里有两点值得注意其一传入值会先经过OperatorValidationUtils.validateParallelism校验其二parallelismConfigured标志用于区分用户显式配置与继承默认值当传入值等于ExecutionConfig.PARALLELISM_DEFAULT即未配置态时不会标记为已显式配置这保证了显式算子并行度在作业图生成时能正确覆盖执行环境默认值。2.2 执行环境层次env.setParallelism()如 Flink 程序结构文档 所述Flink 程序运行在执行环境的上下文中。执行环境为所有算子、数据源、数据接收器定义了一个默认并行度算子层次的显式配置会覆盖这个默认值。如果希望以并行度3执行所有算子可以在执行环境上设置默认并行度final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); DataStreamString text [...]; DataStreamTuple2String, Integer wordCounts [...]; wordCounts.print(); env.execute(Word Count Example);val env StreamExecutionEnvironment.getExecutionEnvironment env.setParallelism(3) val text [...] val wordCounts text .flatMap{ _.split( ) map { (_, 1) } } .keyBy(_._1) .window(TumblingEventTimeWindows.of(Time.seconds(5))) .sum(1) wordCounts.print() env.execute(Word Count Example)env StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(3) text [...] word_counts text .flat_map(lambda x: x.split( )) \ .map(lambda i: (i, 1), output_typeTypes.TUPLE([Types.STRING(), Types.INT()])) \ .key_by(lambda i: i[0]) \ .window(TumblingEventTimeWindows.of(Time.seconds(5))) \ .reduce(lambda i, j: (i[0], i[1] j[1])) word_counts.print() env.execute(Word Count Example)2.3 客户端层次提交作业时指定将作业提交到 Flink 时可在客户端设定并行度。客户端可以是 Java 或 Scala 程序Flink 的命令行接口CLI就是一种典型的客户端。CLI 方式通过-p参数指定并行度例如./bin/flink run -p 10 ../examples/*WordCount-java*.jarJava 程序方式try { PackagedProgram program new PackagedProgram(file, args); InetSocketAddress jobManagerAddress RemoteExecutor.getInetFromHostport(localhost:6123); Configuration config new Configuration(); Client client new Client(jobManagerAddress, config, program.getUserCodeClassLoader()); // set the parallelism to 10 here client.run(program, 10, true); } catch (ProgramInvocationException e) { e.printStackTrace(); }Scala 等价写法try { PackagedProgram program new PackagedProgram(file, args) InetSocketAddress jobManagerAddress RemoteExecutor.getInetFromHostport(localhost:6123) Configuration config new Configuration() Client client new Client(jobManagerAddress, new Configuration(), program.getUserCodeClassLoader()) // set the parallelism to 10 here client.run(program, 10, true) } catch { case e: Exception e.printStackTrace }需要注意Python API 中尚不支持在客户端层次设置并行度这一特性。2.4 系统层次parallelism.default 配置项可以通过设置 Flink 配置文件 中的parallelism.default参数在系统层次指定所有执行环境的默认并行度。从源码可以确认该配置项的定义与默认值见 CoreOptionspublic static final ConfigOptionInteger DEFAULT_PARALLELISM ConfigOptions.key(parallelism.default) .intType() .defaultValue(1) .withDescription(Default parallelism for jobs.);即parallelism.default的默认值为1。各层次并行度可以理解为一种默认值链算子显式配置 执行环境显式配置 系统配置parallelism.default。在没有任何显式配置且未设置系统参数时作业各算子的并行度为 1。3. 设置最大并行度setMaxParallelism() 与计算公式3.1 设置方式与默认值公式最大并行度可以在所有能设置并行度的地方进行设定客户端层次和系统层次除外。与调用setParallelism()修改并行度相似可以通过调用setMaxParallelism()方法设定最大并行度。文档给出的默认最大并行度计算规则是将operatorParallelism (operatorParallelism / 2)向上取整到大于等于该值的 2 的幂次方且下限为128、上限为32768。这一规则在运行时源码中有精确对应见 KeyGroupRangeAssignmentpublic static int computeDefaultMaxParallelism(int operatorParallelism) { checkParallelismPreconditions(operatorParallelism); return Math.min( Math.max( MathUtils.roundUpToPowerOfTwo( operatorParallelism (operatorParallelism / 2)), DEFAULT_LOWER_BOUND_MAX_PARALLELISM), UPPER_BOUND_MAX_PARALLELISM); }逻辑与文档完全一致operatorParallelism operatorParallelism / 2在当前并行度基础上预留 50% 的扩容空间MathUtils.roundUpToPowerOfTwo(...)向上取整到 2 的幂次方便于 key-group 均分Math.max(..., DEFAULT_LOWER_BOUND_MAX_PARALLELISM)下限钳制为 128Math.min(..., UPPER_BOUND_MAX_PARALLELISM)上限钳制为 32768。上限常量在 Transformation 中定义注释明确要求其与流图生成侧的取值保持一致// Has to be equal to StreamGraphGenerator.UPPER_BOUND_MAX_PARALLELISM public static final int UPPER_BOUND_MAX_PARALLELISM 1 15; // 即 32768并且 KeyGroupRangeAssignment 直接复用了同一个常量public static final int UPPER_BOUND_MAX_PARALLELISM Transformation.UPPER_BOUND_MAX_PARALLELISM;按公式推算几个典型值当前并行度为 1~85 时结果落入下限 128并行度为 100 时100 50 150向上取整为 256并行度为 20000 时20000 10000 30000向上取整为 32768。也就是说只有在并行度足够大时默认最大并行度才会逼近上限 32768。setMaxParallelism()在 Transformation 中同样带有校验public void setMaxParallelism(int maxParallelism) { OperatorValidationUtils.validateMaxParallelism(maxParallelism, UPPER_BOUND_MAX_PARALLELISM); this.maxParallelism maxParallelism; }3.2 为什么最大并行度不能设得太大官方文档对此给出了两条重要警告警告为最大并行度设置一个非常大的值将会降低性能因为一些 state backend 需要维持内部的数据结构而这些数据结构将会随着 key-group 的数目而扩张key-group 是状态重新分配的最小单元。从之前的作业恢复时改变该作业的最大并发度将会导致状态不兼容。第二条警告背后的原理可以从 key-group 分配算法得到印证。Flink 恢复 savepoint 时需要把旧作业中每个 key-group 重新映射到新的算子实例上映射公式见 KeyGroupRangeAssignmentpublic static int computeOperatorIndexForKeyGroup( int maxParallelism, int parallelism, int keyGroupId) { return keyGroupId * parallelism / maxParallelism; }可以看到maxParallelism是状态切分与路由的基础分母状态当初是按旧的最大并行度切成 key-group 的恢复时再按新的并行度在这些 key-group 上进行重新划分。一旦 maxParallelism 变化key-group 的划分边界随之改变旧的状态布局就无法与新的划分对齐因此判定为状态不兼容、恢复失败。这也解释了为什么文档反复强调最大并行度应该在作业第一次启动前就想清楚之后不要改动。4. 配置实践建议结合以上源码与文档可以归纳出以下实践要点日常并行度调整优先使用算子层次的setParallelism()做精细调优对整体作业用env.setParallelism()或系统配置parallelism.default设置全局默认值即可使用 savepoint 的作业在首次上线前显式调用setMaxParallelism()或按默认公式估算选择一个覆盖未来扩容需求但不夸张的值例如 2 的幂次方避免日后因扩容需求被迫改变它而导致状态不兼容警惕过大的最大并行度key-group 数目等于最大并行度state backend尤其是 RocksDB 类后端需要维护随 key-group 数量扩张的内部结构过大的值会增加状态开销恢复作业时的并行度变化只要新并行度不超过最大并行度恢复时 key-group 会按keyGroupId * parallelism / maxParallelism公式自动重分配到各算子实例无需手工干预。5. 小结Flink 的并行度按 task算子粒度定义可在算子、执行环境、客户端CLI-p参数或Client.run、系统配置parallelism.default默认值 1四个层次设置最大并行度决定 key-group 的总数是状态切分与扩容的边界可通过setMaxParallelism()显式设定默认值为roundUpToPowerOfTwo(operatorParallelism operatorParallelism / 2)钳制在 [128, 32768] 区间最大并行度一经确定不可随意更改否则从 savepoint 恢复时将因 key-group 布局变化导致状态不兼容相关实现可在 Transformation、KeyGroupRangeAssignment、CoreOptions 中进一步查证。【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表