【问题标题】:Invoking a utility(external) inside Spark streaming job在 Spark 流式传输作业中调用实用程序(外部)
【发布时间】:2017-05-21 07:45:39
【问题描述】:

我有一个使用 Kafka 的流式传输作业(使用 createDstream)。 它的“id”流

[id1,id2,id3 ..]

我有一个实用程序或 api,它接受一个 id 数组并进行一些外部调用并接收一些信息,例如每个 id 的“t”

[id:t1,id2:t2,id3:t3...]

我想在调用实用程序保留 Dstream 时保留 DStream。我不能在 Dstream rdd 上使用地图转换,因为它会调用每个 id,而且该实用程序正在接受 id 的集合。

Dstream.map(x=> myutility(x)) -- ruled out

如果我使用

Dstream.foreachrdd(rdd=> myutility(rdd.collect.toarray))

我失去了DStream。我需要保留DStream 用于下游处理。

【问题讨论】:

  • 重新设计 myutility 使其能够正确并行工作?在 Spark 中拥有单个本地集合是不行的。
  • @user7337271 并行是通过低于 Dstream.foreachrdd(rdd=> myutility(rdd.collect.toarray)) 但丢失 DStream 实现的
  • 这里没有并行性。整个主体 foreachrdd(rdd=> myutility(rdd.collect.toarray)) 在驱动程序上本地执行。你可以transform(rdd=> sc.parallelize(myutility(rdd.collect.toarray))),但它不能解决这个问题。
  • @user7337271 你是对的,我做了错误的假设

标签: scala apache-spark spark-streaming rdd dstream


【解决方案1】:

实现外部批量调用的方法是直接在分区级别转换DStream中的RDD。

图案如下所示:

val transformedStream = dstream.transform{rdd => 
    rdd.mapPartitions{iterator => 
      val externalService = Service.instance() // point to reserve local resources or make server connections.
      val data = iterator.toList // to act in bulk. Need to tune partitioning to avoid huge data loads at this level
      val resultCollection = externalService(data)
      resultCollection.iterator
    }
 }

这种方法使用集群中可用的资源并行处理底层 RDD 的每个分区。请注意,需要为每个分区(而不是每个元素)实例化与外部系统的连接。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-05-09
    • 2014-11-11
    • 1970-01-01
    • 2018-10-17
    • 1970-01-01
    • 2019-05-19
    • 2015-03-21
    • 2018-05-27
    相关资源
    最近更新 更多