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

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
...
问题原因
  1. 默认哈希分区的巧合:Flink默认采用MurmurHash对key进行哈希计算后分配分区,对于Long类型的0和1,其哈希值对3取模的结果相同,导致两个key被分配到同一个子任务。
  2. 数据源并行度限制: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:42:42