sparkstreaming with updatestatebykey getting slower and slower

Viewed 8

Consume multiple topics,Different kafka topic fields are processed differently and return the same result type,and union together,use updateStateByKey updateStateByKey Want to get the earliest and latest results corresponding to each key。

But, The program will run slower and slower ,sometimes spark Caused by: java.lang.OutOfMemoryError: Java heap space

If there is something that is unclear, I will add it in time, thank you for your enthusiastic answer

//1. set checkpoint
scc.sparkContext.setCheckpointDir("checkpoint")
//2. kafka Direct
val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> "192.168.44.10:9092",
      "key.deserializer" -> classOf[StringDeserializer],
      "value.deserializer" -> classOf[StringDeserializer],
      "group.id" -> "use_a_separate_group_id_for_each_stream",
      "auto.offset.reset" -> "latest",
      "enable.auto.commit" -> (true: java.lang.Boolean)
)

val topics = Array("test1", "test2", "test3", "test4")

val kafkaStream = KafkaUtils.createDirectStream[String, String](
      scc,
      LocationStrategies.PreferConsistent,
      ConsumerStrategies.Subscribe[String, String](
        topics,
        kafkaParams
      )
    )    val topics = Array("test")

    val kafkaStream = KafkaUtils.createDirectStream[String, String](
      scc,
      LocationStrategies.PreferConsistent,
      ConsumerStrategies.Subscribe[String, String](
        topics,
        kafkaParams
      )
    )    

//3. Different kafka topics  are handled differently and return the same result

dstream1 = kafkaStream.filter("test1").map()...
dstream1 = kafkaStream.filter("test2").map()...
dstream1 = kafkaStream.filter("test3").map()...
dstream1 = kafkaStream.filter("test4").map()...

//4. then combine all results together

dstream  = dstream1.union(dstream2).union(dstream3).union(dstream4)

//5. updatestatebykey
 
dstream.updatestatebykey()

//6. save result to hdfs

0 Answers
Related