spark-streaming-flume
This commit is contained in:
parent
cbc7c485a8
commit
485ddabe55
@ -18,9 +18,8 @@ object PullBasedWordCount {
|
|||||||
// 1.获取输入流
|
// 1.获取输入流
|
||||||
val flumeStream = FlumeUtils.createPollingStream(ssc, "hadoop001", 8888)
|
val flumeStream = FlumeUtils.createPollingStream(ssc, "hadoop001", 8888)
|
||||||
|
|
||||||
// 2.词频统计
|
// 2.打印输入流中的数据
|
||||||
flumeStream.map(line => new String(line.event.getBody.array()).trim)
|
flumeStream.map(line => new String(line.event.getBody.array()).trim).print()
|
||||||
.flatMap(_.split(" ")).map((_, 1)).reduceByKey(_ + _).print()
|
|
||||||
|
|
||||||
ssc.start()
|
ssc.start()
|
||||||
ssc.awaitTermination()
|
ssc.awaitTermination()
|
||||||
|
@ -19,9 +19,8 @@ object PushBasedWordCount {
|
|||||||
// 1.获取输入流
|
// 1.获取输入流
|
||||||
val flumeStream = FlumeUtils.createStream(ssc, "hadoop001", 8888)
|
val flumeStream = FlumeUtils.createStream(ssc, "hadoop001", 8888)
|
||||||
|
|
||||||
// 2.词频统计
|
// 2.打印输入流的数据
|
||||||
flumeStream.map(line => new String(line.event.getBody.array()).trim)
|
flumeStream.map(line => new String(line.event.getBody.array()).trim).print()
|
||||||
.flatMap(_.split(" ")).map((_, 1)).reduceByKey(_ + _).print()
|
|
||||||
|
|
||||||
ssc.start()
|
ssc.start()
|
||||||
ssc.awaitTermination()
|
ssc.awaitTermination()
|
||||||
|
BIN
pictures/spark-flume-console.png
Normal file
BIN
pictures/spark-flume-console.png
Normal file
Binary file not shown.
After Width: | Height: | Size: 24 KiB |
BIN
pictures/spark-flume-input.png
Normal file
BIN
pictures/spark-flume-input.png
Normal file
Binary file not shown.
After Width: | Height: | Size: 16 KiB |
Loading…
x
Reference in New Issue
Block a user