如何在Apache Flink中让每个算子独占一个顶点?
如何让Flink算子独立成顶点(拆分算子链)
Flink默认会通过**算子链(Operator Chaining)**优化把上下游符合条件的算子合并到同一个Task中,以此减少数据传输开销、提升性能。你遇到的GroupBy-Aggregation和第二个Filter被合并的情况,就是这个默认优化导致的。
如果要让每个算子独占一个顶点,以便单独测量吞吐量和延迟,可以通过以下几种方式拆分算子链:
针对单个算子禁用链
在需要独立的算子(这里是第二个Filter)上调用disableChaining()方法,它就会和上游的Aggregation算子断开链,成为独立顶点:
// 示例Java代码 stream.keyBy(...) .aggregate(new MyAggregationFunction()) .filter(new MySecondFilterFunction()) .disableChaining(); // 禁用该Filter与上游算子的链合并
开启新的算子链
在Aggregation算子之后调用startNewChain(),让后续的Filter作为新链的起点,同样能实现拆分:
// 示例Java代码 stream.keyBy(...) .aggregate(new MyAggregationFunction()) .startNewChain() .filter(new MySecondFilterFunction());
全局禁用算子链(不推荐)
如果需要所有算子都独立成顶点,可以在执行环境全局禁用算子链,但这会牺牲性能,仅在特殊场景使用:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.disableOperatorChaining(); // 全局禁用所有算子链
完成上述设置后,重新提交作业,就能在作业图中看到GroupBy-Aggregation和第二个Filter成为独立的节点,此时可以分别监控它们的吞吐量、延迟等指标。
内容的提问来源于stack exchange,提问作者Gledi
相关产品推荐
相关产品推荐

