aboutsummaryrefslogtreecommitdiff
path: root/examples
diff options
context:
space:
mode:
authorTathagata Das <tathagata.das1565@gmail.com>2013-02-18 15:18:34 -0800
committerTathagata Das <tathagata.das1565@gmail.com>2013-02-18 15:18:34 -0800
commit12ea14c211da908a278ab19fd1e9f6acd45daae8 (patch)
tree4f76d48f589f23185b680164cedaa9204af8784d /examples
parent6a6e6bda5713ccc6da9ca977321a1fcc6d38a1c1 (diff)
downloadspark-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')
-rw-r--r--examples/src/main/java/spark/streaming/examples/JavaNetworkWordCount.java2
-rw-r--r--examples/src/main/scala/spark/streaming/examples/AkkaActorWordCount.scala12
-rw-r--r--examples/src/main/scala/spark/streaming/examples/NetworkWordCount.scala2
-rw-r--r--examples/src/main/scala/spark/streaming/examples/clickstream/PageViewStream.scala2
4 files changed, 9 insertions, 9 deletions
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<String> lines = ssc.networkTextStream(args[1], Integer.parseInt(args[2]));
+ JavaDStream<String> lines = ssc.socketTextStream(args[1], Integer.parseInt(args[2]));
JavaDStream<String> words = lines.flatMap(new FlatMapFunction<String, String>() {
@Override
public Iterable<String> 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(_))