【问题标题】:flink job is not distributed across machinesflink 作业不是跨机器分布的
【发布时间】:2017-10-02 10:30:35
【问题描述】:

我在 Apache flink 中有一个小用例,即批处理系统。我需要处理一组文件。每个文件的处理必须由一台机器处理。我有下面的代码。一直只有一个任务槽被占用,文件一个接一个地处理。我有 6 个节点(所以 6 个任务管理器)并在每个节点中配置了 4 个任务槽。所以,我希望一次处理 24 个文件。

class MyMapPartitionFunction extends RichMapPartitionFunction[java.io.File, Int] {
  override def mapPartition(
      myfiles: java.lang.Iterable[java.io.File],
      out:org.apache.flink.util.Collector[Int])
    : Unit  =  {
    var temp = myfiles.iterator()
    while(temp.hasNext()){
      val fp1 = getRuntimeContext.getDistributedCache.getFile("hadoopRun.sh")
      val file = new File(temp.next().toURI)
      Process(
        "/bin/bash ./run.sh  " + argumentsList(3)+ "/" + file.getName + " " + argumentsList(7) + "/" + file.getName + ".csv",
        new File(fp1.getAbsoluteFile.getParent))
        .lines
        .foreach{println}
      out.collect(1)
    }
  }
}

我以 ./bin/start-cluster.sh 命令启动了 flink,Web 用户界面显示它有 6 个任务管理器,24 个任务槽。

这些文件夹包含大约 49 个文件。当我在这个集合上创建 mapPartition 时,我预计会跨越 49 个并行进程。但是,在我的基础设施中,它们都是一个接一个地处理的。这意味着只有一台机器(一个任务管理器)处理所有 49 个文件名。我想要的是,每个插槽配置 2 个任务,我希望同时处理 24 个文件。

任何指针在这里肯定会有所帮助。我在 flink-conf.yaml 文件中有这些参数

jobmanager.heap.mb: 2048
taskmanager.heap.mb: 1024
taskmanager.numberOfTaskSlots: 4
taskmanager.memory.preallocate: false
parallelism.default: 24

提前致谢。有人可以告诉我我哪里出错了吗?

【问题讨论】:

  • 尝试在mapPartition(new MyMapPartitionFunction())之后添加setParallelism(49)env.fromCollection() 将创建一个并行度为 1 的算子(即使您在 flink-conf.yaml 中将作业并行度配置为 24,因为它使用 NonParallelInput 输入格式)。在不设置并行度的情况下,partition map 运算符将从源继承其并行度。

标签: scala batch-processing apache-flink


【解决方案1】:

正如大卫所描述的,问题是env.fromCollection(Iterable[T]) 创建了一个DataSource 与一个非并行的InputFormat。因此,DataSource1 的并行度执行。后续的操作符 (mapPartition) 从源头继承了这种并行性,因此它们可以被链接起来(这为我们节省了一次网络洗牌)。

解决此问题的方法是通过显式重新平衡源DataSet

env.fromCollection(folders).rebalance()

或在后续运算符 (mapPartition) 中明确设置希望的并行度:

env.fromCollection(folders).mapPartition(...).setParallelism(49)

【讨论】:

  • 非常感谢 Rohrmann 和 David。 rebalance() 看起来更干净,它也很有效!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-02-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-02-23
  • 1970-01-01
相关资源
最近更新 更多