关于Flink 1.12.0自动添加rebalance的源码位置问询
Flink 1.12.0 自动添加Rebalance的源码实现位置
当上下游算子并行度不同且未指定其他分区策略时,Flink自动替换为Rebalance分区器的逻辑,主要集中在以下两个核心位置:
1. StreamGraph构建时的分区策略替换逻辑
核心逻辑在org.apache.flink.streaming.api.graph.StreamGraphGenerator类的configurePartitioning方法中。
Flink默认会给算子间的边分配ForwardPartitioner(仅适用于上下游并行度相同的场景),该方法会检查当前边的分区器类型以及上下游算子的并行度:
- 如果当前分区器是
ForwardPartitioner,且上下游并行度不一致,就会将分区器替换为RebalancePartitioner。
关键代码片段:
// 检查是否需要替换默认的Forward分区器 if (partitioner instanceof ForwardPartitioner) { int upstreamParallelism = upstreamNode.getParallelism(); int downstreamParallelism = downstreamNode.getParallelism(); if (upstreamParallelism != downstreamParallelism) { partitioner = new RebalancePartitioner<>(); } }
2. Rebalance轮询分发的具体实现
轮询(round-robin)的分发逻辑在org.apache.flink.streaming.runtime.partitioner.RebalancePartitioner类中,核心是selectChannel方法:
该方法通过原子变量记录当前分发的通道索引,每次调用时递增索引并取模下游并行度,实现轮询分配。
关键代码片段:
private final AtomicInteger nextChannelToSendTo = new AtomicInteger(0); @Override public int selectChannel(SerializationDelegate<StreamRecord<T>> record) { // 递增索引并取模,实现轮询 return nextChannelToSendTo.getAndIncrement() % numberOfChannels; }
整个流程总结:用户定义的算子链在生成StreamGraph时,StreamGraphGenerator会自动校验上下游并行度,当并行度不匹配且无自定义分区策略时,自动将默认的Forward分区替换为Rebalance分区,最终由RebalancePartitioner实现轮询分发。
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

