Hi,
I'm trying to make a stateful stream of Tuple2[String, Dataset[Record]] :

val kafkaStream = KafkaUtils.createDirectStream[String, String,
StringDecoder, StringDecoder](ssc, kafkaParams, topicSet)
val stateStream: DStream[RDD[(String, Record)]] = kafkaStream.map(x=>
{  sqlContext.read.json(x._2).as[Record]}).map(x=>{x.map(r=>(r.iid,r)).rdd})


Because stateStream is a DStream[RDD[(String, Record)]] I can't call
mapWithState on it.
How can I map it to a DStream[(String,Record)] ?

Thank you,
Daniel

Reply via email to