Added StreamingContext.awaitTermination to streaming examples.
This commit is contained in:
parent
792d9084e2
commit
2e95174c45
|
@ -70,5 +70,6 @@ public final class JavaFlumeEventCount {
|
|||
}).print();
|
||||
|
||||
ssc.start();
|
||||
ssc.awaitTermination();
|
||||
}
|
||||
}
|
||||
|
|
|
@ -104,5 +104,6 @@ public final class JavaKafkaWordCount {
|
|||
|
||||
wordCounts.print();
|
||||
jssc.start();
|
||||
jssc.awaitTermination();
|
||||
}
|
||||
}
|
||||
|
|
|
@ -84,5 +84,6 @@ public final class JavaNetworkWordCount {
|
|||
|
||||
wordCounts.print();
|
||||
ssc.start();
|
||||
ssc.awaitTermination();
|
||||
}
|
||||
}
|
||||
|
|
|
@ -80,5 +80,6 @@ public final class JavaQueueStream {
|
|||
|
||||
reducedStream.print();
|
||||
ssc.start();
|
||||
ssc.awaitTermination();
|
||||
}
|
||||
}
|
||||
|
|
|
@ -171,5 +171,6 @@ object ActorWordCount {
|
|||
lines.flatMap(_.split("\\s+")).map(x => (x, 1)).reduceByKey(_ + _).print()
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -60,5 +60,6 @@ object FlumeEventCount {
|
|||
stream.count().map(cnt => "Received " + cnt + " flume events." ).print()
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -50,6 +50,7 @@ object HdfsWordCount {
|
|||
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
|
||||
wordCounts.print()
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
@ -61,6 +61,7 @@ object KafkaWordCount {
|
|||
wordCounts.print()
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
@ -101,5 +101,6 @@ object MQTTWordCount {
|
|||
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
|
||||
wordCounts.print()
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -54,5 +54,6 @@ object NetworkWordCount {
|
|||
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
|
||||
wordCounts.print()
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -61,5 +61,6 @@ object RawNetworkGrep {
|
|||
union.filter(_.contains("the")).count().foreachRDD(r =>
|
||||
println("Grep count: " + r.collect().mkString))
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -114,5 +114,6 @@ object RecoverableNetworkWordCount {
|
|||
createContext(master, ip, port, outputPath)
|
||||
})
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -65,5 +65,6 @@ object StatefulNetworkWordCount {
|
|||
val stateDstream = wordDstream.updateStateByKey[Int](updateFunc)
|
||||
stateDstream.print()
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -110,5 +110,6 @@ object TwitterAlgebirdCMS {
|
|||
})
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -87,5 +87,6 @@ object TwitterAlgebirdHLL {
|
|||
})
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -69,5 +69,6 @@ object TwitterPopularTags {
|
|||
})
|
||||
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
|
@ -91,5 +91,6 @@ object ZeroMQWordCount {
|
|||
val wordCounts = words.map(x => (x, 1)).reduceByKey(_ + _)
|
||||
wordCounts.print()
|
||||
ssc.start()
|
||||
ssc.awaitTermination()
|
||||
}
|
||||
}
|
||||
|
|
Loading…
Reference in a new issue