From 4afa3902823aea7b00eedda5f2a029a874abd7e3 Mon Sep 17 00:00:00 2001 From: giwa Date: Sun, 31 Aug 2014 15:56:14 +0900 Subject: [PATCH] clean up code --- .../apache/spark/streaming/api/python/PythonDStream.scala | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) 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() {}