You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

关于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 04:22:13