Apache Flink多线程本地执行环境配置及分区不生效问题排查
问题描述
我想用Scala测试Apache Flink的分区功能,编写了一个简单程序:数据源循环生成0和1,按key分组后在处理器中打印线程ID、子任务ID及数据值,期望0和1由不同子任务处理。但无论设置多少并行度,所有数据都被转发到最后一个子任务(输出显示均由子任务2/3处理),无法实现并行。使用Flink版本为1.17.0,程序代码如下:
case class Input(var count: Long = 0) extends Iterator[Long] { override def hasNext: Boolean = count < 1000 override def next(): Long = { count += 1 count % 2 } } case class Process() extends RichFlatMapFunction[Long, Long] { override def flatMap(value: Long, out: Collector[Long]): Unit = { println(s"${Thread.currentThread().getId} ${getRuntimeContext.getIndexOfThisSubtask}/${getRuntimeContext.getNumberOfParallelSubtasks} $value") Thread.sleep(1000) } } def exp003(): Unit = { val env = StreamExecutionEnvironment.createLocalEnvironment(parallelism = 3) env.fromCollection(Input()).keyBy[Long]{x: Long => x}.flatMap(Process()).print() env.execute() }
输出示例:
84 2/3 1 84 2/3 0 84 2/3 1 84 2/3 0 ...
问题原因
- 默认哈希分区的巧合:Flink默认采用
MurmurHash对key进行哈希计算后分配分区,对于Long类型的0和1,其哈希值对3取模的结果相同,导致两个key被分配到同一个子任务。 - 数据源并行度限制:
fromCollection(Input())使用Iterator作为输入时,Flink会将其视为非并行数据源,仅由单个子任务生成数据,但这不是数据集中到同一子任务的核心原因——即使数据源是单并行度,keyBy也应将不同key分发到不同下游子任务。
解决方案
可以通过自定义分区规则替代默认的哈希分区,确保0和1被分配到不同子任务,示例代码如下:
case class Input(var count: Long = 0) extends Iterator[Long] { override def hasNext: Boolean = count < 1000 override def next(): Long = { count += 1 count % 2 } } case class Process() extends RichFlatMapFunction[Long, Long] { override def flatMap(value: Long, out: Collector[Long]): Unit = { println(s"${Thread.currentThread().getId} ${getRuntimeContext.getIndexOfThisSubtask}/${getRuntimeContext.getNumberOfParallelSubtasks} $value") Thread.sleep(1000) } } def exp003(): Unit = { val env = StreamExecutionEnvironment.createLocalEnvironment(parallelism = 3) // 使用partitionCustom自定义分区规则,直接根据key值取模并行度分配 env.fromCollection(Input()) .partitionCustom( new Partitioner[Long] { override def partition(key: Long, numPartitions: Int): Int = { key.toInt % numPartitions } }, x => x ) .flatMap(Process()) .print() env.execute() }
修改后,0会被分配到子任务0/3,1会被分配到子任务1/3,实现预期的并行处理效果。
内容的提问来源于stack exchange,提问作者Xiang Zhang
相关产品推荐
相关产品推荐

