溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊(cè)×
其他方式登錄
點(diǎn)擊 登錄注冊(cè) 即表示同意《億速云用戶服務(wù)條款》

java/scala如何實(shí)現(xiàn)WordCount程序

發(fā)布時(shí)間:2021-12-08 15:20:48 來(lái)源:億速云 閱讀:162 作者:iii 欄目:大數(shù)據(jù)

本篇內(nèi)容介紹了“java/scala如何實(shí)現(xiàn)WordCount程序”的有關(guān)知識(shí),在實(shí)際案例的操作過(guò)程中,不少人都會(huì)遇到這樣的困境,接下來(lái)就讓小編帶領(lǐng)大家學(xué)習(xí)一下如何處理這些情況吧!希望大家仔細(xì)閱讀,能夠?qū)W有所成!

    程序從windows一個(gè)socket端的9999端口接收以換行符分隔的多行文本,每?jī)擅胍粋€(gè)時(shí)間窗口,打印字?jǐn)?shù)統(tǒng)計(jì)。

Socket數(shù)據(jù)發(fā)送命令

window發(fā)送命令 nc -l  -p  9999
linux 發(fā)送命令  nc -lk 9999

Java版本:

package com.unicom.ljs.spark220.study.streaming;
import org.apache.spark.SparkConf;import org.apache.spark.api.java.JavaPairRDD;import org.apache.spark.api.java.function.*;import org.apache.spark.streaming.Durations;import org.apache.spark.streaming.Time;import org.apache.spark.streaming.api.java.JavaDStream;import org.apache.spark.streaming.api.java.JavaPairDStream;import org.apache.spark.streaming.api.java.JavaReceiverInputDStream;import org.apache.spark.streaming.api.java.JavaStreamingContext;import scala.Tuple2;
import java.util.Arrays;import java.util.Iterator;
/** * @author: Created By lujisen * @company ChinaUnicom Software JiNan * @date: 2020-01-30 22:21 * @version: v1.0 * @description: com.unicom.ljs.spark220.study.streaming */public class StreamingWordCount {    public static void main(String[] args) throws InterruptedException {
       SparkConf sparkConf = new SparkConf().setMaster("local[*]").setAppName("StreamingWordCount");        /*這里JavaStreamingContext類似sparkCore的SparkContext        * 帶有兩個(gè)參數(shù)        *           第一個(gè)參數(shù):SparkConf 配置        *           第二個(gè)參數(shù): 每次收取的數(shù)據(jù)流的時(shí)間間隔  作為一個(gè)批次進(jìn)行處理        */        JavaStreamingContext jsc=new JavaStreamingContext(sparkConf, Durations.seconds(2));        /*指定從socket數(shù)據(jù)源接收數(shù)據(jù)        * 指定兩個(gè)參數(shù) 1:主機(jī)名   2:端口        * window發(fā)送命令 nc -l  -p  9999        * linux 發(fā)送命令  nc -lk 9999*/
       JavaReceiverInputDStream<String> sourceDStream = jsc.socketTextStream("localhost", 9999);
       /*接下來(lái)就是對(duì)每個(gè)批次就行處理  這里是每2秒鐘一個(gè)批次 這樣一行行的數(shù)據(jù)流都被拆分為一個(gè)個(gè)的單詞流*/        JavaDStream<String> wordDStream = sourceDStream.flatMap(new FlatMapFunction<String, String>() {            @Override            public Iterator<String> call(String line) throws Exception {                return Arrays.asList(line.split(" ")).iterator();            }        });        /*轉(zhuǎn)換成  hello  1        *         world  1        *         a      1        *         b      1 格式*/        JavaPairDStream<String, Integer> wordPairDStream = wordDStream.mapToPair(new PairFunction<String, String, Integer>() {            @Override            public Tuple2<String, Integer> call(String word) throws Exception {                return new Tuple2<>(word, 1);            }        });        JavaPairDStream<String, Integer>  wordCountResult = wordPairDStream.reduceByKey(new Function2<Integer, Integer, Integer>() {            @Override            public Integer call(Integer v1, Integer v2) throws Exception {                return v1+v2;            }        });
       /*打印結(jié)果*/        wordCountResult.print();
       /*jsc這里必須要調(diào)用start()函數(shù)application才會(huì)啟動(dòng)執(zhí)行,接收數(shù)據(jù)*/        jsc.start();        jsc.awaitTermination();        /*停止*/        jsc.stop();    }}

Scala版本:

package com.unicom.ljs.study.streaming
import org.apache.spark.SparkConfimport org.apache.spark.streaming.StreamingContextimport org.apache.spark.streaming.Seconds
/**  * @author: Created By lujisen  * @company ChinaUnicom Software JiNan  * @date: 2020-01-31 08:59  * @version: v1.0  * @description: com.unicom.ljs.study.streaming  */object StreamingWordCount {  def main(args: Array[String]): Unit = {
   /*構(gòu)建SparkConf配置*/    val sparkConf  =new SparkConf().setMaster("local[*]").setAppName("StreamingWordCountScala")    val ssc =new StreamingContext(sparkConf,Seconds(2))
   /*指定socket數(shù)據(jù)源*/    val sourceDStream=ssc.socketTextStream("localhost",9999)
   val  wordDStream=sourceDStream.flatMap(x=>x.split(" "))
   val  wordPairDStream=wordDStream.map(x=>(x,1))    val  wordCountResult=wordPairDStream.reduceByKey(_+_)
   /*打印結(jié)果*/    wordCountResult.print()    /*啟動(dòng)*/    ssc.start()    ssc.awaitTermination()    /*停止*/    ssc.stop()  }}

java/scala如何實(shí)現(xiàn)WordCount程序

“java/scala如何實(shí)現(xiàn)WordCount程序”的內(nèi)容就介紹到這里了,感謝大家的閱讀。如果想了解更多行業(yè)相關(guān)的知識(shí)可以關(guān)注億速云網(wǎng)站,小編將為大家輸出更多高質(zhì)量的實(shí)用文章!

向AI問(wèn)一下細(xì)節(jié)

免責(zé)聲明:本站發(fā)布的內(nèi)容(圖片、視頻和文字)以原創(chuàng)、轉(zhuǎn)載和分享為主,文章觀點(diǎn)不代表本網(wǎng)站立場(chǎng),如果涉及侵權(quán)請(qǐng)聯(lián)系站長(zhǎng)郵箱:is@yisu.com進(jìn)行舉報(bào),并提供相關(guān)證據(jù),一經(jīng)查實(shí),將立刻刪除涉嫌侵權(quán)內(nèi)容。

AI