From 12ea14c211da908a278ab19fd1e9f6acd45daae8 Mon Sep 17 00:00:00 2001 From: Tathagata Das Date: Mon, 18 Feb 2013 15:18:34 -0800 Subject: Changed networkStream to socketStream and pluggableNetworkStream to become networkStream as a way to create streams from arbitrary network receiver. --- .../java/spark/streaming/examples/JavaNetworkWordCount.java | 2 +- .../scala/spark/streaming/examples/AkkaActorWordCount.scala | 12 ++++++------ .../scala/spark/streaming/examples/NetworkWordCount.scala | 2 +- .../streaming/examples/clickstream/PageViewStream.scala | 2 +- 4 files changed, 9 insertions(+), 9 deletions(-) (limited to 'examples/src') diff --git a/examples/src/main/java/spark/streaming/examples/JavaNetworkWordCount.java b/examples/src/main/java/spark/streaming/examples/JavaNetworkWordCount.java index 4299febfd6..07342beb02 100644 --- a/examples/src/main/java/spark/streaming/examples/JavaNetworkWordCount.java +++ b/examples/src/main/java/spark/streaming/examples/JavaNetworkWordCount.java @@ -35,7 +35,7 @@ public class JavaNetworkWordCount { // Create a NetworkInputDStream on target ip:port and count the // words in input stream of \n delimited test (eg. generated by 'nc') - JavaDStream lines = ssc.networkTextStream(args[1], Integer.parseInt(args[2])); + JavaDStream lines = ssc.socketTextStream(args[1], Integer.parseInt(args[2])); JavaDStream words = lines.flatMap(new FlatMapFunction() { @Override public Iterable call(String x) { 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(_)) -- cgit v1.2.3