
在Apache Flink中处理数据流并将其分配到不同的分区partition是实现并行处理的关键手段之一。Flink提供了灵活的机制来控制数据如何分配到不同的并行任务subtasks上。下面是一些使用Flink DataStream API调用8种分区策略的实例这些策略可以帮助你根据不同的需求来控制数据的分区。1. 默认分区Global Partitioning默认情况下当你使用DataStream的keyBy方法但没有指定特定的分区器时Flink会使用全局Global分区策略即所有的数据都会被发送到同一个并行任务上。DataStreamTuple2String, Integer stream ...; stream.keyBy(value - value.f0) .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为82. 哈希分区Hash Partitioning通过keyBy方法可以实现哈希分区这是最常用的分区方式之一。DataStreamTuple2String, Integer stream ...; stream.keyBy(value - value.f0) .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为83. 重新平衡分区Rebalance Partitioning使用rebalance()方法可以将数据均匀地重新分配到下游的所有并行任务中。DataStreamTuple2String, Integer stream ...; stream.rebalance() .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为84. 重缩放分区Rescale Partitioning与rebalance()类似但rescale()主要用于上游和下游的并行度相同时。它会尝试最小化网络传输。DataStreamTuple2String, Integer stream ...; stream.rescale() .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为85. 广播分区Broadcast Partitioning使用broadcast()方法可以将数据流广播到所有下游任务。这在某些类型的全局状态更新场景中很有用。DataStreamTuple2String, Integer stream ...; stream.broadcast() .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为86. 自定义分区Custom Partitioning你可以实现自定义的分区逻辑通过partitionCustom()方法。这需要你提供一个自定义的分区器。DataStreamTuple2String, Integer stream ...; stream.partitionCustom(new MyCustomPartitioner(), keySelector) .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为8其中MyCustomPartitioner是一个实现了org.apache.flink.api.common.operators.base.PartitionerDescriptor接口的类。7. 范围分区Range Partitioning范围分区通常用于有序的数据流通过keyBy()后跟一个有序的数据类型来实现。例如使用元组的第一个字段作为键。ataStreamTuple2String, Integer stream ...; stream.keyBy(0) // 基于元组的第一个字段进行范围分区 .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为8确保下游并行度与上游一致或更大以充分利用范围分区特性。8. 随机分区Random Partitioning使用shuffle()方法可以将数据随机分配到下游的各个任务中。这通常用于需要随机打乱数据顺序的场景。DataStreamTuple2String, Integer stream ...; stream.shuffle() // 随机分区 .map(value - value) // 示例操作 .setParallelism(8); // 设置并行度为89. Forward分区Forward Partitioning使用shuffle()方法可以将数据随机分配到下游的各个任务中。这通常用于需要随机打乱数据顺序的场景。DataStreamTuple2String, Integer stream ...;stream.forward() // Forward分区.map(value - value) // 示例操作.setParallelism(8); // 设置并行度为8Flink 分区是为了解决并行计算时的数据重分布与负载均衡问题分区后的数据合并并非自动发生需通过Union/Connect/Join/CoGroup等算子显式组合或依赖KeyedStream 的状态聚合逻辑 。为什么要分区适配并行度变化上下游算子并行度不一致时必须重分区才能将数据正确路由到所有下游子任务避免数据丢失或资源浪费 。实现负载均衡防止数据倾斜通过 Shuffle/Rebalance 等策略将数据均匀分发提升集群资源利用率 。满足业务逻辑需求如keyBy按键分组确保相同 Key 进入同一子任务以进行状态计算广播分区将配置流分发给所有实例全局分区用于汇总到单点 。跨节点通信基础分布式环境下数据需通过网络传输到特定 TaskManager 的特定 Slot分区器定义了具体的路由规则 。分区后数据如何“合并”Flink 中“合并”指逻辑上的流汇聚不同场景对应不同算子物理分区本身不自动合并数据简单拼接Union场景多条数据类型完全相同的流直接合并如多源日志。机制stream.union(otherStreams...)数据按 FIFO 混合水位线取最小值不进行去重或关联。注意仅合并流结构不改变数据内容 。异构连接Connect场景两条数据类型不同的流需关联处理如订单流 用户画像流。机制stream1.connect(stream2)生成ConnectedStreams需配合CoProcessFunction自定义逻辑可共享状态但流内部独立 。关键常先对双流执行keyBy将相同 Key 路由到同一子任务再在函数内匹配处理 。时间窗口关联Join / Interval Join场景基于时间窗口和 Key 匹配两条流中的事件如点击流 转化流。机制stream1.join(stream2).where(...).equalTo(...).window(...).apply(...)仅在窗口内且 Key 匹配的数据对才会输出非匹配数据丢弃Inner Join 逻辑。前提必须定义 Watermark 以处理事件时间 。分组聚合KeyedStream Reduce/Aggregate场景分区keyBy后对同一 Key 的数据进行统计如求和、计数。机制keyBy将相同 Key 强制路由到同一子任务后续调用reduce/aggregate/sum在该子任务内部逻辑合并状态输出单条结果。这是最典型的“分区后合并计算”模式 。全局汇总Global Sink场景将所有数据强制发送到单个子任务进行最终汇总。机制使用.global()分区算子将所有记录路由到下游第一个子任务并行度实际失效为 1随后在该子任务内聚合或写入 Sink。易造成单点瓶颈慎用 。核心区别Union/Connect是流的物理拼接数据量不变keyBy 聚合是逻辑归并数据量通常减少Join是条件匹配数据量取决于匹配结果。选择取决于业务是需“拼接数据”、“关联分析”还是“统计聚合”。