【发布时间】:2016-02-05 07:29:46
【问题描述】:
我遇到了以下在 Spark Streaming 中处理消息的代码:
val listRDD = ssc.socketTextStream(host, port)
listRDD.foreachRDD(rdd => {
rdd.foreachPartition(partition => {
// Should I start a separate thread for each RDD and/or Partition?
partition.foreach(message => {
Processor.processMessage(message)
})
})
})
这对我有用,但我不确定这是否是最好的方法。我知道 DStream 由“一对多”的 RDD 组成,但是这段代码一个接一个地依次处理 RDD,对吗?难道没有更好的方法 - 一种方法或函数 - 我可以使用以便并行处理 DStream 中的所有 RDD 吗?我应该为每个 RDD 和/或分区启动一个单独的线程吗?我是否误解了这段代码在 Spark 下的工作方式?
不知何故,我认为这段代码没有利用 Spark 中的并行性。
【问题讨论】:
标签: java scala apache-spark spark-streaming