【问题标题】:stopping spark streaming after reading first batch of data读取第一批数据后停止火花流
【发布时间】:2015-01-30 19:06:28
【问题描述】:

我正在使用火花流来消费 kafka 消息。我想从 kafka 获取一些消息作为样本,而不是阅读所有消息。所以我想读取一批消息,将它们返回给调用者并停止火花流。目前我在 spark 流上下文方法的 awaitTermination 方法中传递 batchInterval 时间。我现在不知道如何将处理后的数据从火花流返回给调用者。这是我目前正在使用的代码

def getsample(params: scala.collection.immutable.Map[String, String]): Unit = {
    if (params.contains("zookeeperQourum"))
      zkQuorum = params.get("zookeeperQourum").get
    if (params.contains("userGroup"))
      group = params.get("userGroup").get
    if (params.contains("topics"))
      topics = params.get("topics").get
    if (params.contains("numberOfThreads"))
      numThreads = params.get("numberOfThreads").get
    if (params.contains("sink"))
      sink = params.get("sink").get
    if (params.contains("batchInterval"))
      interval = params.get("batchInterval").get.toInt
    val sparkConf = new SparkConf().setAppName("KafkaConsumer").setMaster("spark://cloud2-server:7077")
    val ssc = new StreamingContext(sparkConf, Seconds(interval))
    val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
    var consumerConfig = scala.collection.immutable.Map.empty[String, String]
    consumerConfig += ("auto.offset.reset" -> "smallest")
    consumerConfig += ("zookeeper.connect" -> zkQuorum)
    consumerConfig += ("group.id" -> group)
    var data = KafkaUtils.createStream[Array[Byte], Array[Byte], DefaultDecoder, DefaultDecoder](ssc, consumerConfig, topicMap, StorageLevel.MEMORY_ONLY).map(_._2)
    val streams = data.window(Seconds(interval), Seconds(interval)).map(x => new String(x))
    streams.foreach(rdd => rdd.foreachPartition(itr => {
      while (itr.hasNext && size >= 0) {
        var msg=itr.next
        println(msg)
        sample.append(msg)
        sample.append("\n")
        size -= 1
      }
    }))
    ssc.start()
    ssc.awaitTermination(5000)
    ssc.stop(true)
  }

因此,我不想将消息保存在名为“sample”的字符串构建器中,而是返回给调用者。

【问题讨论】:

    标签: apache-kafka spark-skinning


    【解决方案1】:

    你可以实现一个 StreamingListener 然后在里面, onBatchCompleted 你可以调用 ssc.stop()

    private class MyJobListener(ssc: StreamingContext) extends StreamingListener {
    
      override def onBatchCompleted(batchCompleted: StreamingListenerBatchCompleted) = synchronized {
    
        ssc.stop(true)
    
      }
    
    }
    

    这是您将 SparkStreaming 附加到 JobListener 的方式:

    val listen = new MyJobListener(ssc)
    ssc.addStreamingListener(listen)
    
    ssc.start()
    ssc.awaitTermination()
    

    【讨论】:

    • 在 Spark 1.6.1 中,尝试使用您的解决方案时出现以下异常:org.apache.spark.SparkException: Cannot stop StreamingContext within listener thread of AsynchronousListenerBus。任何想法如何解决这个问题?
    • 嗨,我在 spark 2.1 中也面临同样的问题。任何想法,如果你能解决它?谢谢
    • 使用 Spark 2.1,尝试使用您的解决方案时出现以下异常:无法在 org.apache.spark.SparkContext.stop(SparkContext.scala:1785) 的 SparkListenerBus 侦听器线程内停止 SparkContext /跨度>
    【解决方案2】:

    我们可以使用以下代码获取示例消息

    var sampleMessages=streams.repartition(1).mapPartitions(x=>x.take(10))
    

    如果我们想在第一批之后停止,那么我们应该实现自己的 StreamingListener 接口,并且应该在 onBatchCompleted 方法中停止流式传输。

    【讨论】:

      猜你喜欢
      • 2020-05-29
      • 1970-01-01
      • 2016-05-07
      • 1970-01-01
      • 2015-12-11
      • 1970-01-01
      • 1970-01-01
      • 2018-08-20
      • 2016-06-25
      相关资源
      最近更新 更多