Flink 分区是指数据流在算子子任务(SubTask)间的分发机制,核心作用是决定上游数据如何路由到下游并行实例,以支持并行计算、负载均衡及状态一致性 。
核心概念
- 物理本质:Flink 中的“分区”对应下游算子的并行子任务(SubTask)。一个流(Stream)被切分为多个分区,每个分区由一个子任务处理,子任务运行在不同的线程、Slot 或节点上 。
- 逻辑作用:控制数据流向。当上下游并行度不一致或需要重分布数据时(如
keyBy、shuffle),必须通过分区器将数据重新分配;若并行度一致且无需重分布,则默认采用直通(Forward)模式 。 - 两类模式:
- One-to-One(窄依赖):数据不跨节点重分布,仅本地转发(如
map、filter默认行为)。 - Redistributing(宽依赖):数据需重新洗牌分发到不同子任务(如
keyBy、rebalance)。
- One-to-One(窄依赖):数据不跨节点重分布,仅本地转发(如
内置分区策略
Flink 提供多种内置分区器,通过 API 调用指定分发规则:
- Forward(转发):默认策略,仅当上下游并行度一致时使用,数据发往本地对应的下游子任务,无网络开销 。
- Rebalance(轮询):循环均匀分发到所有下游子任务,强制跨节点负载均衡,解决数据倾斜 。
- Shuffle(随机):随机选择下游通道,近似均匀分布,但网络开销较大 。
- Rescale(重缩放):类似轮询但仅在本地组内分发,减少跨节点网络 IO,要求上下游并行度成倍数关系 。
- KeyGroupStream(按键分区):
keyBy底层实现,相同 Key 的数据哈希后落入同一分区,保障状态聚合一致性 。 - Broadcast(广播):每条数据复制并发送给所有下游子任务,常用于维表关联 。
- Global(全局):所有数据强制发往下游第一个子任务(ID=0),易导致瓶颈,慎用 。
- Custom(自定义):通过
partitionCustom实现业务特定的分发逻辑 。
关键区别:分区 vs 分组
- 分区(Partitioning):物理/逻辑上将流切分给不同子任务处理,关注数据去哪算(并行度维度)。
- 分组(Grouping):逻辑上将相同 Key 的数据归集,关注哪些数据在一起算(业务维度)。
keyBy同时实现两者:相同 Key 必同分区,但同分区未必同 Key 。