spark流媒体,kafka:java.lang.StackOverflowerr

ql3eal8s  于 2021-06-07  发布在  Kafka
关注(0)|答案(1)|浏览(256)

我在spark流应用程序中得到以下错误,我使用kafka作为输入流。当我用插座的时候,它工作得很好。但当我换成Kafka的时候,我就犯了错误。有人知道为什么会抛出错误吗?我需要更改批处理时间和检查指针时间吗?
error streamingcontext:启动上下文时出错,将其标记为stopped java.lang.StackOverflower错误
我的程序:

def main(args: Array[String]): Unit = {

    // Function to create and setup a new StreamingContext
    def functionToCreateContext(): StreamingContext = {
      val conf = new SparkConf().setAppName("HBaseStream")
      val sc = new SparkContext(conf)
      // create a StreamingContext, the main entry point for all streaming functionality
      val ssc = new StreamingContext(sc, Seconds(5))
      val brokers = args(0)
      val topics= args(1)
      val topicsSet = topics.split(",").toSet
      val kafkaParams = Map[String, String]("metadata.broker.list" -> brokers)
      val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](
        ssc, kafkaParams, topicsSet)

      val inputStream = messages.map(_._2)
//    val inputStream = ssc.socketTextStream(args(0), args(1).toInt)
      ssc.checkpoint(checkpointDirectory)
      inputStream.print(1)
      val parsedStream = inputStream
        .map(line => {
          val splitLines = line.split(",")
          (splitLines(1), splitLines.slice(2, splitLines.length).map((_.trim.toLong)))
        })
      import breeze.linalg.{DenseVector => BDV}
      import scala.util.Try

      val state: DStream[(String, Array[Long])] = parsedStream.updateStateByKey(
        (current: Seq[Array[Long]], prev: Option[Array[Long]]) =>  {
          prev.map(_ +: current).orElse(Some(current))
            .flatMap(as => Try(as.map(BDV(_)).reduce(_ + _).toArray).toOption)
        })

      state.checkpoint(Duration(10000))
      state.foreachRDD(rdd => rdd.foreach(Blaher.blah))
      ssc
    }
    // Get StreamingContext from checkpoint data or create a new one
    val context = StreamingContext.getOrCreate(checkpointDirectory, functionToCreateContext _)
  }
}
a8jjtwal

a8jjtwal1#

尝试删除检查点目录。
我不确定,但您的流式处理上下文似乎无法从检查点恢复。
不管怎么说,这对我很有效。

相关问题