diff --git a/streaming/src/main/scala/org/apache/spark/streaming/api/python/PythonDStream.scala b/streaming/src/main/scala/org/apache/spark/streaming/api/python/PythonDStream.scala index 990feacbdc598..a6f4181547f02 100644 --- a/streaming/src/main/scala/org/apache/spark/streaming/api/python/PythonDStream.scala +++ b/streaming/src/main/scala/org/apache/spark/streaming/api/python/PythonDStream.scala @@ -116,11 +116,13 @@ class PythonForeachDStream( /** * This is a input stream just for the unitest. This is equivalent to a checkpointable, - * replayable, reliable message queue like Kafka. It requires a sequence as input, and - * returns the i_th element at the i_th batch under manual clock. + * replayable, reliable message queue like Kafka. It requires a JArrayList input of JavaRDD, + * and returns the i_th element at the i_th batch under manual clock. */ -class PythonTestInputStream(ssc_ : JavaStreamingContext, inputRDDs: JArrayList[JavaRDD[Array[Byte]]]) +class PythonTestInputStream( + ssc_ : JavaStreamingContext, + inputRDDs: JArrayList[JavaRDD[Array[Byte]]]) extends InputDStream[Array[Byte]](JavaStreamingContext.toStreamingContext(ssc_)) { def start() {}