diff options
author | Tathagata Das <tathagata.das1565@gmail.com> | 2013-02-18 15:18:34 -0800 |
---|---|---|
committer | Tathagata Das <tathagata.das1565@gmail.com> | 2013-02-18 15:18:34 -0800 |
commit | 12ea14c211da908a278ab19fd1e9f6acd45daae8 (patch) | |
tree | 4f76d48f589f23185b680164cedaa9204af8784d /examples/src/main/scala | |
parent | 6a6e6bda5713ccc6da9ca977321a1fcc6d38a1c1 (diff) | |
download | spark-12ea14c211da908a278ab19fd1e9f6acd45daae8.tar.gz spark-12ea14c211da908a278ab19fd1e9f6acd45daae8.tar.bz2 spark-12ea14c211da908a278ab19fd1e9f6acd45daae8.zip |
Changed networkStream to socketStream and pluggableNetworkStream to become networkStream as a way to create streams from arbitrary network receiver.
Diffstat (limited to 'examples/src/main/scala')
3 files changed, 8 insertions, 8 deletions
diff --git a/examples/src/main/scala/spark/streaming/examples/AkkaActorWordCount.scala b/examples/src/main/scala/spark/streaming/examples/AkkaActorWordCount.scala index ff05842c71..553afc2024 100644 --- a/examples/src/main/scala/spark/streaming/examples/AkkaActorWordCount.scala +++ b/examples/src/main/scala/spark/streaming/examples/AkkaActorWordCount.scala @@ -36,8 +36,8 @@ class SampleActorReceiver[T: ClassManifest](urlOfPublisher: String) } /** - * A sample word count program demonstrating the use of plugging in - * AkkaActor as Receiver + * A sample word count program demonstrating the use of Akka actor stream. + * */ object AkkaActorWordCount { def main(args: Array[String]) { @@ -56,18 +56,18 @@ object AkkaActorWordCount { Seconds(batchDuration.toLong)) /* - * Following is the use of pluggableActorStream to plug in custom actor as receiver + * Following is the use of actorStream to plug in custom actor as receiver * * An important point to note: * Since Actor may exist outside the spark framework, It is thus user's responsibility - * to ensure the type safety, i.e type of data received and PluggableInputDstream + * to ensure the type safety, i.e type of data received and actorStream * should be same. * - * For example: Both pluggableActorStream and SampleActorReceiver are parameterized + * For example: Both actorStream and SampleActorReceiver are parameterized * to same type to ensure type safety. */ - val lines = ssc.pluggableActorStream[String]( + val lines = ssc.actorStream[String]( Props(new SampleActorReceiver[String]("akka://spark@%s:%s/user/FeederActor".format( remoteAkkaHost, remoteAkkaPort.toInt))), "SampleReceiver") diff --git a/examples/src/main/scala/spark/streaming/examples/NetworkWordCount.scala b/examples/src/main/scala/spark/streaming/examples/NetworkWordCount.scala index 32f7d57bea..7ff70ae2e5 100644 --- a/examples/src/main/scala/spark/streaming/examples/NetworkWordCount.scala +++ b/examples/src/main/scala/spark/streaming/examples/NetworkWordCount.scala @@ -27,7 +27,7 @@ object NetworkWordCount { // Create a NetworkInputDStream on target ip:port and count the // words in input stream of \n delimited test (eg. generated by 'nc') - val lines = ssc.networkTextStream(args(1), args(2).toInt) + val lines = ssc.socketTextStream(args(1), args(2).toInt) val words = lines.flatMap(_.split(" ")) val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _) wordCounts.print() diff --git a/examples/src/main/scala/spark/streaming/examples/clickstream/PageViewStream.scala b/examples/src/main/scala/spark/streaming/examples/clickstream/PageViewStream.scala index 60f228b8ad..fba72519a9 100644 --- a/examples/src/main/scala/spark/streaming/examples/clickstream/PageViewStream.scala +++ b/examples/src/main/scala/spark/streaming/examples/clickstream/PageViewStream.scala @@ -27,7 +27,7 @@ object PageViewStream { val ssc = new StreamingContext("local[2]", "PageViewStream", Seconds(1)) // Create a NetworkInputDStream on target host:port and convert each line to a PageView - val pageViews = ssc.networkTextStream(host, port) + val pageViews = ssc.socketTextStream(host, port) .flatMap(_.split("\n")) .map(PageView.fromString(_)) |