最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spark Streaming編程初級實踐詳解

 更新時間:2023年04月20日 09:31:38   作者:WHYBIGDATA  
這篇文章主要為大家介紹了Spark Streaming編程初級實踐詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

寫在前面

  • Linux:CentOS7.5
  • Spark: spark-3.0.0-bin-hadoop3.2
  • Flume:Flume-1.9.0
  • IDE:IntelliJ IDEA2020.2.3

1. 安裝Flume

Flume是Cloudera提供的一個分布式、可靠、可用的系統(tǒng),它能夠?qū)⒉煌瑪?shù)據(jù)源的海量日志數(shù)據(jù)進(jìn)行高效收集、聚合、移動,最后存儲到一個中心化數(shù)據(jù)存儲系統(tǒng)中。Flume 的核心是把數(shù)據(jù)從數(shù)據(jù)源收集過來,再送到目的地。請到Flume官網(wǎng)下載Flume1.7.0安裝文件,下載地址如下:

www.apache.org/dyn/closer.…

或者也可以直接到本教程官網(wǎng)的“下載專區(qū)”中的“軟件”目錄中下載apache-flume-1.7.0-bin.tar.gz。

下載后,把Flume1.7.0安裝到Linux系統(tǒng)的“/usr/local/flume”目錄下,具體安裝和使用方法可以參考教程官網(wǎng)的“實驗指南”欄目中的“日志采集工具Flume的安裝與使用方法。

安裝命令

tar -zxvf apache-flume-1.9.0-bin.tar.gz -C /export/server/
mv apache-flume-1.9.0-bin/ flume-1.9.0
sudo vi /etc/profile
export FLUME_HOME=/usr/local/flume
export PATH=$PATH:$FLUME_HOME/bin
source /etc/profile
mv flume-env.sh.template flume-env.sh
  • 查看版本號
bin/flume-ng version

2.使用Avro數(shù)據(jù)源測試Flume

題目描述

Avro可以發(fā)送一個給定的文件給Flume,Avro 源使用AVRO RPC機(jī)制。請對Flume的相關(guān)配置文件進(jìn)行設(shè)置,從而可以實現(xiàn)如下功能:在一個終端中新建一個文件helloworld.txt(里面包含一行文本“Hello World”),在另外一個終端中啟動Flume以后,可以把helloworld.txt中的文本內(nèi)容顯示出來。

Flume配置文件

al.sources = r1
a1.sinks = k1
a1.channels = c1
a1.sources.r1.type = avro
a1.sources.r1.channels= c1
a1.sources.r1.bind = 0.0.0.0
al.sources.r1.port = 4141
a1.sinks.k1.type = logger
a1.channels.c1.type = memory
al.channels.c1.capacity = 1000
a1.channels.c1.transaction = 100
al.sources.r1.channels = c1
a1.sinks.k1.channel=c1

執(zhí)行命令

  • 先進(jìn)入到Flume安裝目錄,執(zhí)行以下第一行命令;
  • 開始新的一個會話窗口,執(zhí)行第二行命令寫入數(shù)據(jù)到指定的文件中
  • 查看上一步驟中指定的文件內(nèi)容
./bin/flume-ng agent -c . -f ./conf/avro.conf -n a1 -Dflume.root.logger=INFO,console
echo 'hello,world' >> ./log.00
bin/flume-ng avro-client --conf conf -H localhost -p 4141 -F ./log.00

執(zhí)行結(jié)果如下

3. 使用netcat數(shù)據(jù)源測試Flume

題目描述

請對Flume的相關(guān)配置文件進(jìn)行設(shè)置,從而可以實現(xiàn)如下功能:在一個Linux終端(這里稱為“Flume終端”)中,啟動Flume,在另一個終端(這里稱為“Telnet終端”)中,輸入命令“telnet localhost 44444”,然后,在Telnet終端中輸入任何字符,讓這些字符可以順利地在Flume終端中顯示出來。

編寫Flume配置文件

al.sources = r1
a1.sinks = k1
a1.channels = c1
al.sources.r1.type = netcat
al.sources.r1.channels = c1
a1.sources.r1.bind = localhost
al.sources.r1.port = 44444
a1.sinks.k1.type = logger
a1.channels.c1.type = memory
a1.channels.c1.capacity = 1000
al.channels.c1.transaction = 100
al.sources.r1.channels = c1
a1.sinks.k1.channel = c1
  • 執(zhí)行以下命令
./bin/flume-ng agent -c . -f ./netcatExample.conf -n a1 -Dflume.root.logger=INFO,console
telnet localhost 44444
  • 會話窗口成功得到數(shù)據(jù)

4. 使用Flume作為Spark Streaming數(shù)據(jù)源

題目描述

Flume是非常流行的日志采集系統(tǒng),可以作為Spark Streaming的高級數(shù)據(jù)源。請把Flume Source設(shè)置為netcat類型,從終端上不斷給Flume Source發(fā)送各種消息,F(xiàn)lume把消息匯集到Sink,這里把Sink類型設(shè)置為avro,由Sink把消息推送給Spark Streaming,由自己編寫的Spark Streaming應(yīng)用程序?qū)ο⑦M(jìn)行處理。

編寫Flume配置文件

al.sources = r1
a1.sinks = k1
a1.channels =  c1
al.sources.r1.type = netcat
al.sources.r1.bind = localhost
a1.sources.r1.port = 33333
a1.sinks.k1.type = avro
al.sinks.k1.hostname = localhost
a1.sinks.k1.port = 44444
a1.channels.c1.type = memory
al.channels.c1.capacity = 1000000
a1.channels.c1.transactionCapacity = 1000000
al.sources.r1.channels = c1
a1.sinks.k1.channel = c1

