【发布时间】: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