【问题标题】:Processing RDDs in a DStream in parallel并行处理 DStream 中的 RDD
【发布时间】: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


    【解决方案1】:

    为了方便和高效,流被划分为小的 RDD(请查看 micro-batching。但您确实不需要将每个 RDD 划分为分区,甚至不需要将流划分为 RDD。

    这完全取决于Processor.processMessage 的真正含义。如果它是单个转换函数,您只需执行 listRDD.map(Processor.processMessage) 即可获得处理消息的任何结果的流,并行计算,无需您做很多其他事情。

    如果Processor 是一个保持状态的可变对象(例如,计算消息的数量),那么事情就更复杂了,因为您需要定义许多这样的对象来解释并行性,并且还需要以某种方式合并结果稍后。

    【讨论】:

    • 在我将其更改为 'listRDD.map(Processor.processMessage)' 后,它停止接收消息。我是否必须使用其他函数调用来“激活”?
    • 好吧,我还是不知道processMessage是个什么样的函数,不过这里就扯淡了。您是否将结果分配给变量,然后对其进行处理(保存、打印等)?
    • 为了快速调试(针对少量数据),请尝试在该语句的末尾添加.print()
    • processMessage 将从消息中“提取”信息,对其进行转换并保存。每条消息都是单独处理的。它不保持状态。它是一个单一的转换函数。
    • “保存”部分可能会迫使您确实使用foreachRDD。如果您依赖提供的save 函数,Spark 会更易于使用。检查这个:spark.apache.org/docs/latest/…
    猜你喜欢
    • 2019-03-31
    • 2020-06-03
    • 2016-05-11
    • 1970-01-01
    • 2017-04-05
    • 2016-06-12
    • 1970-01-01
    • 2017-06-29
    • 1970-01-01
    相关资源
    最近更新 更多