如何指定Flink作业的不同部分运行在不同的TaskManager上
多算子分散到不同TaskManager配置方案
- 第一步:拆分算子链,关闭算子合并逻辑
Flink图优化器默认会把上下游可以链式执行的算子合并到同一个线程运行,你需要在需要拆分的算子后调用disableChaining()方法断开算子链,示例代码如下:
如果是SQL作业,可以配置对应作业的DataStream<String> sourceStream = env.addSource(new CustomSource()); // 第一个算子 DataStream<String> processed1 = sourceStream.map(new FirstProcessFunc()) .name("算子1") .disableChaining(); // 第二个算子 DataStream<String> processed2 = processed1.map(new SecondProcessFunc()) .name("算子2") .disableChaining(); // 第三个算子 processed2.addSink(new CustomSink()) .name("算子3");pipeline.operator-chaining参数为false,或者对单个SQL节点加算子链关闭的hint。 - 第二步:为三个算子配置独立的槽位共享组
Flink默认所有算子同属default槽位共享组,会尽可能把同作业的算子调度到同一个slot中。你需要给三个算子分别指定不同的槽位共享组,强制每个算子占用独立的slot,因你每个TaskManager仅配置1个slot,自然会调度到不同的TM上,配置示例:
SQL作业可以通过查询hint指定共享组:DataStream<String> sourceStream = env.addSource(new CustomSource()); // 第一个算子指定共享组1 DataStream<String> processed1 = sourceStream.map(new FirstProcessFunc()) .name("算子1") .slotSharingGroup("group1") .disableChaining(); // 第二个算子指定共享组2 DataStream<String> processed2 = processed1.map(new SecondProcessFunc()) .name("算子2") .slotSharingGroup("group2") .disableChaining(); // 第三个算子指定共享组3 processed2.addSink(new CustomSink()) .name("算子3") .slotSharingGroup("group3");SELECT /*+ SLOT_SHARING_GROUP('group1') */ * FROM source_table - 第三步:可选调度优化配置
提交作业时添加参数-D cluster.evenly-spread-out-slots=true,强制Flink调度slot时尽可能均匀分散到不同的TaskManager节点,避免极端情况多个slot被分配到同一台TM。 - 验证方式
作业启动后进入Flink Web UI,打开作业的「任务管理器」标签页,查看三个算子对应的TaskManager地址,确认是三个不同节点即可。
内容的提问来源于stack exchange,提问作者diegoruizbarbero
相关产品推荐
相关产品推荐