主程序代碼

import org.apache.spark.SparkConf
import org.apache.spark.storage.StorageLevel
import org.apache.spark.streaming._
import org.apache.spark.streaming.Milliseconds
import org.apache.spark.streaming.flume._
import org.apache.spark.util.IntParam
object FlumeEventCount {
    def main(args: Array[String]): Unit = {
        if (args.length < 2) {
            System.err.println( "Usage: FlumeEventCount <host> <port>")
            System.exit(1)
        }
        StreamingExamples.setStreamingLogLevels()
        val Array(host, IntParam(port)) = args
        val batchInterval = Milliseconds(2000)
        val sc = new SparkConf()
          .setAppName("FlumeEventCount")
//          .setMaster("local[2]")
        val ssc = new StreamingContext(sc, batchInterval)
        val stream = FlumeUtils.createStream(ssc, host, port, StorageLevel.MEMORY_ONLY_SER_2)
        stream.count().map(cnt => "Received " + cnt + " flume events." ).print()
        ssc.start()
        ssc.awaitTermination()
    }
}

執(zhí)行結(jié)果1

import org.apache.log4j.{Level, Logger}
import org.apache.spark.internal.Logging
object StreamingExamples extends Logging {
    def setStreamingLogLevels(): Unit = {
        val log4jInitialized = Logger.getRootLogger.getAllAppenders.hasMoreElements
        if (!log4jInitialized) {
            logInfo("Setting log level to [WARN] for streaming example." + " To override add a custom log4j.properties to the classpath.")
            Logger.getRootLogger.setLevel(Level.WARN)
        }
    }
}

執(zhí)行結(jié)果2

以上就是Spark Streaming編程初級實踐詳解的詳細(xì)內(nèi)容,更多關(guān)于Spark Streaming編程的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • java中分組統(tǒng)計的三種實現(xiàn)方式

    java中分組統(tǒng)計的三種實現(xiàn)方式

    這篇文章主要介紹了java中分組統(tǒng)計的三種實現(xiàn)方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-07-07
  • 關(guān)于idea中ssm框架的編碼問題分析

    關(guān)于idea中ssm框架的編碼問題分析

    在實際開發(fā)中需要將操作系統(tǒng)編碼、文件編碼、頁面編碼以及tomcat服務(wù)器編碼保持一致,而tomcat在默認(rèn)情況下是使用UTF-8,這就使得其打印的日志文件出現(xiàn)中文亂碼,因此在一般情況下,只需要將tomcat服務(wù)器的編碼改為GBK即可
    2021-06-06
  • 解讀為何java中的boolean類型是32位的

    解讀為何java中的boolean類型是32位的

    這篇文章主要介紹了為何java中的boolean類型是32位的問題,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-04-04
  • java使用zookeeper實現(xiàn)的分布式鎖示例

    java使用zookeeper實現(xiàn)的分布式鎖示例

    這篇文章主要介紹了java使用zookeeper實現(xiàn)的分布式鎖示例,需要的朋友可以參考下
    2014-05-05
  • 詳細(xì)分析JAVA加解密算法

    詳細(xì)分析JAVA加解密算法

    這篇文章主要介紹了JAVA加解密算法的的相關(guān)資料,文中講解非常詳細(xì),代碼幫助大家更好的理解和學(xué)習(xí),感興趣的朋友可以了解下
    2020-06-06
  • Java8?Stream之groupingBy分組使用解讀

    Java8?Stream之groupingBy分組使用解讀

    這篇文章主要介紹了Java8?Stream之groupingBy分組使用,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • Java?講解兩種找二叉樹的最近公共祖先的方法

    Java?講解兩種找二叉樹的最近公共祖先的方法

    樹是一種非線性的數(shù)據(jù)結(jié)構(gòu),它是由n(n>=0)個有限結(jié)點組成一個具有層次關(guān)系的集合,這篇文章主要給大家介紹了關(guān)于Java求解二叉樹的最近公共祖先的相關(guān)資料,需要的朋友可以參考下
    2022-04-04
  • springboot配置mybatis和事務(wù)管理方式

    springboot配置mybatis和事務(wù)管理方式

    這篇文章主要介紹了springboot配置mybatis和事務(wù)管理方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2023-04-04
  • Maven install 報錯

    Maven install 報錯"程序包不存在"問題的解決方法

    這篇文章主要介紹了Maven install 報錯"程序包不存在"問題的解決方法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-07-07
  • Java中的自動拆裝箱、基本類型的轉(zhuǎn)換、包裝類的緩存詳解

    Java中的自動拆裝箱、基本類型的轉(zhuǎn)換、包裝類的緩存詳解

    文章詳細(xì)介紹了Java中數(shù)據(jù)類型的拆裝箱、自動拆箱和裝箱,以及包裝類的緩存機(jī)制,包括基本數(shù)據(jù)類型的容量大小、轉(zhuǎn)換規(guī)則和自動類型轉(zhuǎn)換等
    2024-12-12

最新評論

共和县| 额济纳旗| 易门县| 长岭县| 海盐县| 大洼县| 买车| 尚志市| 诸暨市| 房产| 梅州市| 济南市| 眉山市| 五常市| 永年县| 万州区| 盱眙县| 墨竹工卡县| 沭阳县| 文昌市| 景宁| 水富县| 赤壁市| 华池县| 来凤县| 惠水县| 新余市| 揭东县| 岳普湖县| 万州区| 吉木乃县| 龙门县| 定安县| 定日县| 白水县| 青海省| 昌图县| 虎林市| 同江市| 应用必备| 福鼎市|