【问题标题】:Why foreachRDD in Spark Streaming + Kafka is slow, should it be?为什么 Spark Streaming + Kafka 中的 foreachRDD 很慢,应该是这样吗?
【发布时间】:2017-04-13 19:51:19
【问题描述】:

我正在使用 Spark 2.1 和 Kafka 0.08.xx 来执行 Spark Streaming 工作。这是一个文本过滤工作,在这个过程中大部分文本都会被过滤掉。我以两种不同的方式实现:

  1. 直接对 DirectStream 的输出进行过滤:

    val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topics)
    val jsonMsg = messages.map(_._2)
    val filteredMsg = jsonMsg.filter(x=>x.contains(TEXT1) && x.contains(TEXT2) && x.contains(TEXT3))
    
  2. 使用foreachRDD函数

     messages.foreachRDD { rdd => 
               val record = rdd.map(_.2).filter(x => x.contains(TEXT1) &&
                                                     x.contains(TEXT2) &&
                                                     x.contains(TEXT3) )} 
    

我发现第一种方法明显比第二种方法快,但我不确定这是不是常见的情况。

方法一和方法二有区别吗?

【问题讨论】:

    标签: scala apache-spark apache-kafka spark-streaming


    【解决方案1】:

    filter 是一种转换。转换是惰性求值的,也就是说,在您执行某个操作(例如foreachRDD、写入数据等)之前,它们不会执行任何操作。

    所以在 1. 实际上什么都没有发生,因此比 2. 快得多,后者使用动作 foreachRDD 来做某事。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-08-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多