Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- val inputData: mutable.Queue[RDD[String]] = mutable.Queue()
- val inputStream: InputDStream[String] = streamingContext.queueStream(inputData)
- inputData += sparkContext.parallelize(Seq("Test1"))
- inputData += sparkContext.parallelize(Seq("Test2"))
- val iDStream = inputStream.cache()
- println("Size is: "+iDStream.count())
- iDStream.foreachRDD(rdd => println(rdd))
- streamingContext.start()
Add Comment
Please, Sign In to add comment