diff options
author | Tathagata Das <tathagata.das1565@gmail.com> | 2014-01-13 23:23:46 -0800 |
---|---|---|
committer | Tathagata Das <tathagata.das1565@gmail.com> | 2014-01-13 23:23:46 -0800 |
commit | 4e497db8f3826cf5142b2165a08d02c6f3c2cd90 (patch) | |
tree | 9ec25e86ccf8986035215e51f7b0e1ba1b96dad6 /streaming/src/test/java | |
parent | 1233b3de01be1ff57910786f5f3e2e2a23e228ab (diff) | |
download | spark-4e497db8f3826cf5142b2165a08d02c6f3c2cd90.tar.gz spark-4e497db8f3826cf5142b2165a08d02c6f3c2cd90.tar.bz2 spark-4e497db8f3826cf5142b2165a08d02c6f3c2cd90.zip |
Removed StreamingContext.registerInputStream and registerOutputStream - they were useless as InputDStream has been made to register itself. Also made DStream.register() private[streaming] - not useful to expose the confusing function. Updated a lot of documentation.
Diffstat (limited to 'streaming/src/test/java')
-rw-r--r-- | streaming/src/test/java/org/apache/spark/streaming/JavaTestUtils.scala | 3 |
1 files changed, 1 insertions, 2 deletions
diff --git a/streaming/src/test/java/org/apache/spark/streaming/JavaTestUtils.scala b/streaming/src/test/java/org/apache/spark/streaming/JavaTestUtils.scala index 42ab9590d6..33f6df8f88 100644 --- a/streaming/src/test/java/org/apache/spark/streaming/JavaTestUtils.scala +++ b/streaming/src/test/java/org/apache/spark/streaming/JavaTestUtils.scala @@ -43,7 +43,6 @@ trait JavaTestBase extends TestSuiteBase { implicit val cm: ClassTag[T] = implicitly[ClassTag[AnyRef]].asInstanceOf[ClassTag[T]] val dstream = new TestInputStream[T](ssc.ssc, seqData, numPartitions) - ssc.ssc.registerInputStream(dstream) new JavaDStream[T](dstream) } @@ -57,7 +56,7 @@ trait JavaTestBase extends TestSuiteBase { implicit val cm: ClassTag[T] = implicitly[ClassTag[AnyRef]].asInstanceOf[ClassTag[T]] val ostream = new TestOutputStreamWithPartitions(dstream.dstream) - dstream.dstream.ssc.registerOutputStream(ostream) + ostream.register() } /** |