【问题标题】:Substitute for thread executor pool in Scala在 Scala 中替代线程执行器池
【发布时间】:2015-01-21 07:35:34
【问题描述】:

我的应用程序要求我有多个线程在运行以从各种 HDFS 节点获取数据。为此,我正在使用线程执行器池和分叉线程。 分叉于:

val pathSuffixList = fileStatuses.getOrElse("FileStatus", List[Any]()).asInstanceOf[List[Map[String, Any]]]
  pathSuffixList.foreach(block => {
    ConsumptionExecutor.execute(new Consumption(webHdfsUri,block))
  })

我的班级消费:

class Consumption(webHdfsUri: String, block:Map[String,Any]) extends Runnable {

      override def run(): Unit = {
        val uriSplit = webHdfsUri.split("\\?")
        val fileOpenUri = uriSplit(0) + "/" + block.getOrElse("pathSuffix", "").toString + "?op=OPEN"
        val inputStream = new URL(fileOpenUri).openStream()
        val datumReader = new GenericDatumReader[Void]()
        val dataStreamReader = new DataFileStream(inputStream, datumReader)
        //        val schema = dataStreamReader.getSchema()
        val dataIterator = dataStreamReader.iterator()
        while (dataIterator.hasNext) {
          println(" data : " + dataStreamReader.next())
        }
      }

    }

消费执行者:

object ConsumptionExecutor{

  val counter: AtomicLong = new AtomicLong()

  val executionContext: ExecutorService = Executors.newCachedThreadPool(new ThreadFactory {
    def newThread(r: Runnable): Thread = {
      val thread: Thread = new Thread(r)
      thread.setName("ConsumptionExecutor-" + counter.incrementAndGet())
      thread
    }
  })
  executionContext.asInstanceOf[ThreadPoolExecutor].setMaximumPoolSize(200)

  def execute(trigger: Runnable) {
    executionContext.execute(trigger)
  }

}

但是,我想在不需要提供固定线程池大小的地方使用 Akka 流式传输/Akka 演员,而 Akka 会处理所有事情。 我对 Akka 以及 Streaming 和 actor 的概念还很陌生。有人可以以示例代码的形式给我任何线索以适合我的用例吗? 提前致谢!

【问题讨论】:

    标签: java multithreading scala akka akka-stream


    【解决方案1】:

    一个想法是为您正在读取的每个 HDFS 节点创建一个ActorPublisher 的(子类)实例,然后将它们作为多个Sources 放入FlowGraph 中。

    类似这样的伪代码,其中省略了ActorPublisher 源的详细信息:

    val g = PartialFlowGraph { implicit b =>
      import FlowGraphImplicits._
      val in1 = actorSource1
      val in2 = actorSource2
      // etc.
    
      val out = UndefinedSink[T]
      val merge = Merge[T]
    
      in1 ~> merge ~> out
      in2 ~> merge
      // etc.
    }
    

    这可以通过迭代它们并为每个actor源添加一个边缘到merge来改进一组actor源,但这给出了这个想法。

    【讨论】:

      猜你喜欢
      • 2011-06-17
      • 2014-08-27
      • 2017-12-29
      • 2021-04-11
      • 1970-01-01
      • 1970-01-01
      • 2017-10-31
      • 1970-01-01
      • 2018-07-28
      相关资源
      最近更新 更多