diff options
author | unknown <skumar@SKUMAR01.Trueffect.local> | 2013-05-06 08:05:45 -0600 |
---|---|---|
committer | unknown <skumar@SKUMAR01.Trueffect.local> | 2013-05-06 08:05:45 -0600 |
commit | cbf6a5ee1e7d290d04a0c5dac78d360266d415a4 (patch) | |
tree | 21cf89eec382c69e8494c45f39d0eec62671a5bc /examples | |
parent | 1d54401d7e41095d8cbeeefd42c9d39ee500cd9f (diff) | |
download | spark-cbf6a5ee1e7d290d04a0c5dac78d360266d415a4.tar.gz spark-cbf6a5ee1e7d290d04a0c5dac78d360266d415a4.tar.bz2 spark-cbf6a5ee1e7d290d04a0c5dac78d360266d415a4.zip |
Removed unused code, clarified intent of the program, batch size to 1 second
Diffstat (limited to 'examples')
-rw-r--r-- | examples/src/main/scala/spark/streaming/examples/StatefulNetworkWordCount.scala | 8 |
1 files changed, 3 insertions, 5 deletions
diff --git a/examples/src/main/scala/spark/streaming/examples/StatefulNetworkWordCount.scala b/examples/src/main/scala/spark/streaming/examples/StatefulNetworkWordCount.scala index b662cb1162..51c3c9f9b4 100644 --- a/examples/src/main/scala/spark/streaming/examples/StatefulNetworkWordCount.scala +++ b/examples/src/main/scala/spark/streaming/examples/StatefulNetworkWordCount.scala @@ -4,7 +4,7 @@ import spark.streaming._ import spark.streaming.StreamingContext._ /** - * Counts words in UTF8 encoded, '\n' delimited text received from the network every second. + * Counts words cumulatively in UTF8 encoded, '\n' delimited text received from the network every second. * Usage: StatefulNetworkWordCount <master> <hostname> <port> * <master> is the Spark master URL. In local mode, <master> should be 'local[n]' with n > 1. * <hostname> and <port> describe the TCP server that Spark Streaming would connect to receive data. @@ -15,8 +15,6 @@ import spark.streaming.StreamingContext._ * `$ ./run spark.streaming.examples.StatefulNetworkWordCount local[2] localhost 9999` */ object StatefulNetworkWordCount { - private def className[A](a: A)(implicit m: Manifest[A]) = m.toString - def main(args: Array[String]) { if (args.length < 3) { System.err.println("Usage: StatefulNetworkWordCount <master> <hostname> <port>\n" + @@ -32,8 +30,8 @@ object StatefulNetworkWordCount { Some(currentCount + previousCount) } - // Create the context with a 10 second batch size - val ssc = new StreamingContext(args(0), "NetworkWordCumulativeCountUpdateStateByKey", Seconds(10), + // Create the context with a 1 second batch size + val ssc = new StreamingContext(args(0), "NetworkWordCumulativeCountUpdateStateByKey", Seconds(1), System.getenv("SPARK_HOME"), Seq(System.getenv("SPARK_EXAMPLES_JAR"))) ssc.checkpoint(".") |